Browse Source

media_delivery: маршрутизация всегда через CHAT-группу

media_download/медиа-ответы и суперузловая репликация теперь берут группу
из контекста (dl->group_id, group_id из пакета, peer->group_id) вместо
хардкода TOPO_GROUP_UTUN. node_props_changed получает channel_id.
SUPER_HELLO несёт group_id.
v2
evgeny 4 weeks ago
parent
commit
326dfd7e2c
  1. 47
      src/chat/chat_member.c
  2. 3
      src/chat/chat_msg.c
  3. 8
      src/chat/member_sync.c
  4. 3
      src/chat/member_sync.h
  5. 60
      src/media_delivery/media_delivery.c
  6. 1
      src/media_delivery/media_delivery.h
  7. 1
      src/media_delivery/media_delivery_proto.h
  8. 31
      src/media_delivery/media_download.c
  9. 2
      tests/test_media_delivery_full.c
  10. 5
      tests/test_media_delivery_sql.c

47
src/chat/chat_member.c

@ -322,37 +322,28 @@ void chat_core_request_member_rtt_trampoline(void* arg) {
/* ─── node_props_changed callback (member_sync → chat_event) ─── */
static void on_adm_tags_changed(uint64_t node_id, const char* adm_tags, void* arg) {
static void on_adm_tags_changed(uint64_t node_id, const char* adm_tags, const char* channel_id, void* arg) {
(void)adm_tags; (void)arg;
if (!g_cc.initialized || !g_cc.inst || !g_cc.db) return;
struct chat_core_ctx* cc = &g_cc;
for (int i = 0; i < cc->si_count; i++) {
const char* ch_id = cc->si_ch_id[i];
if (!ch_id) continue;
char peers_tbl[80]; peers_table_name(ch_id, peers_tbl, sizeof(peers_tbl));
char sql[256]; snprintf(sql, sizeof(sql),
"SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(cc->db, sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
if (sqlite3_step(st) == SQLITE_ROW) {
sqlite3_finalize(st);
size_t cl = strlen(ch_id);
uint8_t evt[1 + 64 + 1 + CHAT_MEMBER_DISPLAY_SIZE];
evt[0] = (uint8_t)cl;
memcpy(evt + 1, ch_id, cl);
evt[1 + cl] = 1;
if (chat_core_get_single_member(ch_id, node_id,
evt + 1 + cl + 1) == 0)
chat_event_post(CHAT_EVT_MEMBER_UPDATED, evt,
1 + (int)cl + 1 + CHAT_MEMBER_DISPLAY_SIZE);
return;
}
if (!g_cc.initialized || !g_cc.inst || !g_cc.db || !channel_id || !channel_id[0]) return;
char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl));
char sql[256]; snprintf(sql, sizeof(sql), "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
if (sqlite3_step(st) == SQLITE_ROW) {
sqlite3_finalize(st);
size_t cl = strlen(channel_id);
uint8_t evt[1 + 64 + 1 + CHAT_MEMBER_DISPLAY_SIZE];
evt[0] = (uint8_t)cl;
memcpy(evt + 1, channel_id, cl);
evt[1 + cl] = 1;
if (chat_core_get_single_member(channel_id, node_id, evt + 1 + cl + 1) == 0)
chat_event_post(CHAT_EVT_MEMBER_UPDATED, evt, 1 + (int)cl + 1 + CHAT_MEMBER_DISPLAY_SIZE);
return;
}
sqlite3_finalize(st);
}
}

3
src/chat/chat_msg.c

@ -633,7 +633,8 @@ static int md_start_download(struct UTUN_INSTANCE* inst,
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download start ch=%s file=%s blocks=%d size=%lld dest=%s",
CC_ID, ch_id, base_filename, nb, (long long)fsize, dest);
media_download_start(inst, 0, &result, dest, media_base, author_node_id, md_download_done_cb, ctx, md_download_progress_cb, ctx);
uint64_t gid = strtoull(ch_id, NULL, 10);
media_download_start(inst, gid, &result, dest, media_base, author_node_id, md_download_done_cb, ctx, md_download_progress_cb, ctx);
media_index_result_free(&result);
return 0;
}

8
src/chat/member_sync.c

@ -230,9 +230,9 @@ static int _member_get_items(void* ctx, const char* ns, uint8_t level,
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;
static void _fire_props_changed(uint64_t node_id, const char* adm_tags) {
static void _fire_props_changed(uint64_t node_id, const char* adm_tags, const char* channel_id) {
struct ms_props_cbk* pc = g_props_cbks;
while (pc) { pc->fn(node_id, adm_tags, pc->arg); pc = pc->next; }
while (pc) { pc->fn(node_id, adm_tags, channel_id, pc->arg); pc = pc->next; }
}
void member_sync_add_props_cbk(node_props_changed_fn fn, void* arg) {
@ -283,7 +283,7 @@ static void _ms_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int even
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: bgp_node ev=%d nid=0x%016llx tags=%s",
MS_ID, event, (unsigned long long)node_id, tags);
_fire_props_changed(node_id, tags[0] ? tags : NULL);
_fire_props_changed(node_id, tags[0] ? tags : NULL, group->channel_id);
}
static int _sig_is_zero64(const uint8_t* sig) {
@ -689,7 +689,7 @@ static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer,
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; }
while (pc) { pc->fn(nid, tags_final, ns, pc->arg); pc = pc->next; }
}
}
if (r & MS_APPLY_STALE) {

3
src/chat/member_sync.h

@ -217,10 +217,11 @@ void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb);
/*
* Коллбэк: изменились adm_tags любого узла (включая себя).
* channel_id — канал (namespace), в котором произошло изменение.
* Вызывается при успешной обработке MSG_ITEM_UPDATE с adm_tags.
* Многоподписочный — можно добавить несколько подписчиков.
*/
typedef void (*node_props_changed_fn)(uint64_t node_id, const char* adm_tags, void* arg);
typedef void (*node_props_changed_fn)(uint64_t node_id, const char* adm_tags, const char* channel_id, 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);

60
src/media_delivery/media_delivery.c

@ -24,13 +24,13 @@
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 void md_on_props_changed(uint64_t node_id, const char* adm_tags, const char* channel_id, 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 void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id, uint64_t group_id);
static int md_send(struct UTUN_INSTANCE* inst, uint64_t group_id, 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);
static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id, uint64_t group_id);
/* ── helpers ── */
@ -198,7 +198,7 @@ static struct media_super_peer* md_super_peer_find(struct media_delivery_ctx* md
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) {
static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id, uint64_t group_id) {
struct media_super_peer* sp = md_super_peer_find(md, node_id);
if (sp) return sp;
if (!md->super_peers) {
@ -210,6 +210,7 @@ static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md,
sp = (struct media_super_peer*)qe->data;
memset(sp, 0, sizeof(*sp));
sp->peer_node_id = node_id;
sp->group_id = group_id;
sp->timeout_tb = MEDIA_REPL_TIMEOUT_TB;
sp->md = md;
memcpy(qe->data, &node_id, 8);
@ -346,7 +347,7 @@ static void md_reconnect_cb(void* arg) {
while (peer) {
if (!peer->connected && peer->connect_timer) {
peer->connect_timer = NULL;
md_super_connect(md, peer->peer_node_id);
md_super_connect(md, peer->peer_node_id, peer->group_id);
}
peer = peer->connected ? NULL : peer;
}
@ -384,7 +385,7 @@ static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super
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, TOPO_GROUP_UTUN, peer->peer_node_id, pkt, (size_t)off) == 0) {
if (md_send(md->inst, peer->group_id, peer->peer_node_id, pkt, (size_t)off) == 0) {
peer->inflight_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: SUPER_REPL sent to 0x%016llx seq=%u entries=%d",
MD_ID, (unsigned long long)peer->peer_node_id, seq, num_entries);
@ -413,9 +414,8 @@ static int md_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_n
memcpy(e->dgram + 1, data, len); e->len = (uint16_t)(len + 1);
int rc = etcp_route_send(inst, group_id, dst_node_id, e, 1);
if (rc != 0) {
struct ETCP_CONN* c = topo_group_find_conn_for_node(topo_groups_find(inst->topo_groups, group_id), dst_node_id);
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "%s: md_send FAIL rc=%d dst=0x%016llx group=0x%016llx route_conn=%p",
MD_ID, rc, (unsigned long long)dst_node_id, (unsigned long long)group_id, (void*)c);
DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "%s: md_send FAIL rc=%d dst=0x%016llx group=0x%016llx",
MD_ID, rc, (unsigned long long)dst_node_id, (unsigned long long)group_id);
u_free(e->dgram); queue_entry_free(e);
}
return rc;
@ -497,10 +497,10 @@ static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node,
uint16_t nc = (uint16_t)count;
memcpy(resp + num_pos, &nc, 2);
struct ETCP_CONN* qconn = topo_group_find_conn_for_node(topo_groups_find(md->inst->topo_groups, TOPO_GROUP_UTUN), from_node);
int snd_rc = md_send(md->inst, TOPO_GROUP_UTUN, from_node, resp, (size_t)off);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY from 0x%016llx → %d entries (super=%d) resp_len=%zu send_rc=%d route_conn=%p",
MD_ID, (unsigned long long)from_node, count, md->is_supernode, (size_t)off, snd_rc, (void*)qconn);
uint64_t resp_group = q->group_id ? q->group_id : TOPO_GROUP_UTUN;
md_send(md->inst, resp_group, from_node, resp, (size_t)off);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY from 0x%016llx → %d entries (super=%d)",
MD_ID, (unsigned long long)from_node, count, md->is_supernode);
}
static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_node,
@ -519,7 +519,7 @@ static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_no
memcpy(ha->block_id, hb->block_id, 16);
ha->status = (rc == 0) ? 0 : 1;
md_send(md->inst, TOPO_GROUP_UTUN, from_node, ack, sizeof(ack));
md_send(md->inst, hb->group_id, from_node, ack, sizeof(ack));
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: HAVE_BLOCK from 0x%016llx rc=%d",
MD_ID, (unsigned long long)from_node, rc);
@ -606,7 +606,8 @@ static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_no
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, TOPO_GROUP_UTUN, from_node, ack, sizeof(ack));
struct media_super_peer* peer = md_super_peer_find(md, from_node);
md_send(md->inst, peer ? peer->group_id : TOPO_GROUP_UTUN, from_node, ack, sizeof(ack));
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: SUPER_REPL from 0x%016llx entries=%d max_id=%lld",
MD_ID, (unsigned long long)from_node, num, (long long)max_id);
@ -636,7 +637,7 @@ static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_n
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);
struct media_super_peer* peer = md_super_peer_add(md, from_node, sh->group_id);
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;
@ -648,7 +649,8 @@ static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_n
struct media_pkt_super_hello resp;
resp.subcmd = MEDIA_SUBCMD_SUPER_HELLO;
resp.last_recv_id = my_last_recv;
md_send(md->inst, TOPO_GROUP_UTUN, from_node, (const uint8_t*)&resp, sizeof(resp));
resp.group_id = peer->group_id;
md_send(md->inst, peer->group_id, from_node, (const uint8_t*)&resp, sizeof(resp));
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%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);
@ -681,8 +683,9 @@ static void md_super_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64
struct media_pkt_super_hello hello;
hello.subcmd = MEDIA_SUBCMD_SUPER_HELLO;
hello.last_recv_id = my_last;
hello.group_id = peer->group_id;
if (md_send(md->inst, TOPO_GROUP_UTUN, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) {
if (md_send(md->inst, peer->group_id, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: SUPER_HELLO sent to 0x%016llx last_recv=%lld",
MD_ID, (unsigned long long)node_id, (long long)my_last);
} else {
@ -1036,18 +1039,11 @@ 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--; md->stream_completed = 1; }
}
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);
static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id, uint64_t group_id) {
struct media_super_peer* peer = md_super_peer_add(md, peer_node_id, group_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;
if (g->conn_mgr) { grp = g; break; }
gle = gle->next;
}
struct TOPO_GROUP* grp = topo_groups_find(md->inst->topo_groups, group_id);
if (!grp || !grp->conn_mgr) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: no group/conn_mgr for supernode connect", MD_ID);
return;
@ -1083,7 +1079,7 @@ static int md_super_start(struct media_delivery_ctx* md) {
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);
md_super_connect(md, pid, g->group_id);
}
sqlite3_finalize(st);
}
@ -1150,19 +1146,19 @@ static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event
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);
if (nt == 4) md_super_connect(md, node_id, group->group_id);
}
}
}
static void md_on_props_changed(uint64_t node_id, const char* adm_tags, void* arg) {
static void md_on_props_changed(uint64_t node_id, const char* adm_tags, const char* channel_id, 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);
if (is_super) md_super_connect(md, node_id, strtoull(channel_id, NULL, 10));
else md_super_peer_remove(md, node_id);
}
}

1
src/media_delivery/media_delivery.h

@ -28,6 +28,7 @@ struct UTUN_INSTANCE;
struct media_super_peer {
struct ll_entry ll; // индекс по peer_node_id (8 байт)
uint64_t peer_node_id;
uint64_t group_id; // CHAT-группа, в которой установлена связь с пиром
uint64_t peer_last_recv_id; // что пир говорит он получил от нас (из SUPER_HELLO / SUPER_ACK)
uint32_t timeout_tb; // текущий таймаут (растёт при ошибках)
uint8_t inflight_count; // пакетов в полёте (макс 4)

1
src/media_delivery/media_delivery_proto.h

@ -156,6 +156,7 @@ struct media_pkt_cancel {
struct media_pkt_super_hello {
uint8_t subcmd; // MEDIA_SUBCMD_SUPER_HELLO
uint64_t last_recv_id; // ID в МОЕЙ базе до которого пир подтвердил приём
uint64_t group_id; // CHAT-группа, в которой установлена связь
};
struct media_pkt_block_overloaded {

31
src/media_delivery/media_download.c

@ -59,7 +59,7 @@ static void md_mkdir_parent(const char* filepath) {
}
}
static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, size_t len) {
static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst, const uint8_t* data, size_t len) {
if (!inst || !inst->topo_groups || !inst->connections) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: md_dl_send: no routing (inst=%p topo=%p conn=%p)", MDL_ID, (void*)inst, (void*)(inst ? inst->topo_groups : NULL), (void*)(inst ? inst->connections : NULL)); return -1; }
struct ll_entry* e = queue_entry_new(0);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new failed for send to 0x%016llx", MDL_ID, (unsigned long long)dst); return -1; }
@ -67,8 +67,8 @@ static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* d
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_malloc(%zu) failed for send to 0x%016llx", MDL_ID, len + 1, (unsigned long long)dst); 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, e, 1);
if (rc != 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: etcp_route_send failed rc=%d to 0x%016llx", MDL_ID, rc, (unsigned long long)dst); u_free(e->dgram); queue_entry_free(e); }
int rc = etcp_route_send(inst, group_id, dst, e, 1);
if (rc != 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: etcp_route_send failed rc=%d to 0x%016llx group=0x%016llx", MDL_ID, rc, (unsigned long long)dst, (unsigned long long)group_id); u_free(e->dgram); queue_entry_free(e); }
return rc;
}
@ -89,7 +89,7 @@ static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst,
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_REQ to 0x%016llx group=%016llx bi=%d dl->num_blocks=%d block_id=%02x%02x%02x%02x...",
MDL_ID, (unsigned long long)dst, (unsigned long long)dl->group_id, peer_local_bi, dl->num_blocks,
block_id[0], block_id[1], block_id[2], block_id[3]);
return md_dl_send(inst, dst, (const uint8_t*)&req, sizeof(req));
return md_dl_send(inst, dl->group_id, dst, (const uint8_t*)&req, sizeof(req));
}
/* ── forward decl ── */
@ -144,7 +144,7 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%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));
md_dl_send(inst, dl->group_id, super, (const uint8_t*)&hb, sizeof(hb));
}
/* ── send BLOCK_PROCESSING to supernode ── */
@ -174,7 +174,7 @@ static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_PROCESSING to super 0x%016llx block=%d",
MDL_ID, (unsigned long long)super, bi);
md_dl_send(inst, super, (const uint8_t*)&hb, sizeof(hb));
md_dl_send(inst, dl->group_id, super, (const uint8_t*)&hb, sizeof(hb));
}
/* ── conn_mgr callback ── */
@ -266,7 +266,7 @@ static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download*
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%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);
md_dl_send(inst, dl->group_id, super, pkt, pkt_len);
u_free(pkt);
}
@ -274,7 +274,6 @@ static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download*
void media_download_handle_query_resp(struct UTUN_INSTANCE* inst,
const uint8_t* data, size_t len) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY_RESP received len=%zu data=%p", MDL_ID, len, (void*)data);
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);
@ -339,14 +338,8 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst,
return;
}
/* find a group with conn_mgr for direct connection (best-effort) */
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;
if (g->conn_mgr) { grp = g; break; }
gle = gle->next;
}
/* find the chat group for direct connection (best-effort) */
struct TOPO_GROUP* grp = topo_groups_find(inst->topo_groups, dl->group_id);
for (int pi = 0; pi < dl->num_peers; pi++) {
dl->peers[pi].connected = 1;
@ -422,7 +415,7 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst,
rch->offset = (uint32_t)ds->sent_offset;
rch->data_len = (uint16_t)rd;
memcpy(rpkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, rbuf, rd);
if (md_dl_send(inst, ds->node_id, rpkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd) == 0)
if (md_dl_send(inst, ds->group_id, ds->node_id, rpkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd) == 0)
ds->sent_offset += rd;
else break;
}
@ -513,7 +506,7 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
memcpy(rd.block_id, dl->block_ids + bi * 16, 16);
rd.chunk = (uint32_t)bi;
rd.total_size = (uint32_t)rc->file_offset;
md_dl_send(inst, rc->downstream[i].node_id, (const uint8_t*)&rd, sizeof(rd));
md_dl_send(inst, rc->downstream[i].group_id, rc->downstream[i].node_id, (const uint8_t*)&rd, sizeof(rd));
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: relay BLOCK_DONE forwarded to 0x%016llx block=%d",
MDL_ID, (unsigned long long)rc->downstream[i].node_id, bi);
}
@ -695,7 +688,7 @@ int media_download_cancel(struct UTUN_INSTANCE* inst,
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));
md_dl_send(inst, dl->group_id, dl->peers[pi].node_id, (const uint8_t*)&can, sizeof(can));
}
}

2
tests/test_media_delivery_full.c

@ -172,7 +172,7 @@ static void phase_a1_super_hello(void) {
if (g_inst[2]->md.is_supernode && g_inst[3]->md.is_supernode) OK(); else FAIL();
}
TEST("SUPER_HELLO s1↔s2"); {
struct media_pkt_super_hello h; h.subcmd = MEDIA_SUBCMD_SUPER_HELLO; h.last_recv_id = 0;
struct media_pkt_super_hello h; h.subcmd = MEDIA_SUBCMD_SUPER_HELLO; h.last_recv_id = 0; h.group_id = TOPO_GROUP_UTUN;
msend(g_inst[2], g_nid[3], (const uint8_t*)&h, sizeof(h));
msend(g_inst[3], g_nid[2], (const uint8_t*)&h, sizeof(h));
int a = 0; while (a < 500) { uasync_poll(g_ua, POLL_MS); a++; }

5
tests/test_media_delivery_sql.c

@ -398,7 +398,7 @@ static void test_proto_sizes(void) {
int ok = 1;
if (sizeof(struct media_pkt_query) != 27) { FAIL("query: %zu != 27", sizeof(struct media_pkt_query)); ok = 0; }
if (sizeof(struct media_pkt_have_block) != 117) { FAIL("have_block: %zu != 117", sizeof(struct media_pkt_have_block)); ok = 0; }
if (sizeof(struct media_pkt_super_hello) != 9) { FAIL("super_hello: %zu != 9", sizeof(struct media_pkt_super_hello)); ok = 0; }
if (sizeof(struct media_pkt_super_hello) != 17) { FAIL("super_hello: %zu != 17", sizeof(struct media_pkt_super_hello)); ok = 0; }
if (sizeof(struct media_pkt_super_ack) != 5) { FAIL("super_ack: %zu != 5", sizeof(struct media_pkt_super_ack)); ok = 0; }
if (sizeof(struct media_pkt_block_req) != 53) { FAIL("block_req: %zu != 53", sizeof(struct media_pkt_block_req)); ok = 0; }
if (sizeof(struct media_pkt_cancel) != 37) { FAIL("cancel: %zu != 37", sizeof(struct media_pkt_cancel)); ok = 0; }
@ -443,7 +443,8 @@ static void test_proto_build(void) {
struct media_pkt_super_hello sh;
sh.subcmd = MEDIA_SUBCMD_SUPER_HELLO;
sh.last_recv_id = 12345;
if (sh.subcmd == MEDIA_SUBCMD_SUPER_HELLO && sh.last_recv_id == 12345) OK();
sh.group_id = 0x11223344;
if (sh.subcmd == MEDIA_SUBCMD_SUPER_HELLO && sh.last_recv_id == 12345 && sh.group_id == 0x11223344) OK();
else FAIL();
}
}

Loading…
Cancel
Save