diff --git a/src/chat/chat_member.c b/src/chat/chat_member.c index 26a5e0bd..288787f6 100644 --- a/src/chat/chat_member.c +++ b/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); } } diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index d79d1931..f9bf2242 100644 --- a/src/chat/chat_msg.c +++ b/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; } diff --git a/src/chat/member_sync.c b/src/chat/member_sync.c index b5dccb3b..b98b191e 100644 --- a/src/chat/member_sync.c +++ b/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) { diff --git a/src/chat/member_sync.h b/src/chat/member_sync.h index ae6f720b..c99a350d 100644 --- a/src/chat/member_sync.h +++ b/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); diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index f1047f33..1ee2346c 100644 --- a/src/media_delivery/media_delivery.c +++ b/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); } } diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index 4d7491dc..56f967de 100644 --- a/src/media_delivery/media_delivery.h +++ b/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) diff --git a/src/media_delivery/media_delivery_proto.h b/src/media_delivery/media_delivery_proto.h index 7e596002..a94a47ea 100644 --- a/src/media_delivery/media_delivery_proto.h +++ b/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 { diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index fd46687e..3ef28c5a 100644 --- a/src/media_delivery/media_download.c +++ b/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)); } } diff --git a/tests/test_media_delivery_full.c b/tests/test_media_delivery_full.c index 7e338aca..4c57b3dd 100644 --- a/tests/test_media_delivery_full.c +++ b/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++; } diff --git a/tests/test_media_delivery_sql.c b/tests/test_media_delivery_sql.c index 8b0e6772..9d63c214 100644 --- a/tests/test_media_delivery_sql.c +++ b/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(); } }