#include "media_delivery.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 #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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert prepare failed: %s", MD_ID, sqlite3_errmsg(db)); 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); if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert step failed rc=%d: %s", MD_ID, rc, sqlite3_errmsg(db)); return -1; } return 0; } 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_find_nodes prepare failed: %s", MD_ID, sqlite3_errmsg(db)); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_get_since prepare failed: %s", MD_ID, sqlite3_errmsg(db)); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ss_get prepare failed: %s", MD_ID, sqlite3_errmsg(db)); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ss_set prepare failed: %s", MD_ID, sqlite3_errmsg(db)); 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(md->inst->ua, 64, offsetof(struct media_super_peer, peer_node_id), 8, "md_super"); if (!md->super_peers) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(super_peers) failed", MD_ID); return NULL; } } struct ll_entry* qe = queue_entry_new(sizeof(struct media_super_peer)); if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(super_peer) failed", MD_ID); 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 + 1); if (!e->dgram) { queue_entry_free(e); return -1; } e->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY; memcpy(e->dgram + 1, data, len); e->len = (uint16_t)(len + 1); int rc = etcp_route_send(inst, TOPO_GROUP_UTUN, 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(md->inst->ua, 64, offsetof(struct media_served_node, node_id), 8, "md_served"); if (!md->served_nodes) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(served_nodes) failed", MD_ID); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(served_node) failed", MD_ID); 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) return; if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: QUERY but no DB", MD_ID); 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) return; if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: HAVE_BLOCK but no DB", MD_ID); 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) return; if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_REPL but no DB", MD_ID); 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) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_ACK from unknown peer 0x%016llx", MD_ID, (unsigned long long)from_node); 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)) return; if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_HELLO but no DB", MD_ID); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)from_node); 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) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: conn_cb for unknown peer 0x%016llx", MD_ID, (unsigned long long)node_id); 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); } else { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_HELLO send failed to 0x%016llx", MD_ID, (unsigned long long)node_id); } /* wait for HELLO response — md_handle_super_hello will start replication */ } /* ── handle incoming BLOCK_REQ (admission control + streaming) ── */ #define MD_MAX_STREAMS 3 #define MD_CHUNK_SIZE 1024 struct stream_ctx { struct media_delivery_ctx* md; uint64_t dst_node_id; uint8_t media_id[16]; uint8_t block_id[16]; uint32_t chunk; FILE* file; uint64_t offset; uint64_t remaining; uint64_t block_start; size_t block_data_len; struct queue_waiter_handle waiter; }; static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { (void)q; struct stream_ctx* sc = (struct stream_ctx*)arg; struct media_delivery_ctx* md = sc->md; if (!md || sc->remaining == 0) { u_free(sc); return; } /* read up to MD_CHUNK_SIZE bytes */ size_t to_read = sc->remaining < MD_CHUNK_SIZE ? (size_t)sc->remaining : MD_CHUNK_SIZE; uint8_t buf[MD_CHUNK_SIZE]; fseeko(sc->file, (off_t)sc->offset, SEEK_SET); size_t rd = fread(buf, 1, to_read, sc->file); if (rd == 0) { /* EOF or error */ u_free(sc); return; } /* build CHUNK packet */ uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + MD_CHUNK_SIZE]; struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt; memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE); ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK; memcpy(ch->media_id, sc->media_id, 16); memcpy(ch->block_id, sc->block_id, 16); ch->chunk = sc->chunk; ch->offset = (uint32_t)(sc->offset - sc->block_start); ch->data_len = (uint16_t)rd; memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd); md_send(md->inst, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); sc->offset += rd; sc->remaining -= rd; if (sc->remaining == 0) { /* all data sent — send BLOCK_DONE */ uint8_t done_pkt[MEDIA_BLOCK_DONE_SIZE]; struct media_pkt_block_done* bd = (struct media_pkt_block_done*)done_pkt; memset(bd, 0, sizeof(*bd)); bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE; memcpy(bd->media_id, sc->media_id, 16); memcpy(bd->block_id, sc->block_id, 16); bd->chunk = sc->chunk; bd->total_size = (uint32_t)sc->block_data_len; /* block_sig is zero — receiver will verify vs block_sigs from media_index_result */ md_send(md->inst, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE); media_delivery_stream_done(md->inst); fclose(sc->file); u_free(sc); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE sent chunk=%d total=%u", MD_ID, sc->chunk, bd->total_size); } else { /* continue streaming — register waiter for backpressure */ etcp_router_on_send_ready(md->inst, TOPO_GROUP_UTUN, sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY, &sc->waiter, stream_send_chunk_cb, sc); } } static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_node, const uint8_t* data, size_t len) { if (len < MEDIA_BLOCK_REQ_SIZE) return; struct media_pkt_block_req* req = (struct media_pkt_block_req*)data; DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ from 0x%016llx block=%02x%02x... active_streams=%d", MD_ID, (unsigned long long)from_node, req->block_id[0], req->block_id[1], md->active_streams); int limit_reached = md->active_streams >= MD_MAX_STREAMS; md->active_streams++; if (limit_reached) { md->active_streams--; /* undo */ struct media_pkt_block_overloaded ov; ov.subcmd = MEDIA_SUBCMD_BLOCK_OVERLOADED; memcpy(ov.media_id, req->media_id, 16); memcpy(ov.block_id, req->block_id, 16); ov.chunk = req->chunk; ov.retry_after_ms = 2000; md_send(md->inst, from_node, (const uint8_t*)&ov, sizeof(ov)); DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ OVERLOADED (streams=%d) to 0x%016llx", MD_ID, md->active_streams, (unsigned long long)from_node); return; } /* find file in media_files DB */ sqlite3* db = md->db; if (!db) return; const char* sql = "SELECT location,chunk_size,chunk,offset,file_size FROM media_files WHERE block_id=? AND node_id=? LIMIT 1"; sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return; sqlite3_bind_blob(st, 1, req->block_id, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)md->inst->node_id); if (sqlite3_step(st) != SQLITE_ROW) { sqlite3_finalize(st); return; } const char* loc = (const char*)sqlite3_column_text(st, 0); int64_t chunk_size = sqlite3_column_int64(st, 1); int chunk = sqlite3_column_int(st, 2); int64_t offset = sqlite3_column_int64(st, 3); int64_t file_size = sqlite3_column_int64(st, 4); /* compute chunk start and size */ int64_t block_start = (int64_t)chunk * chunk_size; int64_t block_end = block_start + chunk_size; if (block_end > file_size) block_end = file_size; int64_t block_length = block_end - block_start; /* find media_base */ char media_base[512] = "/tmp/utun_media"; { sqlite3_stmt* ms = NULL; sqlite3_prepare_v2(db, "SELECT value FROM ui_state WHERE key='media_base'", -1, &ms, NULL); if (ms && sqlite3_step(ms) == SQLITE_ROW) { const char* mb = (const char*)sqlite3_column_text(ms, 0); if (mb) snprintf(media_base, sizeof(media_base), "%s", mb); } if (ms) sqlite3_finalize(ms); } (void)media_base; char path[2048]; snprintf(path, sizeof(path), "%s/%s", media_base, loc); FILE* f = fopen(path, "rb"); sqlite3_finalize(st); if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: cannot open file %s for streaming", MD_ID, path); md->active_streams--; /* undo pre-increment */ return; } struct stream_ctx* sc = u_calloc(1, sizeof(*sc)); if (!sc) { fclose(f); md->active_streams--; return; } sc->md = md; sc->dst_node_id = from_node; memcpy(sc->media_id, req->media_id, 16); memcpy(sc->block_id, req->block_id, 16); sc->chunk = req->chunk; sc->file = f; sc->offset = (uint64_t)block_start; sc->remaining = (uint64_t)block_length; sc->block_start = (uint64_t)block_start; sc->block_data_len = (size_t)block_length; DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: streaming file=%s chunk=%d start=%lld len=%lld to 0x%016llx", MD_ID, path, chunk, (long long)block_start, (long long)block_length, (unsigned long long)from_node); /* start streaming — first chunk immediately, then backpressure */ memset(&sc->waiter, 0, sizeof(sc->waiter)); stream_send_chunk_cb(NULL, sc); } /* ── decrement stream counter (called when stream finished/cancelled internally) ── */ void media_delivery_stream_done(struct UTUN_INSTANCE* inst) { if (inst) { struct media_delivery_ctx* md = &inst->md; if (md->active_streams > 0) md->active_streams--; } } 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)peer_node_id); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_start missing inst/db/groups", MD_ID); 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) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: super_stop but no DB", MD_ID); 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) { media_delivery_set_supernode(md->inst, is_super); } 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) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: recv before init", MD_ID); return; } const uint8_t* data = entry->dgram; size_t len = entry->len; /* router_deliver prepends svc_id byte → data[0]=svc_id, data[1]=наш subcmd. conn_mgr uses two-tier (cmd+subcmd) hence +2; we have single-tier so +1. */ if (len < 2) return; uint8_t subcmd = data[1]; data += 1; len -= 1; 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: DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SERVE_LEAVE from 0x%016llx", MD_ID, (unsigned long long)from_node); 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); else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: QUERY ignored — not supernode", MD_ID); break; case MEDIA_SUBCMD_HAVE_BLOCK: if (md->is_supernode) md_handle_have_block(md, from_node, data, len); else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK ignored — not supernode", MD_ID); break; case MEDIA_SUBCMD_SUPER_REPL: if (md->is_supernode) md_handle_super_repl(md, from_node, data, len); else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL ignored — not supernode", MD_ID); break; case MEDIA_SUBCMD_SUPER_ACK: if (md->is_supernode) md_handle_super_ack(md, from_node, data, len); else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_ACK ignored — not supernode", MD_ID); break; case MEDIA_SUBCMD_SUPER_HELLO: if (md->is_supernode) md_handle_super_hello(md, from_node, data, len); else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO ignored — not supernode", MD_ID); 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_REQ: md_handle_block_req(md, from_node, data, len); break; case MEDIA_SUBCMD_BLOCK_OVERLOADED: media_download_handle_overloaded(inst, data, len); 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) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: init with NULL inst", MD_ID); 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, 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, "%s: initialized node_id=0x%016llx super=%d", MD_ID, (unsigned long long)md->self_node_id, md->is_supernode); return 0; } void media_delivery_destroy(struct UTUN_INSTANCE* inst) { if (!inst || !inst->md.initialized) return; 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, "%s: destroyed", MD_ID); } int media_delivery_bind(struct UTUN_INSTANCE* inst) { return media_delivery_init(inst); } void media_delivery_set_supernode(struct UTUN_INSTANCE* inst, int is_supernode) { if (!inst) return; struct media_delivery_ctx* md = &inst->md; if (is_supernode && !md->is_supernode) { md->is_supernode = 1; md_super_start(md); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode ENABLED", MD_ID); } else if (!is_supernode && md->is_supernode) { md->is_supernode = 0; md_super_stop(md); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode DISABLED", MD_ID); } }