Browse Source

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)
topo_upd
Evgeny 3 months ago
parent
commit
6902f5fd95
  1. 9
      src/db_sync.c
  2. 1
      src/utun_instance.h
  3. 11
      tools/chatgui/transport/chat_core.c
  4. 17
      tools/chatgui/transport/chat_sync.c

9
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,

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

11
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",

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

Loading…
Cancel
Save