From 6902f5fd95a0659264beb6e64d58cdc91160c4ae Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 15 Jul 2026 18:15:16 +0300 Subject: [PATCH] chat: add ll_queue with index for P2P connections (replaces linear search) - utun_instance.h: add chat_connections queue (indexed by node_id, 8 bytes) - chat_sync.c: queue_new in init, queue_free in destroy, queue_find_data_by_index in cs_find_conn_for_node - chat_sync.c: remove from queue in cs_on_conn_down - chat_core.c: add to queue in connect_from_invite and connect_auto - chat_core.c: remove from queue in cc_parallel_cleanup and ca_cleanup - db_sync.c: replace linear conn search with queue_find_data_by_index (db_sync_send_hash) --- src/db_sync.c | 9 ++++----- src/utun_instance.h | 1 + tools/chatgui/transport/chat_core.c | 11 +++++++++++ tools/chatgui/transport/chat_sync.c | 17 ++++++++++++----- 4 files changed, 28 insertions(+), 10 deletions(-) diff --git a/src/db_sync.c b/src/db_sync.c index e72d00ff..32876ad2 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -487,11 +487,10 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, entry->dgram = buf; entry->len = plen + 9; - struct ETCP_CONN* conn = db->inst->connections; - while (conn) { - if (conn->peer_node_id == node_id && conn->links_up > 0) - break; - conn = conn->next; + struct ETCP_CONN* conn = NULL; + if (db->inst->chat_connections) { + struct ll_entry* e = queue_find_data_by_index(db->inst->chat_connections, (const uint8_t*)&node_id); + if (e) { struct { uint64_t node_id; struct ETCP_CONN* conn; }* ce = e->data; if (ce->conn->links_up) conn = ce->conn; } } if (!conn) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, diff --git a/src/utun_instance.h b/src/utun_instance.h index 9251c559..1e705669 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -135,6 +135,7 @@ struct UTUN_INSTANCE { // etcp_router bindings и seq-connections (per-instance service routing) struct ETCP_ROUTER_BINDINGS router_bindings; struct ll_queue* router_conns; + struct ll_queue* chat_connections; // chat P2P connections indexed by node_id struct CONN_MGR* conn_mgr; // Connection Manager (может быть NULL) struct DB_SYNC* db_sync; // Distributed DB sync (может быть NULL) diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 3c34763a..0016561a 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -564,6 +564,10 @@ struct cc_parallel_ctx { static void cc_parallel_cleanup(struct cc_parallel_state* st) { if (!st) return; + if (st->inst && st->inst->chat_connections) { + struct ll_entry* e = queue_find_data_by_index(st->inst->chat_connections, (const uint8_t*)&st->node_id); + if (e) { queue_remove_data(st->inst->chat_connections, e); queue_entry_free(e); } + } u_free(st->conns); u_free(st->timers); u_free(st); @@ -760,6 +764,8 @@ void chat_core_connect_from_invite(struct chat_invite* inv) { pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; } pst->conns[idx] = conn; + conn->peer_node_id = inv->node_id; + { struct ll_entry* qe = queue_entry_new(16); if (qe) { *(uint64_t*)qe->data = inv->node_id; *(struct ETCP_CONN**)(qe->data+8) = conn; queue_data_put_with_index(g_cc.inst->chat_connections, qe); } } pst->timers[idx] = uasync_set_timeout(g_cc.inst->ua, CC_PARALLEL_CONNECT_TIMEOUT_MS * 10, pctx, cc_parallel_timeout_cb, "cc_parallel"); DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect attempt %d/%d to %d.%d.%d.%d:%d", @@ -804,6 +810,10 @@ struct ca_state { static void ca_cleanup(struct ca_state* st) { if (!st) return; + if (st->inst && st->inst->chat_connections) { + struct ll_entry* e = queue_find_data_by_index(st->inst->chat_connections, (const uint8_t*)&st->node_id); + if (e) { queue_remove_data(st->inst->chat_connections, e); queue_entry_free(e); } + } if (st->ctxs) { for (int i = 0; i < st->addr_count; i++) u_free(st->ctxs[i]); u_free(st->ctxs); } u_free(st->conns); u_free(st->timers); @@ -970,6 +980,7 @@ void chat_core_connect_auto(uint64_t node_id, pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; } pst->conns[i] = conn; + { struct ll_entry* qe = queue_entry_new(16); if (qe) { *(uint64_t*)qe->data = node_id; *(struct ETCP_CONN**)(qe->data+8) = conn; queue_data_put_with_index(g_cc.inst->chat_connections, qe); } } pst->timers[i] = uasync_set_timeout(g_cc.inst->ua, CA_CONNECT_TIMEOUT_MS * 10, pctx, ca_timeout_cb, "ca_timeout"); DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect attempt %d/%d to %d.%d.%d.%d:%d", diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index fd8057bd..a6a679d8 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -433,12 +433,11 @@ static void cs_join_timeout_cb(void* arg) { /* ── helper: find active ETCP_CONN for node ── */ +struct CHAT_CONN_ENTRY { uint64_t node_id; struct ETCP_CONN* conn; }; + static struct ETCP_CONN* cs_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { - struct ETCP_CONN* c = inst->connections; - while (c) { - if (c->peer_node_id == node_id && c->initialized && c->links_up) return c; - c = c->next; - } + struct ll_entry* e = queue_find_data_by_index(inst->chat_connections, (const uint8_t*)&node_id); + if (e) { struct CHAT_CONN_ENTRY* ce = (struct CHAT_CONN_ENTRY*)e->data; if (ce->conn->initialized && ce->conn->links_up) return ce->conn; } return NULL; } @@ -510,6 +509,10 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) { (void)arg; if (!conn || !g_cs) return; uint64_t peer = conn->peer_node_id; + if (g_cs->inst->chat_connections) { + struct ll_entry* e = queue_find_data_by_index(g_cs->inst->chat_connections, (const uint8_t*)&peer); + if (e) { queue_remove_data(g_cs->inst->chat_connections, e); queue_entry_free(e); } + } uint16_t rtt = conn->rtt_avg_100; if (rtt > 0 && peer != 0 && g_cs->inst->topo_groups && g_cs->inst->topo_groups->topo_sqlite_db) { sqlite3* db = g_cs->inst->topo_groups->topo_sqlite_db; @@ -614,6 +617,8 @@ int chat_sync_init(struct UTUN_INSTANCE* inst, etcp_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb); etcp_add_new_conn_cbk(inst, cs_on_new_conn, NULL); + inst->chat_connections = queue_new(inst->ua, 16, 0, 8, "chat_conns"); + struct ETCP_CONN* c = inst->connections; while (c) { etcp_conn_add_up_cbk(c, cs_on_conn_up, NULL); @@ -641,6 +646,8 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { member_sync_destroy(inst); cs->initialized = 0; g_cs = NULL; + queue_free(inst->chat_connections); inst->chat_connections = NULL; + etcp_unbind(inst, ETCP_RT_ID_CHAT_SYNC); if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; } cs_cancel_proto_timers(cs);