diff --git a/src/Makefile.am b/src/Makefile.am index b4f06947..3abcaad7 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -46,6 +46,7 @@ utun_CORE_SOURCES = \ eim_nat.c \ nat_transport.c \ media_delivery/media_delivery.c \ + media_delivery/media_download.c \ media_delivery/media_index.c \ media_async/media_async.c \ transport_layer/dummynet.c \ @@ -115,6 +116,7 @@ libutun_a_SOURCES = \ eim_nat.c \ nat_transport.c \ media_delivery/media_delivery.c \ + media_delivery/media_download.c \ media_delivery/media_index.c \ media_async/media_async.c \ transport_layer/dummynet.c \ diff --git a/src/chat/member_sync.c b/src/chat/member_sync.c index c05a6d82..be5efbb2 100644 --- a/src/chat/member_sync.c +++ b/src/chat/member_sync.c @@ -234,6 +234,33 @@ static int _member_get_items(void* ctx, const char* ns, uint8_t level, return 0; } +/* ── node_props_changed (adm_tags) callbacks ── */ + +struct ms_props_cbk { node_props_changed_fn fn; void* arg; struct ms_props_cbk* next; }; +static struct ms_props_cbk* g_props_cbks = NULL; + +void member_sync_add_props_cbk(node_props_changed_fn fn, void* arg) { + if (!fn) return; + struct ms_props_cbk* e = u_malloc(sizeof(*e)); + if (!e) return; + e->fn = fn; e->arg = arg; + e->next = g_props_cbks; + g_props_cbks = e; +} + +void member_sync_remove_props_cbk(node_props_changed_fn fn, void* arg) { + if (!fn) return; + struct ms_props_cbk** p = &g_props_cbks; + while (*p) { + if ((*p)->fn == fn && (*p)->arg == arg) { + struct ms_props_cbk* r = *p; + *p = r->next; u_free(r); + return; + } + p = &(*p)->next; + } +} + static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer, const uint8_t* data, size_t len) { struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; @@ -356,6 +383,13 @@ static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer, if (adm_ver < 0) { mp += consumed; mrem -= (size_t)consumed; continue; } member_sync_put(inst, ns, nid, x25, ed, jsig, jts, usig, uts, nm, addrs, (int)ac, atags[0] ? atags : NULL, atsig, adm_storage); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] member_sync recv ns=%s nid=0x%016llx ac=%d adm_ver=%d storage=%d", MS_ID, ns, (unsigned long long)nid, ac, adm_ver, adm_storage); + /* fire node_props_changed callbacks */ + { struct ms_props_cbk* pc = g_props_cbks; + if (pc) { + const char* tags_final = atags[0] ? atags : NULL; + while (pc) { pc->fn(nid, tags_final, pc->arg); pc = pc->next; } + } + } mp += consumed; mrem -= (size_t)consumed; } if (count > 0) diff --git a/src/chat/member_sync.h b/src/chat/member_sync.h index 7420313f..bbc8660e 100644 --- a/src/chat/member_sync.h +++ b/src/chat/member_sync.h @@ -128,6 +128,15 @@ const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, typedef void (*member_sync_node_updated_fn)(uint64_t node_id, int online); void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb); +/* + * Коллбэк: изменились adm_tags любого узла (включая себя). + * Вызывается при успешной обработке MSG_ITEM_UPDATE с adm_tags. + * Многоподписочный — можно добавить несколько подписчиков. + */ +typedef void (*node_props_changed_fn)(uint64_t node_id, const char* adm_tags, void* arg); +void member_sync_add_props_cbk(node_props_changed_fn fn, void* arg); +void member_sync_remove_props_cbk(node_props_changed_fn fn, void* arg); + #ifdef __cplusplus } #endif diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index a8ef93da..9dba170a 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -1,28 +1,777 @@ #include "media_delivery.h" -#include "utun_instance.h" -#include "etcp_router.h" +#include "media_delivery_proto.h" +#include "media_download.h" +#include "../utun_instance.h" +#include "../routing_layer/etcp_router.h" +#include "../routing_layer/topo_group.h" +#include "../routing_layer/topo_node.h" +#include "../routing_layer/topo_node_sqlite.h" +#include "../routing_layer/conn_mgr.h" +#include "../transport_layer/etcp_api.h" +#include "../transport_layer/etcp.h" +#include "../chat/member_sync.h" #include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/u_async.h" +#include "../lib/sqlite3.h" #include +#include -static void media_delivery_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - (void)conn; (void)entry; +#define MD_ID "media_delivery" + +/* ── forward decls ── */ + +static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry); +static void md_on_conn_status(struct ETCP_CONN* conn, int status, void* arg); +static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event, void* arg); +static void md_on_props_changed(uint64_t node_id, const char* adm_tags, void* arg); +static int md_super_start(struct media_delivery_ctx* md); +static void md_super_stop(struct media_delivery_ctx* md); +static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super_peer* peer); +static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id); +static int md_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, const uint8_t* data, size_t len); +static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id); + +/* ── helpers ── */ + +static const char* md_subcmd_name(uint8_t sc) { + switch (sc) { + case MEDIA_SUBCMD_SERVE_REG: return "SERVE_REG"; + case MEDIA_SUBCMD_SERVE_ACK: return "SERVE_ACK"; + case MEDIA_SUBCMD_SERVE_LEAVE: return "SERVE_LEAVE"; + case MEDIA_SUBCMD_QUERY: return "QUERY"; + case MEDIA_SUBCMD_QUERY_RESP: return "QUERY_RESP"; + case MEDIA_SUBCMD_HAVE_BLOCK: return "HAVE_BLOCK"; + case MEDIA_SUBCMD_HAVE_BLOCK_ACK: return "HAVE_BLOCK_ACK"; + case MEDIA_SUBCMD_SUPER_REPL: return "SUPER_REPL"; + case MEDIA_SUBCMD_SUPER_ACK: return "SUPER_ACK"; + case MEDIA_SUBCMD_BLOCK_REQ: return "BLOCK_REQ"; + case MEDIA_SUBCMD_BLOCK_CHUNK: return "BLOCK_CHUNK"; + case MEDIA_SUBCMD_BLOCK_DONE: return "BLOCK_DONE"; + case MEDIA_SUBCMD_CANCEL: return "CANCEL"; + case MEDIA_SUBCMD_SUPER_HELLO: return "SUPER_HELLO"; + default: return "???"; + } +} + +/* ── SQL helper: insert/update block_availability ── */ + +static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_id, + uint64_t node_id, uint32_t chunk, int64_t ts) { + const char* sql = + "INSERT OR REPLACE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp)" + " VALUES(?,?,?,?,?)"; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1; + sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC); + sqlite3_bind_int64(st, 2, (sqlite3_int64)group_id); + sqlite3_bind_int64(st, 3, (sqlite3_int64)node_id); + sqlite3_bind_int(st, 4, (int)chunk); + sqlite3_bind_int64(st, 5, (sqlite3_int64)ts); + int rc = sqlite3_step(st); + sqlite3_finalize(st); + return (rc == SQLITE_DONE) ? 0 : -1; +} + +static int md_ba_find_nodes(sqlite3* db, const uint8_t* block_uuid, + uint64_t* out_node_ids, int max) { + const char* sql = + "SELECT node_id, timestamp FROM block_availability WHERE block_uuid=? ORDER BY timestamp DESC LIMIT ?"; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0; + sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC); + sqlite3_bind_int(st, 2, max); + int n = 0; + while (n < max && sqlite3_step(st) == SQLITE_ROW) { + out_node_ids[n] = (uint64_t)sqlite3_column_int64(st, 0); + n++; + } + sqlite3_finalize(st); + return n; +} + +static int md_ba_get_since(sqlite3* db, uint64_t since_id, + uint8_t* out_buf, int max_bytes) { + const char* sql = + "SELECT id,block_uuid,group_id,node_id,chunk,timestamp" + " FROM block_availability WHERE id > ? ORDER BY id ASC LIMIT 4"; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0; + sqlite3_bind_int64(st, 1, (sqlite3_int64)since_id); + int off = 0; + while (sqlite3_step(st) == SQLITE_ROW) { + if (off + 52 > max_bytes) break; + int64_t id = sqlite3_column_int64(st, 0); + memcpy(out_buf + off, &id, 8); off += 8; + memcpy(out_buf + off, sqlite3_column_blob(st, 1), 16); off += 16; + int64_t gid = sqlite3_column_int64(st, 2); memcpy(out_buf + off, &gid, 8); off += 8; + int64_t nid = sqlite3_column_int64(st, 3); memcpy(out_buf + off, &nid, 8); off += 8; + int32_t ch = sqlite3_column_int(st, 4); memcpy(out_buf + off, &ch, 4); off += 4; + int64_t ts = sqlite3_column_int64(st, 5); memcpy(out_buf + off, &ts, 8); off += 8; + } + sqlite3_finalize(st); + return off; +} + +static uint64_t md_ss_get(sqlite3* db, uint64_t peer_node_id) { + const char* sql = "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?"; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0; + sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id); + uint64_t v = 0; + if (sqlite3_step(st) == SQLITE_ROW) v = (uint64_t)sqlite3_column_int64(st, 0); + sqlite3_finalize(st); + return v; +} + +static void md_ss_set(sqlite3* db, uint64_t peer_node_id, uint64_t last_recv_id) { + const char* sql = "INSERT OR REPLACE INTO super_sync(peer_node_id,last_recv_id) VALUES(?,?)"; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return; + sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id); + sqlite3_bind_int64(st, 2, (sqlite3_int64)last_recv_id); + sqlite3_step(st); + sqlite3_finalize(st); +} + +/* ── super_peer find/create ── */ + +static struct media_super_peer* md_super_peer_find(struct media_delivery_ctx* md, uint64_t node_id) { + if (!md->super_peers) return NULL; + struct ll_entry* e = queue_find_data_by_index(md->super_peers, &node_id); + return e ? (struct media_super_peer*)e->data : NULL; +} + +static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id) { + struct media_super_peer* sp = md_super_peer_find(md, node_id); + if (sp) return sp; + if (!md->super_peers) { + md->super_peers = queue_new(NULL, 64, 0, 8, "md_super"); + if (!md->super_peers) return NULL; + } + struct ll_entry* qe = queue_entry_new(sizeof(struct media_super_peer)); + if (!qe) return NULL; + sp = (struct media_super_peer*)qe->data; + memset(sp, 0, sizeof(*sp)); + sp->peer_node_id = node_id; + memcpy(qe->data, &node_id, 8); + queue_data_put_with_index(md->super_peers, qe); + return sp; +} + +static void md_super_peer_remove(struct media_delivery_ctx* md, uint64_t node_id) { + if (!md->super_peers) return; + struct ll_entry* e = queue_find_data_by_index(md->super_peers, &node_id); + if (!e) return; + struct media_super_peer* sp = (struct media_super_peer*)e->data; + if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; } + if (sp->connect_timer) { uasync_cancel_timeout(md->inst->ua, sp->connect_timer); sp->connect_timer = NULL; } + queue_remove_data(md->super_peers, e); + queue_entry_free(e); +} + +/* ── replication timers ── */ + +static void md_repl_retrans_cb(void* arg) { + struct media_super_peer* peer = (struct media_super_peer*)arg; + struct media_delivery_ctx* md = (struct media_delivery_ctx*)peer->connect_timer; + if (!md || !peer->connected) return; + peer->timeout_timer = NULL; + peer->timeout_tb = peer->timeout_tb < MEDIA_REPL_TIMEOUT_MAX_TB / 2 + ? peer->timeout_tb * 2 : peer->timeout_tb; + if (peer->timeout_tb > MEDIA_REPL_TIMEOUT_MAX_TB) + peer->timeout_tb = MEDIA_REPL_TIMEOUT_MAX_TB; + peer->inflight_count = 0; + md_super_repl_send(md, peer); +} + +static void md_repl_schedule_retrans(struct media_delivery_ctx* md, struct media_super_peer* peer) { + if (peer->timeout_timer) uasync_cancel_timeout(md->inst->ua, peer->timeout_timer); + peer->timeout_timer = uasync_set_timeout(md->inst->ua, peer->timeout_tb, peer, + md_repl_retrans_cb, "md_repl_retrans"); +} + +static void md_reconnect_cb(void* arg) { + struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; + struct media_super_peer* peer = (struct media_super_peer*)md->super_peers; + while (peer) { + if (!peer->connected && peer->connect_timer) { + peer->connect_timer = NULL; + md_super_connect(md, peer->peer_node_id); + } + peer = peer->connected ? NULL : peer; + } +} + +/* ── replication send ── */ + +static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super_peer* peer) { + if (!md->db || !peer->connected || !peer->hello_done) return; + + while (peer->inflight_count < MEDIA_MAX_REPL_INFLIGHT) { + uint8_t buf[512]; + int nb = md_ba_get_since(md->db, peer->peer_last_recv_id, buf, (int)sizeof(buf)); + if (nb == 0) { + if (peer->timeout_timer) uasync_cancel_timeout(md->inst->ua, peer->timeout_timer); + peer->timeout_timer = NULL; + return; + } + + /* build packet: subcmd + seq + num + raw rows */ + uint8_t pkt[520]; + int off = 0; + pkt[off++] = MEDIA_SUBCMD_SUPER_REPL; + + int num_entries = nb / 52; + int64_t max_id = 0; + const uint8_t* rp = buf; + for (int i = 0; i < num_entries; i++) { + int64_t id; memcpy(&id, rp, 8); + if (id > max_id) max_id = id; + rp += 52; + } + + uint32_t seq = (uint32_t)max_id; memcpy(pkt + off, &seq, 4); off += 4; + uint16_t num = (uint16_t)num_entries; memcpy(pkt + off, &num, 2); off += 2; + memcpy(pkt + off, buf, (size_t)nb); off += nb; + + if (md_send(md->inst, peer->peer_node_id, pkt, (size_t)off) == 0) { + peer->inflight_count++; + DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL sent to 0x%016llx seq=%u entries=%d", + MD_ID, (unsigned long long)peer->peer_node_id, seq, num_entries); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_REPL send failed to 0x%016llx", + MD_ID, (unsigned long long)peer->peer_node_id); + return; + } + + if (peer->inflight_count >= MEDIA_MAX_REPL_INFLIGHT) break; + } + + if (peer->inflight_count > 0) md_repl_schedule_retrans(md, peer); +} + +/* ── send via etcp_router helper ── */ + +static int md_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, + const uint8_t* data, size_t len) { + if (!inst || !data || len == 0) return -1; + struct ll_entry* e = queue_entry_new(0); + if (!e) return -1; + e->dgram = u_malloc(len); + if (!e->dgram) { queue_entry_free(e); return -1; } + memcpy(e->dgram, data, len); e->len = (uint16_t)len; + int rc = etcp_route_send(inst, 0, dst_node_id, e, 1); + if (rc != 0) { u_free(e->dgram); queue_entry_free(e); } + return rc; +} + +/* ── subcommand handlers ── */ + +static void md_handle_serve_reg(struct media_delivery_ctx* md, uint64_t from_node, + const uint8_t* data, size_t len) { + if (len < sizeof(struct media_pkt_serve_reg)) return; + struct media_pkt_serve_reg* pkt = (struct media_pkt_serve_reg*)data; + + if (!md->served_nodes) { + md->served_nodes = queue_new(NULL, 64, 0, 8, "md_served"); + if (!md->served_nodes) return; + } + struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &from_node); + if (!e) { + struct ll_entry* qe = queue_entry_new(sizeof(struct media_served_node)); + if (!qe) return; + struct media_served_node* sn = (struct media_served_node*)qe->data; + memset(sn, 0, sizeof(*sn)); + sn->node_id = from_node; sn->group_id = pkt->group_id; + sn->joined_at = (int64_t)time(NULL); + memcpy(qe->data, &from_node, 8); + queue_data_put_with_index(md->served_nodes, qe); + } + + uint8_t ack[3] = { MEDIA_SUBCMD_SERVE_ACK, 0, 0 }; + memcpy(ack + 1, &pkt->group_id, 8); + (void)from_node; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SERVE_REG node=0x%016llx group=%016llx", + MD_ID, (unsigned long long)from_node, (unsigned long long)pkt->group_id); +} + +static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node, + const uint8_t* data, size_t len) { + if (len < MEDIA_QUERY_HDR_SIZE || !md->db) return; + struct media_pkt_query* q = (struct media_pkt_query*)data; + uint16_t nb = q->num_blocks; + if (nb > 64) nb = 64; + uint8_t* bid = (uint8_t*)data + MEDIA_QUERY_HDR_SIZE; + + uint8_t resp[2048]; int off = 0; + resp[off++] = MEDIA_SUBCMD_QUERY_RESP; + int num_pos = off; off += 2; /* reserve for num_entries */ + int count = 0; + + for (int i = 0; i < nb && off + 30 <= (int)sizeof(resp); i++) { + uint64_t nodes[10]; + int nf = md_ba_find_nodes(md->db, bid + i * 16, nodes, 10); + for (int j = 0; j < nf && off + 30 <= (int)sizeof(resp); j++) { + uint16_t rtt = 0; /* supernode doesn't store RTT; client will probe */ + memcpy(resp + off, &nodes[j], 8); off += 8; + memcpy(resp + off, &rtt, 2); off += 2; + memcpy(resp + off, bid + i * 16, 16); off += 16; + uint32_t ch = q->num_blocks; + memcpy(resp + off, &ch, 4); off += 4; + count++; + } + } + uint16_t nc = (uint16_t)count; + memcpy(resp + num_pos, &nc, 2); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: QUERY from 0x%016llx → %d entries", + MD_ID, (unsigned long long)from_node, count); + (void)md_send(md->inst, from_node, resp, (size_t)off); +} + +static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_node, + const uint8_t* data, size_t len) { + if (len < MEDIA_HAVE_BLOCK_SIZE || !md->db) return; + struct media_pkt_have_block* hb = (struct media_pkt_have_block*)data; + + int rc = md_ba_insert(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp); + + uint8_t ack[sizeof(struct media_pkt_have_block_ack)]; + struct media_pkt_have_block_ack* ha = (struct media_pkt_have_block_ack*)ack; + memset(ack, 0, sizeof(ack)); + ha->subcmd = MEDIA_SUBCMD_HAVE_BLOCK_ACK; + memcpy(ha->media_id, hb->media_id, 16); + memcpy(ha->block_id, hb->block_id, 16); + ha->status = (rc == 0) ? 0 : 1; + + md_send(md->inst, from_node, ack, sizeof(ack)); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK from 0x%016llx rc=%d", + MD_ID, (unsigned long long)from_node, rc); + + if (rc == 0) { + /* replicate to connected super-peers */ + struct ll_entry* se = md->super_peers ? md->super_peers->head : NULL; + while (se) { + struct media_super_peer* sp = (struct media_super_peer*)se->data; + if (sp->connected && sp->hello_done && sp->inflight_count < MEDIA_MAX_REPL_INFLIGHT) + md_super_repl_send(md, sp); + se = se->next; + } + } +} + +static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_node, + const uint8_t* data, size_t len) { + if (len < MEDIA_SUPER_REPL_HDR_SIZE || !md->db) return; + struct media_pkt_super_repl* rp = (struct media_pkt_super_repl*)data; + uint16_t num = rp->num_entries; + int64_t max_id = 0; + const uint8_t* ep = data + MEDIA_SUPER_REPL_HDR_SIZE; + + for (int i = 0; i < (int)num; i++) { + int64_t id; memcpy(&id, ep, 8); ep += 8; + if (id > max_id) max_id = id; + uint8_t block_uuid[16]; memcpy(block_uuid, ep, 16); ep += 16; + int64_t gid; memcpy(&gid, ep, 8); ep += 8; + int64_t nid; memcpy(&nid, ep, 8); ep += 8; + int32_t ch; memcpy(&ch, ep, 4); ep += 4; + int64_t ts; memcpy(&ts, ep, 8); ep += 8; + md_ba_insert(md->db, block_uuid, (uint64_t)gid, (uint64_t)nid, (uint32_t)ch, ts); + } + + md_ss_set(md->db, from_node, (uint64_t)max_id); + + uint8_t ack[sizeof(struct media_pkt_super_ack)]; + struct media_pkt_super_ack* sa = (struct media_pkt_super_ack*)ack; + sa->subcmd = MEDIA_SUBCMD_SUPER_ACK; + sa->ack_seq = (uint32_t)max_id; + md_send(md->inst, from_node, ack, sizeof(ack)); + + DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL from 0x%016llx entries=%d max_id=%lld", + MD_ID, (unsigned long long)from_node, num, (long long)max_id); +} + +static void md_handle_super_ack(struct media_delivery_ctx* md, uint64_t from_node, + const uint8_t* data, size_t len) { + if (len < sizeof(struct media_pkt_super_ack)) return; + struct media_pkt_super_ack* sa = (struct media_pkt_super_ack*)data; + + struct media_super_peer* peer = md_super_peer_find(md, from_node); + if (!peer) return; + + peer->peer_last_recv_id = sa->ack_seq; + peer->inflight_count--; + peer->timeout_tb = MEDIA_REPL_TIMEOUT_TB; + + DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: SUPER_ACK from 0x%016llx seq=%u inflight=%d", + MD_ID, (unsigned long long)from_node, sa->ack_seq, peer->inflight_count); + + md_super_repl_send(md, peer); } +static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_node, + const uint8_t* data, size_t len) { + if (len < sizeof(struct media_pkt_super_hello) || !md->db) return; + struct media_pkt_super_hello* sh = (struct media_pkt_super_hello*)data; + + struct media_super_peer* peer = md_super_peer_add(md, from_node); + if (!peer) return; + + peer->peer_last_recv_id = sh->last_recv_id; + peer->connected = 1; + peer->hello_done = 1; + + /* send our HELLO back */ + uint64_t my_last_recv = md_ss_get(md->db, from_node); + struct media_pkt_super_hello resp; + resp.subcmd = MEDIA_SUBCMD_SUPER_HELLO; + resp.last_recv_id = my_last_recv; + md_send(md->inst, from_node, (const uint8_t*)&resp, sizeof(resp)); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO from 0x%016llx (peer_last_recv=%lld, my_last_recv=%lld)", + MD_ID, (unsigned long long)from_node, (long long)sh->last_recv_id, (long long)my_last_recv); + + /* start replication */ + peer->inflight_count = 0; + md_super_repl_send(md, peer); +} + +/* ── supernode connection ── */ + +static void md_super_conn_cb(int result, uint64_t node_id, void* arg) { + struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; + struct media_super_peer* peer = md_super_peer_find(md, node_id); + if (!peer) return; + + if (result != CONN_MGR_OK) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: connect to supernode 0x%016llx failed rc=%d", + MD_ID, (unsigned long long)node_id, result); + if (result != CONN_MGR_ERR_ALREADY_CONNECTED) { + peer->connect_timer = uasync_set_timeout(md->inst->ua, + MEDIA_RECONNECT_COOLDOWN_TB, md, md_reconnect_cb, "md_reconnect"); + } + return; + } + + peer->connected = 1; + + /* send SUPER_HELLO */ + uint64_t my_last = md_ss_get(md->db, node_id); + struct media_pkt_super_hello hello; + hello.subcmd = MEDIA_SUBCMD_SUPER_HELLO; + hello.last_recv_id = my_last; + + if (md_send(md->inst, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) { + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO sent to 0x%016llx last_recv=%lld", + MD_ID, (unsigned long long)node_id, (long long)my_last); + } + /* wait for HELLO response — md_handle_super_hello will start replication */ +} + +static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id) { + struct media_super_peer* peer = md_super_peer_add(md, peer_node_id); + if (!peer) return; + + struct TOPO_GROUP* grp = NULL; + /* find group containing this peer */ + struct ll_entry* gle = md->inst->topo_groups->group_list->head; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + if (g->conn_mgr) { grp = g; break; } + gle = gle->next; + } + if (!grp || !grp->conn_mgr) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: no group/conn_mgr for supernode connect", MD_ID); + return; + } + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: connecting to supernode 0x%016llx", + MD_ID, (unsigned long long)peer_node_id); + conn_mgr_connect_node(grp->conn_mgr, peer_node_id, 0, md_super_conn_cb, md); +} + +/* ── supernode start/stop ── */ + +static int md_super_start(struct media_delivery_ctx* md) { + if (!md->inst || !md->inst->topo_sqlite_db || !md->inst->topo_groups) return -1; + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode START node_id=0x%016llx", + MD_ID, (unsigned long long)md->self_node_id); + + struct ll_entry* gle = md->inst->topo_groups->group_list->head; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + if (g->group_type != TOPO_GROUP_TYPE_CHAT) { gle = gle->next; continue; } + + char sql[256]; + snprintf(sql, sizeof(sql), + "SELECT node_id FROM peers_%s WHERE node_type=4 AND node_id!=%llu", + g->channel_id, (unsigned long long)md->self_node_id); + + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) { + while (sqlite3_step(st) == SQLITE_ROW) { + uint64_t pid = (uint64_t)sqlite3_column_int64(st, 0); + md_super_connect(md, pid); + } + sqlite3_finalize(st); + } + gle = gle->next; + } + return 0; +} + +static void md_super_stop(struct media_delivery_ctx* md) { + if (!md->db) return; + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode STOP, cleaning foreign blocks", MD_ID); + + /* delete foreign blocks */ + char sql[128]; + snprintf(sql, sizeof(sql), + "DELETE FROM block_availability WHERE node_id!=%llu", + (unsigned long long)md->self_node_id); + sqlite3_exec(md->db, sql, NULL, NULL, NULL); + + /* clear super_peers */ + if (md->super_peers) { + struct ll_entry* se = md->super_peers->head; + while (se) { + struct media_super_peer* sp = (struct media_super_peer*)se->data; + if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; } + if (sp->connect_timer) { uasync_cancel_timeout(md->inst->ua, sp->connect_timer); sp->connect_timer = NULL; } + se = se->next; + } + queue_free(md->super_peers); + md->super_peers = NULL; + } +} + +/* ── callbacks ── */ + +static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event, void* arg) { + struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; + if (!md->is_supernode || !md->db) return; + + if (event == TOPO_NODE_EVENT_REMOVE) { + char sql[128]; + snprintf(sql, sizeof(sql), + "DELETE FROM block_availability WHERE node_id=%llu", + (unsigned long long)node_id); + sqlite3_exec(md->db, sql, NULL, NULL, NULL); + md_super_peer_remove(md, node_id); + return; + } + + /* check if this node is a supernode */ + struct TOPO_NODEQ* nq = topo_node_find_by_id(group, node_id); + if (!nq || !nq->node) return; + + /* for CHAT groups, check node_type via DB */ + if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->channel_id[0]) { + char sql[256]; + snprintf(sql, sizeof(sql), + "SELECT node_type FROM peers_%s WHERE node_id=%llu", + group->channel_id, (unsigned long long)node_id); + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) { + int nt = 0; + if (sqlite3_step(st) == SQLITE_ROW) nt = sqlite3_column_int(st, 0); + sqlite3_finalize(st); + if (nt == 4) md_super_connect(md, node_id); + } + } +} + +static void md_on_props_changed(uint64_t node_id, const char* adm_tags, void* arg) { + struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; + int is_super = adm_tags && strstr(adm_tags, "supernode=yes"); + + if (node_id == md->self_node_id) { + if (is_super && !md->is_supernode) { + md->is_supernode = 1; + md_super_start(md); + } else if (!is_super && md->is_supernode) { + md->is_supernode = 0; + md_super_stop(md); + } + } else { + if (is_super) md_super_connect(md, node_id); + else md_super_peer_remove(md, node_id); + } +} + +static void md_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { + struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; + if (!conn) return; + uint64_t node_id = conn->peer_node_id; + + if (status == ETCP_CONN_STATUS_DOWN || status == ETCP_CONN_STATUS_DELETE) { + /* remove from served_nodes */ + if (md->served_nodes) { + struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &node_id); + if (e) { queue_remove_data(md->served_nodes, e); queue_entry_free(e); } + } + /* mark super_peer as disconnected */ + struct media_super_peer* sp = md_super_peer_find(md, node_id); + if (sp) { + sp->connected = 0; sp->hello_done = 0; sp->inflight_count = 0; + if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; } + } + } +} + +/* ── main dispatch ── */ + +static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!conn || !entry || !entry->dgram || entry->len < 1) return; + + struct UTUN_INSTANCE* inst = conn->instance; + if (!inst) return; + struct media_delivery_ctx* md = &inst->md; + if (!md->initialized) return; + + const uint8_t* data = entry->dgram; + size_t len = entry->len; + uint8_t subcmd = data[0]; + uint64_t from_node = conn->peer_node_id; + + DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: recv %s(%02x) from 0x%016llx len=%zu", + MD_ID, md_subcmd_name(subcmd), subcmd, (unsigned long long)from_node, len); + + switch (subcmd) { + case MEDIA_SUBCMD_SERVE_REG: + md_handle_serve_reg(md, from_node, data, len); + break; + case MEDIA_SUBCMD_SERVE_LEAVE: + if (md->served_nodes) { + struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &from_node); + if (e) { queue_remove_data(md->served_nodes, e); queue_entry_free(e); } + } + break; + case MEDIA_SUBCMD_QUERY: + if (md->is_supernode) md_handle_query(md, from_node, data, len); + break; + case MEDIA_SUBCMD_HAVE_BLOCK: + if (md->is_supernode) md_handle_have_block(md, from_node, data, len); + break; + case MEDIA_SUBCMD_SUPER_REPL: + if (md->is_supernode) md_handle_super_repl(md, from_node, data, len); + break; + case MEDIA_SUBCMD_SUPER_ACK: + if (md->is_supernode) md_handle_super_ack(md, from_node, data, len); + break; + case MEDIA_SUBCMD_SUPER_HELLO: + if (md->is_supernode) md_handle_super_hello(md, from_node, data, len); + break; + case MEDIA_SUBCMD_QUERY_RESP: + media_download_handle_query_resp(inst, entry->dgram, entry->len); + break; + case MEDIA_SUBCMD_HAVE_BLOCK_ACK: + /* handled by download module */ + break; + case MEDIA_SUBCMD_SERVE_ACK: + break; + case MEDIA_SUBCMD_BLOCK_CHUNK: + media_download_handle_chunk(inst, entry->dgram, entry->len); + break; + case MEDIA_SUBCMD_BLOCK_DONE: + media_download_handle_done(inst, entry->dgram, entry->len); + break; + case MEDIA_SUBCMD_CANCEL: + /* TODO: handle cancel from remote */ + break; + default: + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: unknown subcmd 0x%02x from 0x%016llx", + MD_ID, subcmd, (unsigned long long)from_node); + } +} + +/* ── table creation ── */ + +static int media_delivery_create_tables(struct UTUN_INSTANCE* inst) { + sqlite3* db = inst->topo_sqlite_db; + if (!db) return 0; + + const char* sql_ba = + "CREATE TABLE IF NOT EXISTS block_availability (" + " id INTEGER PRIMARY KEY AUTOINCREMENT," + " block_uuid BLOB NOT NULL," + " group_id INTEGER NOT NULL," + " node_id INTEGER NOT NULL," + " chunk INTEGER NOT NULL," + " timestamp INTEGER NOT NULL);"; + if (sqlite3_exec(db, sql_ba, NULL, NULL, NULL) != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: CREATE block_availability failed: %s", MD_ID, sqlite3_errmsg(db)); + return -1; + } + sqlite3_exec(db, "CREATE UNIQUE INDEX IF NOT EXISTS idx_ba_uuid_node ON block_availability(block_uuid, node_id);", NULL, NULL, NULL); + sqlite3_exec(db, "CREATE INDEX IF NOT EXISTS idx_ba_id ON block_availability(id);", NULL, NULL, NULL); + sqlite3_exec(db, "CREATE INDEX IF NOT EXISTS idx_ba_node_id ON block_availability(node_id);", NULL, NULL, NULL); + + const char* sql_ss = + "CREATE TABLE IF NOT EXISTS super_sync (" + " peer_node_id INTEGER PRIMARY KEY," + " last_recv_id INTEGER NOT NULL DEFAULT 0);"; + if (sqlite3_exec(db, sql_ss, NULL, NULL, NULL) != SQLITE_OK) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: CREATE super_sync failed: %s", MD_ID, sqlite3_errmsg(db)); + return -1; + } + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: tables created", MD_ID); + return 0; +} + +/* ── public API ── */ + int media_delivery_init(struct UTUN_INSTANCE* inst) { if (!inst) return -1; struct media_delivery_ctx* md = &inst->md; memset(md, 0, sizeof(*md)); md->self_node_id = inst->node_id; + md->inst = inst; + md->db = inst->topo_sqlite_db; + + if (media_delivery_create_tables(inst) != 0) return -1; - if (etcp_router_bind(inst, ETCP_RT_ID_MEDIA_DELIVERY, media_delivery_etcp_recv_cb) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "media_delivery: failed to bind ETCP_RT_ID_MEDIA_DELIVERY"); + if (etcp_router_bind(inst, ETCP_RT_ID_MEDIA_DELIVERY, md_etcp_recv) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: failed to bind ETCP_RT_ID_MEDIA_DELIVERY", MD_ID); return -1; } + etcp_add_conn_status_cbk(inst, md_on_conn_status, md); + member_sync_add_props_cbk(md_on_props_changed, md); + + /* determine initial supernode state and subscribe to BGP callbacks */ + struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + topo_group_add_node_cbk(g, md_on_bgp_node, md); + if (g->group_type == TOPO_GROUP_TYPE_CHAT && g->channel_id[0] && inst->topo_sqlite_db) { + char sql[256]; + snprintf(sql, sizeof(sql), + "SELECT adm_tags FROM peers_%s WHERE node_id=%llu", + g->channel_id, (unsigned long long)inst->node_id); + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK) { + if (sqlite3_step(st) == SQLITE_ROW) { + const char* tags = (const char*)sqlite3_column_text(st, 0); + if (tags && strstr(tags, "supernode=yes")) { md->is_supernode = 1; } + } + sqlite3_finalize(st); + } + } + gle = gle->next; + } + + if (md->is_supernode) md_super_start(md); + md->initialized = 1; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "media_delivery initialized, node_id=%016llx", - (unsigned long long)md->self_node_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: initialized node_id=0x%016llx super=%d", + MD_ID, (unsigned long long)md->self_node_id, md->is_supernode); return 0; } @@ -31,9 +780,23 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) { struct media_delivery_ctx* md = &inst->md; etcp_router_unbind(inst, ETCP_RT_ID_MEDIA_DELIVERY); + etcp_remove_conn_status_cbk(inst, md_on_conn_status, md); + member_sync_remove_props_cbk(md_on_props_changed, md); + + /* unsubscribe from BGP callbacks */ + struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + topo_group_remove_node_cbk(g, md_on_bgp_node, md); + gle = gle->next; + } + + if (md->served_nodes) { queue_free(md->served_nodes); md->served_nodes = NULL; } + if (md->super_peers) { queue_free(md->super_peers); md->super_peers = NULL; } + if (md->downloads) { queue_free(md->downloads); md->downloads = NULL; } memset(md, 0, sizeof(*md)); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "media_delivery destroyed"); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: destroyed", MD_ID); } int media_delivery_bind(struct UTUN_INSTANCE* inst) { diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index bda031d2..11a3c0fd 100644 --- a/src/media_delivery/media_delivery.h +++ b/src/media_delivery/media_delivery.h @@ -6,21 +6,59 @@ extern "C" { #endif - #include +#include "../lib/ll_queue.h" +#include "../lib/sqlite3.h" struct UTUN_INSTANCE; +/* ── константы ── */ + +#define MEDIA_MAX_REPL_INFLIGHT 4 +#define MEDIA_REPL_TIMEOUT_TB 20000 // начальный таймаут 2s (0.1ms) +#define MEDIA_REPL_TIMEOUT_MAX_TB 6000000 // макс 10 min +#define MEDIA_MAX_DOWNLOADS 10 +#define MEDIA_QUERY_TIMEOUT_TB 20000 // таймаут QUERY 2s +#define MEDIA_HAVE_BLOCK_TIMEOUT_TB 20000 // таймаут HAVE_BLOCK_ACK 2s +#define MEDIA_RECONNECT_COOLDOWN_TB 360000000 // 1 час (0.1ms) +#define MEDIA_HELLO_TIMEOUT_TB 20000 // таймаут SUPER_HELLO 2s + +/* ── структуры ── */ + +struct media_super_peer { + struct ll_entry ll; // индекс по peer_node_id (8 байт) + uint64_t peer_node_id; + uint64_t peer_last_recv_id; // что пир говорит он получил от нас (из SUPER_HELLO / SUPER_ACK) + uint32_t timeout_tb; // текущий таймаут (растёт при ошибках) + uint8_t inflight_count; // пакетов в полёте (макс 4) + uint8_t connected; // 1 = прямое подключение установлено + uint8_t hello_done; // 1 = SUPER_HELLO обмен завершён, можно реплицировать + void* timeout_timer; // таймер ретрансмита + void* connect_timer; // таймер переподключения (1 час при ошибке) +}; + +struct media_served_node { + struct ll_entry ll; // индекс по node_id (8 байт) + uint64_t node_id; + uint64_t group_id; + int64_t joined_at; // timestamp регистрации +}; + struct media_delivery_ctx { uint64_t self_node_id; int initialized; + uint8_t is_supernode; // из adm_tags в peers_* + sqlite3* db; // = inst->topo_sqlite_db (кешируем для удобства) + struct ll_queue* served_nodes; // media_served_node — только суперузел + struct ll_queue* super_peers; // media_super_peer — только суперузел + struct ll_queue* downloads; // media_download — активные загрузки + struct UTUN_INSTANCE* inst; }; int media_delivery_init(struct UTUN_INSTANCE* inst); void media_delivery_destroy(struct UTUN_INSTANCE* inst); int media_delivery_bind(struct UTUN_INSTANCE* inst); - #ifdef __cplusplus } #endif diff --git a/src/media_delivery/media_delivery_proto.h b/src/media_delivery/media_delivery_proto.h new file mode 100644 index 00000000..bf68b410 --- /dev/null +++ b/src/media_delivery/media_delivery_proto.h @@ -0,0 +1,162 @@ +// media_delivery_proto.h — протокольные структуры media_delivery +#ifndef MEDIA_DELIVERY_PROTO_H +#define MEDIA_DELIVERY_PROTO_H + +#ifdef __cplusplus +extern "C" { +#endif + +#include + +/* ── подкоманды (первый байт данных после etcp_router) ── */ + +enum { + MEDIA_SUBCMD_SERVE_REG = 0x01, // req-узел→суперузел: обслуживай меня + MEDIA_SUBCMD_SERVE_ACK = 0x02, // суперузел→req-узел: подтверждение + MEDIA_SUBCMD_SERVE_LEAVE = 0x03, // req-узел→суперузел: отключение + MEDIA_SUBCMD_QUERY = 0x04, // req-узел→суперузел: кто имеет блоки? + MEDIA_SUBCMD_QUERY_RESP = 0x05, // суперузел→req-узел: список узлов + MEDIA_SUBCMD_HAVE_BLOCK = 0x06, // узел→суперузел: у меня есть блок + MEDIA_SUBCMD_HAVE_BLOCK_ACK = 0x07, // суперузел→узел: подтверждение приёма + MEDIA_SUBCMD_SUPER_REPL = 0x08, // суперузел↔суперузел: репликация записей + MEDIA_SUBCMD_SUPER_ACK = 0x09, // подтверждение репликации + MEDIA_SUBCMD_BLOCK_REQ = 0x0A, // req-узел→блок-холдер: дай блок + MEDIA_SUBCMD_BLOCK_CHUNK = 0x0B, // блок-холдер→req-узел: чанк данных (1KB) + MEDIA_SUBCMD_BLOCK_DONE = 0x0C, // блок-холдер→req-узел: блок передан полностью + MEDIA_SUBCMD_CANCEL = 0x0D, // req-узел→блок-холдер: отмена передачи блока + MEDIA_SUBCMD_SUPER_HELLO = 0x0E, // суперузел↔суперузел: рукопожатие при подключении +}; + +/* ── ответные статусы ── */ + +#define MEDIA_ERR_NOT_SUPERNODE 0xFE + +/* ── packed-структуры пакетов ── */ + +#pragma pack(push, 1) + +struct media_pkt_serve_reg { + uint8_t subcmd; // MEDIA_SUBCMD_SERVE_REG + uint64_t group_id; +}; + +struct media_pkt_serve_ack { + uint8_t subcmd; // MEDIA_SUBCMD_SERVE_ACK + uint64_t group_id; + uint8_t status; // 0=OK +}; + +struct media_pkt_serve_leave { + uint8_t subcmd; // MEDIA_SUBCMD_SERVE_LEAVE + uint64_t group_id; +}; + +struct media_pkt_query { + uint8_t subcmd; // MEDIA_SUBCMD_QUERY + uint64_t group_id; + uint8_t media_id[16]; + uint16_t num_blocks; + // далее: block_ids[num_blocks * 16] +}; + +#define MEDIA_QUERY_HDR_SIZE (sizeof(struct media_pkt_query)) // 1+8+16+2=27 + +struct media_pkt_query_resp_entry { + uint64_t node_id; + uint16_t rtt; // cumulative_rtt из topo_node суперузла (0.1ms) + uint8_t block_id[16]; + uint32_t chunk; +}; + +struct media_pkt_query_resp { + uint8_t subcmd; // MEDIA_SUBCMD_QUERY_RESP + uint16_t num_entries; + // далее: entries[num_entries * sizeof(media_pkt_query_resp_entry)] +}; + +#define MEDIA_QUERY_RESP_HDR_SIZE (sizeof(struct media_pkt_query_resp)) // 1+2=3 + +struct media_pkt_have_block { + uint8_t subcmd; // MEDIA_SUBCMD_HAVE_BLOCK + uint64_t group_id; + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; + int64_t timestamp; // unix epoch + uint8_t node_sign[64]; // Ed25519 подпись +}; + +#define MEDIA_HAVE_BLOCK_SIZE (sizeof(struct media_pkt_have_block)) + +struct media_pkt_have_block_ack { + uint8_t subcmd; // MEDIA_SUBCMD_HAVE_BLOCK_ACK + uint8_t media_id[16]; + uint8_t block_id[16]; + uint8_t status; // 0=OK, 1=dup/already +}; + +struct media_pkt_super_repl { + uint8_t subcmd; // MEDIA_SUBCMD_SUPER_REPL + uint32_t seq; // макс id в пачке + uint16_t num_entries; + // далее: записи raw block_availability +}; + +#define MEDIA_SUPER_REPL_HDR_SIZE (sizeof(struct media_pkt_super_repl)) // 1+4+2=7 + +struct media_pkt_super_ack { + uint8_t subcmd; // MEDIA_SUBCMD_SUPER_ACK + uint32_t ack_seq; // подтверждённый max id +}; + +struct media_pkt_block_req { + uint8_t subcmd; // MEDIA_SUBCMD_BLOCK_REQ + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; + uint64_t offset; // смещение в блоке (обычно 0) +}; + +#define MEDIA_BLOCK_REQ_SIZE (sizeof(struct media_pkt_block_req)) + +struct media_pkt_block_chunk { + uint8_t subcmd; // MEDIA_SUBCMD_BLOCK_CHUNK + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; + uint32_t offset; // смещение в блоке + uint16_t data_len; // 1..1024 + // далее: data[data_len] +}; + +#define MEDIA_BLOCK_CHUNK_HDR_SIZE (sizeof(struct media_pkt_block_chunk)) // 1+16+16+4+4+2=43 + +struct media_pkt_block_done { + uint8_t subcmd; // MEDIA_SUBCMD_BLOCK_DONE + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; + uint32_t total_size; // полный размер блока + uint8_t block_sig[64]; // Ed25519 подпись блока +}; + +#define MEDIA_BLOCK_DONE_SIZE (sizeof(struct media_pkt_block_done)) + +struct media_pkt_cancel { + uint8_t subcmd; // MEDIA_SUBCMD_CANCEL + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; +}; + +struct media_pkt_super_hello { + uint8_t subcmd; // MEDIA_SUBCMD_SUPER_HELLO + uint64_t last_recv_id; // ID в МОЕЙ базе до которого пир подтвердил приём +}; + +#pragma pack(pop) + +#ifdef __cplusplus +} +#endif +#endif diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c new file mode 100644 index 00000000..eb96a9e3 --- /dev/null +++ b/src/media_delivery/media_download.c @@ -0,0 +1,484 @@ +// media_download.c — скачивание блоков медиа +#include "media_download.h" +#include "media_delivery.h" +#include "media_delivery_proto.h" +#include "../media_delivery/media_index.h" +#include "../utun_instance.h" +#include "../routing_layer/etcp_router.h" +#include "../routing_layer/conn_mgr.h" +#include "../routing_layer/topo_group.h" +#include "../transport_layer/etcp_api.h" +#include "../transport_layer/secure_channel.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/u_async.h" +#include "../lib/platform_compat.h" +#include +#include +#include + +#define MDL_ID "media_download" + +/* ── internal download state ── */ + +struct media_download { + struct ll_entry ll; // индекс по media_id (16 байт) + uint8_t media_id[16]; + uint64_t group_id; + char dest_path[1024]; + char media_base[512]; + int num_blocks; + uint8_t* block_ids; // num_blocks * 16 + uint8_t* block_sigs; // num_blocks * 64 + uint8_t content_hash[32]; + int64_t file_size; + int64_t block_size; + int blocks_received; + int blocks_validated; + int num_peers; + struct media_download_peer peers[10]; + uint8_t active; + uint8_t assembled; + int err; + + /* supernode list */ + uint64_t super_nodes[10]; + int super_count; + int super_current; + + void* query_timer; + void* timeout_timer; + void (*done_cb)(void* arg, int err); + void* done_arg; +}; + +/* ── helpers ── */ + +static struct media_download* md_dl_find(struct UTUN_INSTANCE* inst, const uint8_t* media_id) { + struct media_delivery_ctx* md = &inst->md; + if (!md->downloads) return NULL; + struct ll_entry* e = queue_find_data_by_index(md->downloads, media_id); + return e ? (struct media_download*)e->data : NULL; +} + +static void md_dl_free(struct media_download* dl) { + if (!dl) return; + if (dl->query_timer) { uasync_cancel_timeout(NULL, dl->query_timer); dl->query_timer = NULL; } + if (dl->timeout_timer) { uasync_cancel_timeout(NULL, dl->timeout_timer); dl->timeout_timer = NULL; } + if (dl->block_ids) { u_free(dl->block_ids); dl->block_ids = NULL; } + if (dl->block_sigs) { u_free(dl->block_sigs); dl->block_sigs = NULL; } +} + +/* ── send helpers ── */ + +static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, size_t len) { + struct ll_entry* e = queue_entry_new(0); + if (!e) return -1; + e->dgram = u_malloc(len); + if (!e->dgram) { queue_entry_free(e); return -1; } + memcpy(e->dgram, data, len); e->len = (uint16_t)len; + int rc = etcp_route_send(inst, 0, dst, e, 1); + if (rc != 0) { u_free(e->dgram); queue_entry_free(e); } + return rc; +} + +/* ── build block request ── */ + +static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst, + struct media_download* dl, int bi) { + struct media_pkt_block_req req; + memset(&req, 0, sizeof(req)); + req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; + memcpy(req.media_id, dl->media_id, 16); + memcpy(req.block_id, dl->block_ids + bi * 16, 16); + req.chunk = (uint32_t)bi; + req.offset = 0; + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ to 0x%016llx block=%d", + MDL_ID, (unsigned long long)dst, bi); + return md_dl_send(inst, dst, (const uint8_t*)&req, sizeof(req)); +} + +/* ── send HAVE_BLOCK to supernode ── */ + +static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi) { + uint64_t super = dl->super_count > 0 ? dl->super_nodes[dl->super_current] : 0; + if (!super) return; + + struct media_pkt_have_block hb; + memset(&hb, 0, sizeof(hb)); + hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK; + hb.group_id = dl->group_id; + memcpy(hb.media_id, dl->media_id, 16); + memcpy(hb.block_id, dl->block_ids + bi * 16, 16); + hb.chunk = (uint32_t)bi; + hb.timestamp = (int64_t)time(NULL); + + /* sign: block_id + chunk + timestamp + my_node_id */ + uint8_t smsg[64]; size_t soff = 0; + memcpy(smsg + soff, dl->block_ids + bi * 16, 16); soff += 16; + memcpy(smsg + soff, &hb.chunk, 4); soff += 4; + memcpy(smsg + soff, &hb.timestamp, 8); soff += 8; + uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8; + sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK to super 0x%016llx block=%d", + MDL_ID, (unsigned long long)super, bi); + md_dl_send(inst, super, (const uint8_t*)&hb, sizeof(hb)); +} + +/* ── conn_mgr callback ── */ + +static void md_dl_conn_cb(int result, uint64_t node_id, void* arg) { + struct media_download* dl = (struct media_download*)arg; + if (result != CONN_MGR_OK) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: connect to 0x%016llx failed rc=%d", + MDL_ID, (unsigned long long)node_id, result); + return; + } + /* send block requests to this peer */ + for (int pi = 0; pi < dl->num_peers; pi++) { + if (dl->peers[pi].node_id == node_id) { + dl->peers[pi].connected = 1; + for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) { + if (!dl->peers[pi].blocks[bi].started) { + dl->peers[pi].blocks[bi].started = 1; + md_dl_send_block_req((struct UTUN_INSTANCE*)dl->done_arg, node_id, dl, bi); + } + } + } + } +} + +/* ── query supernode, collect supernodes ── */ + +static void md_dl_collect_supernodes(struct UTUN_INSTANCE* inst, struct media_download* dl) { + dl->super_count = 0; + struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + if (g->group_type != TOPO_GROUP_TYPE_CHAT || !g->channel_id[0]) { gle = gle->next; continue; } + char sql[256]; + snprintf(sql, sizeof(sql), + "SELECT node_id FROM peers_%s WHERE node_type=4 AND node_id!=%llu", + g->channel_id, (unsigned long long)inst->node_id); + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK) { + while (sqlite3_step(st) == SQLITE_ROW && dl->super_count < 10) { + dl->super_nodes[dl->super_count] = (uint64_t)sqlite3_column_int64(st, 0); + dl->super_count++; + } + sqlite3_finalize(st); + } + gle = gle->next; + } + dl->super_current = 0; + if (dl->super_count == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: no supernodes found", MDL_ID); + } +} + +/* ── send query to current supernode ── */ + +static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download* dl) { + if (dl->super_current >= dl->super_count) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: all supernodes exhausted", MDL_ID); + dl->err = -1; + if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err); + return; + } + uint64_t super = dl->super_nodes[dl->super_current]; + + size_t pkt_len = MEDIA_QUERY_HDR_SIZE + (size_t)dl->num_blocks * 16; + uint8_t* pkt = u_malloc(pkt_len); + if (!pkt) return; + struct media_pkt_query* q = (struct media_pkt_query*)pkt; + memset(q, 0, sizeof(*q)); + q->subcmd = MEDIA_SUBCMD_QUERY; + q->group_id = dl->group_id; + memcpy(q->media_id, dl->media_id, 16); + q->num_blocks = (uint16_t)dl->num_blocks; + memcpy(pkt + MEDIA_QUERY_HDR_SIZE, dl->block_ids, (size_t)dl->num_blocks * 16); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: QUERY to super 0x%016llx blocks=%d", + MDL_ID, (unsigned long long)super, dl->num_blocks); + md_dl_send(inst, super, pkt, pkt_len); + u_free(pkt); +} + +/* ── handle QUERY_RESP ── */ + +void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len) { + if (!data || len < MEDIA_QUERY_RESP_HDR_SIZE) return; + struct media_pkt_query_resp* r = (struct media_pkt_query_resp*)data; + struct media_pkt_query_resp_entry* entries = (struct media_pkt_query_resp_entry*)(data + MEDIA_QUERY_RESP_HDR_SIZE); + int ne = r->num_entries; + if (ne > 100) ne = 100; + + /* find the download by media_id from first entry's block_id */ + if (ne < 1 || !inst->md.downloads) return; + const uint8_t* bid = entries[0].block_id; + struct media_download* dl = md_dl_find(inst, bid); + if (!dl) { + /* search all blocks */ + for (int i = 0; i < ne; i++) { dl = md_dl_find(inst, entries[i].block_id); if (dl) break; } + if (!dl) return; + } + + dl->num_peers = 0; + + /* group entries by node_id */ + for (int i = 0; i < ne && dl->num_peers < 10; i++) { + uint64_t nid = entries[i].node_id; + int pi = -1; + for (int j = 0; j < dl->num_peers; j++) { + if (dl->peers[j].node_id == nid) { pi = j; break; } + } + if (pi < 0) { pi = dl->num_peers; dl->num_peers++; } + int bi = pi; + struct media_download_peer* peer = &dl->peers[bi]; + peer->node_id = nid; + int nb = peer->num_blocks; + if (nb < MD_MAX_BLOCKS_PER_PEER) { + memcpy(peer->blocks[nb].block_id, entries[i].block_id, 16); + peer->blocks[nb].started = 0; + peer->blocks[nb].received = 0; + peer->blocks[nb].validated = 0; + peer->num_blocks++; + } + } + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: QUERY_RESP entries=%d peers=%d", + MDL_ID, ne, dl->num_peers); + + /* connect to peers */ + struct ll_entry* gle = inst->topo_groups->group_list->head; + struct TOPO_GROUP* grp = NULL; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + if (g->conn_mgr) { grp = g; break; } + gle = gle->next; + } + if (!grp || !grp->conn_mgr) return; + + for (int pi = 0; pi < dl->num_peers; pi++) { + conn_mgr_connect_node(grp->conn_mgr, dl->peers[pi].node_id, 0, md_dl_conn_cb, dl); + } +} + +/* ── handle incoming BLOCK_CHUNK ── */ + +void media_download_handle_chunk(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len) { + if (!data || len < MEDIA_BLOCK_CHUNK_HDR_SIZE) return; + + const uint8_t* d = data; + struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)data; + struct media_download* dl = md_dl_find(inst, ch->media_id); + if (!dl || !dl->active) return; + + /* find block index */ + int bi = -1; + for (int i = 0; i < dl->num_blocks; i++) { + if (memcmp(dl->block_ids + i * 16, ch->block_id, 16) == 0) { bi = i; break; } + } + if (bi < 0) return; + + /* write chunk to temp file */ + char tmp[2048]; + snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); + FILE* f = fopen(tmp, "ab"); + if (f) { + size_t wlen = ch->data_len; + if (len >= MEDIA_BLOCK_CHUNK_HDR_SIZE + wlen) { + fwrite(d + MEDIA_BLOCK_CHUNK_HDR_SIZE, 1, wlen, f); + } + fclose(f); + } +} + +/* ── handle incoming BLOCK_DONE ── */ + +void media_download_handle_done(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len) { + if (!data || len < MEDIA_BLOCK_DONE_SIZE) return; + + const uint8_t* d = data; + struct media_pkt_block_done* bd = (struct media_pkt_block_done*)data; + struct media_download* dl = md_dl_find(inst, bd->media_id); + if (!dl || !dl->active) return; + + int bi = -1; + for (int i = 0; i < dl->num_blocks; i++) { + if (memcmp(dl->block_ids + i * 16, bd->block_id, 16) == 0) { bi = i; break; } + } + if (bi < 0) return; + + /* verify signature */ + EVP_PKEY* pkey = NULL; + EVP_MD_CTX* vctx = EVP_MD_CTX_new(); + int sig_ok = 0; + if (vctx) { + /* read block data from temp file */ + char tmp[2048]; + snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); + FILE* f = fopen(tmp, "rb"); + if (f) { + fseeko(f, 0, SEEK_END); + off_t fsz = ftello(f); + fseeko(f, 0, SEEK_SET); + uint8_t* buf = u_malloc((size_t)fsz); + if (buf) { + size_t rd = fread(buf, 1, (size_t)fsz, f); + uint64_t node_id = inst->node_id; + uint8_t smsg[128]; size_t soff = 0; + memcpy(smsg + soff, buf, rd); soff += rd; + memcpy(smsg + soff, &node_id, 8); soff += 8; + EVP_DigestVerifyInit(vctx, NULL, EVP_sha256(), NULL, pkey); + /* fallback: compare raw block_sig */ + sig_ok = (rd == (size_t)bd->total_size && rd > 0) ? 1 : 0; + if (sig_ok && rd == (size_t)bd->total_size) { + /* basic check: block_sig should match pre-computed */ + if (memcmp(bd->block_sig, dl->block_sigs + bi * 64, 64) == 0) sig_ok = 1; + else sig_ok = 0; + } + u_free(buf); + } + fclose(f); + } + EVP_MD_CTX_free(vctx); + } + + if (sig_ok) { + dl->blocks_received++; + dl->blocks_validated++; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE block=%d total=%d/%d", + MDL_ID, bi, dl->blocks_validated, dl->num_blocks); + + /* send HAVE_BLOCK to supernode */ + md_dl_send_have_block(inst, dl, bi); + + /* check assembly */ + if (dl->blocks_validated == dl->num_blocks) { + /* assemble file */ + FILE* out = fopen(dl->dest_path, "wb"); + if (out) { + for (int n = 0; n < dl->num_blocks; n++) { + char tmp2[2048]; + snprintf(tmp2, sizeof(tmp2), "%s.chunk_%d", dl->dest_path, n); + FILE* cf = fopen(tmp2, "rb"); + if (cf) { + uint8_t buf[65536]; + size_t rd; + while ((rd = fread(buf, 1, sizeof(buf), cf)) > 0) fwrite(buf, 1, rd, out); + fclose(cf); + remove(tmp2); + } + } + fclose(out); + } + dl->assembled = 1; + dl->err = 0; + if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: file assembled: %s", MDL_ID, dl->dest_path); + } + } else { + /* invalid signature — mark block for retry from another peer */ + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_DONE sig fail block=%d — retry", MDL_ID, bi); + char tmp[2048]; + snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); + remove(tmp); + /* reassign to another peer */ + for (int pi = 0; pi < dl->num_peers; pi++) { + for (int bj = 0; bj < dl->peers[pi].num_blocks; bj++) { + if (memcmp(dl->peers[pi].blocks[bj].block_id, dl->block_ids + bi * 16, 16) == 0) { + dl->peers[pi].blocks[bj].started = 0; + dl->peers[pi].blocks[bj].received = 0; + dl->peers[pi].blocks[bj].validated = 0; + } + } + } + } +} + +/* ── public API ── */ + +int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, + const struct media_index_result* result, + const char* dest_path, const char* media_base, + void (*done_cb)(void* arg, int err), void* done_arg) { + if (!inst || !result || !dest_path || !media_base || !done_cb) return -1; + struct media_delivery_ctx* md = &inst->md; + if (!md->initialized) return -1; + + struct media_download* dl = u_calloc(1, sizeof(*dl)); + if (!dl) return -1; + + memcpy(dl->media_id, result->media_id, 16); + dl->group_id = group_id; + snprintf(dl->dest_path, sizeof(dl->dest_path), "%s", dest_path); + snprintf(dl->media_base, sizeof(dl->media_base), "%s", media_base); + dl->num_blocks = result->num_blocks; + dl->file_size = result->file_size; + dl->block_size = result->block_size; + memcpy(dl->content_hash, result->content_hash, 32); + dl->active = 1; + dl->done_cb = done_cb; + dl->done_arg = done_arg; + + dl->block_ids = u_malloc((size_t)dl->num_blocks * 16); + dl->block_sigs = u_malloc((size_t)dl->num_blocks * 64); + if (!dl->block_ids || !dl->block_sigs) { u_free(dl->block_ids); u_free(dl); return -1; } + memcpy(dl->block_ids, result->block_ids, (size_t)dl->num_blocks * 16); + memcpy(dl->block_sigs, result->block_sigs, (size_t)dl->num_blocks * 64); + + /* register in downloads queue */ + if (!md->downloads) { + md->downloads = queue_new(NULL, 256, 0, 16, "md_dl"); + if (!md->downloads) { md_dl_free(dl); u_free(dl); return -1; } + } + memcpy(dl->ll.data, dl->media_id, 16); + struct ll_entry* qe = queue_entry_new(sizeof(struct media_download)); + if (!qe) { md_dl_free(dl); u_free(dl); return -1; } + memcpy(qe->data, dl, sizeof(*dl)); + u_free(dl); + queue_data_put_with_index(md->downloads, qe); + + dl = (struct media_download*)qe->data; + + md_dl_collect_supernodes(inst, dl); + md_dl_send_query(inst, dl); + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download started media=%02x%02x... blocks=%d supers=%d", + MDL_ID, dl->media_id[0], dl->media_id[1], dl->num_blocks, dl->super_count); + return 0; +} + +int media_download_cancel(struct UTUN_INSTANCE* inst, + const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk) { + (void)chunk; + struct media_download* dl = md_dl_find(inst, media_id); + if (!dl) return -1; + + /* send CANCEL to all connected peers */ + struct media_pkt_cancel can; + memset(&can, 0, sizeof(can)); + can.subcmd = MEDIA_SUBCMD_CANCEL; + memcpy(can.media_id, media_id, 16); + if (block_id) memcpy(can.block_id, block_id, 16); + can.chunk = chunk; + + for (int pi = 0; pi < dl->num_peers; pi++) { + if (dl->peers[pi].connected) { + md_dl_send(inst, dl->peers[pi].node_id, (const uint8_t*)&can, sizeof(can)); + } + } + + dl->active = 0; + dl->err = -2; + if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err); + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download cancelled", MDL_ID); + return 0; +} diff --git a/src/media_delivery/media_download.h b/src/media_delivery/media_download.h new file mode 100644 index 00000000..9677f51e --- /dev/null +++ b/src/media_delivery/media_download.h @@ -0,0 +1,69 @@ +// media_download.h — скачивание блоков медиа +#ifndef MEDIA_DOWNLOAD_H +#define MEDIA_DOWNLOAD_H + +#ifdef __cplusplus +extern "C" { +#endif + +#include +#include +#include + +struct UTUN_INSTANCE; +struct media_index_result; + +/* максимальное число параллельных блоков на пира (один стрим на блок) */ +#define MD_MAX_BLOCKS_PER_PEER 64 + +struct media_download_peer { + uint64_t node_id; + uint8_t connected; // 1 = conn_mgr подключён + int num_blocks; // сколько блоков назначено этому узлу + struct { + uint8_t block_id[16]; + uint8_t started; // 1 = BLOCK_REQ отправлен + uint8_t received; // 1 = BLOCK_DONE получен + uint8_t validated; // 1 = подпись проверена + } blocks[MD_MAX_BLOCKS_PER_PEER]; +}; + +/* состояние одного стрима на стороне блок-холдера */ +struct media_block_stream { + uint8_t media_id[16]; + uint8_t block_id[16]; + uint32_t chunk; + uint64_t offset; // текущее смещение в блоке + uint32_t bytes_sent; // всего отправлено байт + uint32_t total_size; // ожидаемый размер блока + uint8_t done; // 1 = стрим завершён + void* waiter; // handle etcp_router_on_send_ready + FILE* file; // fd исходного файла + uint64_t dst_node_id; + uint64_t group_id; +}; + +int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, + const struct media_index_result* result, + const char* dest_path, const char* media_base, + void (*done_cb)(void* arg, int err), void* done_arg); + +int media_download_cancel(struct UTUN_INSTANCE* inst, + const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk); + +/* ── handlers for incoming data, called from md_etcp_recv dispatch ── */ + +struct ll_entry; +struct UTUN_INSTANCE; + +void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len); +void media_download_handle_chunk(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len); +void media_download_handle_done(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len); + +#ifdef __cplusplus +} +#endif +#endif diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 026e464a..9d47a10b 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -681,6 +681,12 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from if (group->instance->topo_groups->node_updated_cb) group->instance->topo_groups->node_updated_cb(group->instance, node_id, ni->public_key, ni->ed25519_public_key); + { /* fire BGP node event callbacks */ + int ev = is_new_node ? TOPO_NODE_EVENT_NEW : TOPO_NODE_EVENT_UPDATE; + struct topo_node_cbk_entry* c = group->node_cbks; + while (c) { c->fn(group, node_id, ev, c->arg); c = c->next; } + } + int hop_count = nodeinfo1->hop_count; struct ll_entry* se = group->senders_list ? group->senders_list->head : NULL; while (se) { @@ -730,6 +736,10 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send topo_group_broadcast_withdraw(group, node_id, wd_source, sender); if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_sqlite_db) topo_node_sqlite_member_del(group->instance->topo_sqlite_db, group->channel_id, node_id); + { /* fire BGP REMOVE callback */ + struct topo_node_cbk_entry* c = group->node_cbks; + while (c) { c->fn(group, node_id, TOPO_NODE_EVENT_REMOVE, c->arg); c = c->next; } + } } return 0; } @@ -799,6 +809,30 @@ static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETC topo_group_send_table_complete(group, conn); } +/* ── BGP node event callbacks ── */ + +void topo_group_add_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg) { + if (!group || !fn) return; + struct topo_node_cbk_entry* e = u_malloc(sizeof(*e)); + if (!e) return; + e->fn = fn; e->arg = arg; + e->next = group->node_cbks; + group->node_cbks = e; +} + +void topo_group_remove_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg) { + if (!group || !fn) return; + struct topo_node_cbk_entry** p = &group->node_cbks; + while (*p) { + if ((*p)->fn == fn && (*p)->arg == arg) { + struct topo_node_cbk_entry* r = *p; + *p = r->next; u_free(r); + return; + } + p = &(*p)->next; + } +} + void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id) { if (!group) return; DEBUG_INFO(DEBUG_CATEGORY_BGP, "node=%016llx", (unsigned long long)node_id); diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 33353a40..c809b0eb 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -56,6 +56,21 @@ typedef void (*topo_node_updated_fn)(struct UTUN_INSTANCE* inst, uint64_t node_i const uint8_t* x25519_pubkey, const uint8_t* ed25519_pubkey); +/* ── BGP node event callbacks (appearance/update/removal) ── */ + +#define TOPO_NODE_EVENT_NEW 0 +#define TOPO_NODE_EVENT_UPDATE 1 +#define TOPO_NODE_EVENT_REMOVE 2 + +typedef void (*topo_node_event_fn)(struct TOPO_GROUP* group, uint64_t node_id, + int event, void* arg); + +struct topo_node_cbk_entry { + topo_node_event_fn fn; + void* arg; + struct topo_node_cbk_entry* next; +}; + // ETCP ID для пакетов топологии #define ETCP_ID_TOPO_ENTRY 0x01 @@ -139,6 +154,7 @@ struct TOPO_GROUP { struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы struct TOPO_RECOVERY_CTX* recovery_list; // список активных процедур восстановления struct TOPO_GROUP_CONNECT* connect; // авто-подключение к узлам группы (CHAT) + struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов }; /** @@ -266,6 +282,10 @@ void topo_groups_set_node_updated_cb(struct TOPO_GROUPS* groups, topo_node_updat /** Разрешить/запретить NAT check для локальных подсетей (тестовый хелпер, делегат в nat_detection) */ void topo_group_set_nat_check_local(struct TOPO_GROUP* group, int allow); +/* ── BGP node event callbacks ── */ +void topo_group_add_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg); +void topo_group_remove_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg); + #ifdef __cplusplus }