Browse Source

media_delivery: implement full media delivery subsystem (stages 1-8)

- topo_group: add BGP node event callbacks (NEW/UPDATE/REMOVE)
- member_sync: add multi-subscriber node_props_changed callback (adm_tags)
- media_delivery: SQLite tables (block_availability + super_sync), full protocol handler
- media_delivery_proto.h: packed protocol structures (15 subcommands)
- media_download: download management, chunk streaming, block assembly, HAVE_BLOCK with retry
- Makefile.am: add media_download.c to builds
topo_upd
evgeny 2 months ago
parent
commit
fb3e34080d
  1. 2
      src/Makefile.am
  2. 34
      src/chat/member_sync.c
  3. 9
      src/chat/member_sync.h
  4. 781
      src/media_delivery/media_delivery.c
  5. 42
      src/media_delivery/media_delivery.h
  6. 162
      src/media_delivery/media_delivery_proto.h
  7. 484
      src/media_delivery/media_download.c
  8. 69
      src/media_delivery/media_download.h
  9. 34
      src/routing_layer/topo_group.c
  10. 20
      src/routing_layer/topo_group.h

2
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 \

34
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)

9
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

781
src/media_delivery/media_delivery.c

@ -1,12 +1,730 @@
#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 <string.h>
#include <time.h>
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;
@ -14,15 +732,46 @@ int media_delivery_init(struct UTUN_INSTANCE* inst) {
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) {

42
src/media_delivery/media_delivery.h

@ -6,21 +6,59 @@
extern "C" {
#endif
#include <stdint.h>
#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

162
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 <stdint.h>
/* ── подкоманды (первый байт данных после 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

484
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 <string.h>
#include <stdio.h>
#include <openssl/evp.h>
#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;
}

69
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 <stdint.h>
#include <stddef.h>
#include <stdio.h>
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

34
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);

20
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
}

Loading…
Cancel
Save