From a7ecee01f5dbab0d2a72fc6d12247233a2e7695d Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 15 Jul 2026 19:56:47 +0300 Subject: [PATCH] etcp: replace connections linked-list with indexed ll_queue, add state/ref_count lifecycle MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Remove ETCP_CONN.next linked-list, UTUN_INSTANCE.connections list/count - Unified connections in single ll_queue indexed by peer_node_id (0=pending) - struct conn_queue_entry { peer_node_id, conn* } for uniform queue data - Connection lifecycle via etcp.c: create→pending, conn_ready→indexed, close→detach+deferred free - Close split: phase1 detach from external world, phase2 deferred via uasync_call_soon - etcp_conn_set_peer_node_id: reindex on ID change (WARN) - etcp_conn_ref_take/free: blocking close until refcount returns to 0 - Remove chat_connections — use unified connections queue directly - Update all iteration/search patterns across 40 files (src + tests + chatgui) --- src/control_server.c | 30 ++- src/db_sync.c | 47 ++-- src/etcp.c | 224 +++++++++++++------ src/etcp.h | 33 ++- src/etcp_connect.c | 33 +-- src/etcp_connections.c | 49 ++-- src/etcp_dump.c | 15 +- src/ntp_node_time.c | 7 +- src/topo_group.c | 6 +- src/utun_instance.c | 49 ++-- src/utun_instance.h | 12 +- tests/bbr_integration/test_bbr_integration.c | 101 +++++---- tests/test_bgp_route_exchange.c | 23 +- tests/test_bgp_triangle.c | 22 +- tests/test_conn_mgr.c | 7 +- tests/test_db_sync.c | 7 +- tests/test_etcp_100_packets.c | 18 +- tests/test_etcp_api.c | 14 +- tests/test_etcp_congestion.c | 16 +- tests/test_etcp_connect.c | 7 +- tests/test_etcp_dummynet.c | 10 +- tests/test_etcp_ping.c | 18 +- tests/test_etcp_reconnect.c | 29 +-- tests/test_etcp_reinit_inflight.c | 28 ++- tests/test_etcp_router.c | 8 +- tests/test_etcp_router_reconnect.c | 7 +- tests/test_etcp_simple_traffic.c | 14 +- tests/test_etcp_two_instances.c | 16 +- tests/test_icmp_proxy.c | 10 +- tests/test_ipv6_sockets.c | 10 +- tests/test_nat_detection.c | 41 ++-- tests/test_nat_transport.c | 13 +- tests/test_pkt_normalizer_etcp.c | 20 +- tests/test_route_ping.c | 20 +- tests/test_socks_http_proxy.c | 9 +- tests/test_stcp_traffic.c | 7 +- tests/test_tcp_proxy_remote.c | 10 +- tests/test_udp_proxy.c | 10 +- tools/chatgui/transport/chat_core.c | 17 +- tools/chatgui/transport/chat_sync.c | 80 ++++--- 40 files changed, 649 insertions(+), 448 deletions(-) diff --git a/src/control_server.c b/src/control_server.c index 039db141..1c206492 100644 --- a/src/control_server.c +++ b/src/control_server.c @@ -885,10 +885,10 @@ static void send_conn_list(struct control_server* server, struct control_client* /* Count connections */ uint8_t count = 0; - struct ETCP_CONN* conn = instance->connections; - while (conn && count < ETCPMON_MAX_CONNECTIONS) { + struct ll_entry* entry = instance->connections->head; + while (entry && count < ETCPMON_MAX_CONNECTIONS) { count++; - conn = conn->next; + entry = entry->next; } /* Build response */ @@ -909,12 +909,13 @@ static void send_conn_list(struct control_server* server, struct control_client* rsp->count = count; struct etcpmon_conn_info* info = (struct etcpmon_conn_info*)(buffer + sizeof(*hdr) + sizeof(*rsp)); - conn = instance->connections; - for (uint8_t i = 0; i < count && conn; i++) { - info[i].peer_node_id = conn->peer_node_id; - strncpy(info[i].name, conn->log_name, ETCPMON_MAX_CONN_NAME - 1); + entry = instance->connections->head; + for (uint8_t i = 0; i < count && entry; i++) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + info[i].peer_node_id = ce->conn->peer_node_id; + strncpy(info[i].name, ce->conn->log_name, ETCPMON_MAX_CONN_NAME - 1); info[i].name[ETCPMON_MAX_CONN_NAME - 1] = '\0'; - conn = conn->next; + entry = entry->next; } /* Log and send response */ @@ -1372,14 +1373,11 @@ static void send_debug_config(struct control_server* server, struct control_clie static struct ETCP_CONN* find_connection_by_peer_id(struct UTUN_INSTANCE* instance, uint64_t peer_id) { - struct ETCP_CONN* conn = instance->connections; - while (conn) { - if (conn->peer_node_id == peer_id) { - return conn; - } - conn = conn->next; - } - return NULL; + if (!instance || !instance->connections) return NULL; + struct ll_entry* e = queue_find_data_by_index(instance->connections, (const uint8_t*)&peer_id); + if (!e) return NULL; + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + return ce->conn; } /* ============================================================================ diff --git a/src/db_sync.c b/src/db_sync.c index 32876ad2..681c96e4 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -488,9 +488,9 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, entry->len = plen + 9; 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; } + { + struct ll_entry* e = queue_find_data_by_index(db->inst->connections, (const uint8_t*)&node_id); + if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; if (ce->conn->links_up) conn = ce->conn; } } if (!conn) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, @@ -1430,11 +1430,16 @@ int db_sync_init(struct UTUN_INSTANCE* inst) etcp_set_new_conn_cbk(inst, db_sync_on_new_conn, NULL); // Attach to existing connections - struct ETCP_CONN* conn = inst->connections; - while (conn) { - etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL); - etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL); - conn = conn->next; + { + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce && ce->conn) { + etcp_conn_add_up_cbk(ce->conn, db_sync_on_conn_up, NULL); + etcp_conn_add_down_cbk(ce->conn, db_sync_on_conn_down, NULL); + } + entry = entry->next; + } } if (db->enabled) @@ -1474,11 +1479,16 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) } // Detach from existing connections - struct ETCP_CONN* conn = inst->connections; - while (conn) { - etcp_conn_remove_up_cbk(conn, db_sync_on_conn_up, NULL); - etcp_conn_remove_down_cbk(conn, db_sync_on_conn_down, NULL); - conn = conn->next; + { + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce && ce->conn) { + etcp_conn_remove_up_cbk(ce->conn, db_sync_on_conn_up, NULL); + etcp_conn_remove_down_cbk(ce->conn, db_sync_on_conn_down, NULL); + } + entry = entry->next; + } } db_sqlite_close(db); @@ -1579,11 +1589,12 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, // Initiate sync with already connected peers { - struct ETCP_CONN* conn = inst->connections; - while (conn) { - uint64_t pid = conn->peer_node_id; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + uint64_t pid = ce->conn->peer_node_id; if (pid != 0 && pid != inst->node_id - && conn->links_up > 0 && conn->initialized) + && ce->conn->links_up > 0 && ce->conn->initialized) { struct SI_PEER* p = si_peer_add(si, pid); if (p && p->synced == 0) { @@ -1591,7 +1602,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, db_sync_initiate_sync(si, pid); } } - conn = conn->next; + entry = entry->next; } } diff --git a/src/etcp.c b/src/etcp.c index 44775316..4e1357bf 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -44,6 +44,7 @@ static void input_send_q_cb(struct ll_queue* q, void* arg); static void wait_ack_cb(struct ll_queue* q, void* arg); static void send_ack_req_cb(struct ll_queue* q, void* arg); static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp); +static void etcp_connection_free_deferred(void* arg); struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp); static void clear_queue(struct ll_queue* q) { @@ -217,6 +218,7 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n etcp->links_up=0; etcp->reset_done=0; etcp->callbacks_running=0; + etcp->ref_count=0; etcp->last_rr_link=NULL; etcp->name = u_strdup(name); // Initialize log_name with local node_id (peer will be updated later when known) @@ -256,7 +258,19 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n etcp->link_ready_for_send_fn = etcp_link_ready_callback; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] connection initialized. ETCP=%p mtu=%d, next_tx_id=%u", + // Add to instance's connections queue (peer_node_id=0 for pending, reindexed later) + { + struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry)); + if (!qe) { etcp_connection_close(etcp); return NULL; } + struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; + ce->peer_node_id = 0; + ce->conn = etcp; + etcp->conn_queue_entry = qe; + etcp->conn_queue = instance->connections; + queue_data_put_with_index(instance->connections, qe); + } + + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] connection initialized. ETCP=%p mtu=%d, next_tx_id=%u state=pending", etcp->log_name, etcp, etcp->mtu, etcp->next_tx_id); // Вызываем callback для нового соединения если установлен @@ -286,57 +300,15 @@ static void etcp_on_down(struct ETCP_CONN* etcp) { } -// Close connection with NULL pointer safety (prevents double free) -void etcp_connection_close(struct ETCP_CONN* etcp) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); +// Phase 2 resources cleanup (callable both sync and async via call_soon) +static void etcp_connection_free_resources(struct ETCP_CONN* etcp) { if (!etcp) return; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] freeing resources phase 2", etcp->log_name); - if (etcp->callbacks_running) { - DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "[%s] FATAL: etcp_connection_close called from inside callback chain — SEGFAULTING to show backtrace", - etcp->log_name); - *(volatile int*)0 = 0; - } - - if (etcp->links_up!=0) { - etcp->links_up=0; - etcp_on_down(etcp); - } - - // Cancel active timers to prevent memory leaks - if (etcp->retrans_timer) { - uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer); - etcp->retrans_timer = NULL; - } - if (etcp->ack_resp_timer) { - uasync_cancel_timeout(etcp->instance->ua, etcp->ack_resp_timer); - etcp->ack_resp_timer = NULL; - } - etcp_metrics_stop_timer(etcp); - - routing_del_conn(etcp); - - // NOTE: topo_bgp_remove_conn is already called via down_cbk in etcp_on_down() above - // DO NOT call it again here to avoid double-processing and ref_count corruption - - // Очистить транзитные очереди (до pn_deinit/stcp_link_close — waiter'ы на normalizer->input/tx_queue) - etcp_router_transit_queues_destroy(etcp); - - // Deinitialize packet normalizer (this will call routing_del_conn) - if (etcp->normalizer) { - pn_deinit((struct PKTNORM*)etcp->normalizer); - etcp->normalizer = NULL; - } - - // Clear links BEFORE draining queues: etcp_link_close calls - // etcp_conn_on_inflight_lim_changed -> input_queue_try_resume - // which accesses input_send_q, input_queue, input_wait_ack. + // Close links if (etcp->links) { struct ETCP_LINK* link = etcp->links; - while (link) { - struct ETCP_LINK* next = link->next; - etcp_link_close(link); - link = next; - } + while (link) { struct ETCP_LINK* next = link->next; etcp_link_close(link); link = next; } etcp->links = NULL; } @@ -350,35 +322,86 @@ void etcp_connection_close(struct ETCP_CONN* etcp) { // Free callback chains { struct etcp_cbk_entry* cbe = etcp->ready_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->ready_cbks = NULL; } - { struct etcp_cbk_entry* cbe = etcp->up_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->up_cbks = NULL; } - { struct etcp_cbk_entry* cbe = etcp->down_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->down_cbks = NULL; } + { struct etcp_cbk_entry* cbe = etcp->up_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->up_cbks = NULL; } + { struct etcp_cbk_entry* cbe = etcp->down_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->down_cbks = NULL; } - // Free memory pools after all elements are returned - if (etcp->inflight_pool) { - memory_pool_destroy(etcp->inflight_pool); - etcp->inflight_pool = NULL; - } + if (etcp->inflight_pool) { memory_pool_destroy(etcp->inflight_pool); etcp->inflight_pool = NULL; } + if (etcp->io_pool) { memory_pool_destroy(etcp->io_pool); etcp->io_pool = NULL; } + + u_free(etcp->name); + u_free(etcp); +} + +static void etcp_connection_free_deferred(void* arg) { + etcp_connection_free_resources((struct ETCP_CONN*)arg); +} + +// Close connection: phase 1 detach + deferred phase 2 cleanup +void etcp_connection_close(struct ETCP_CONN* etcp) { + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); + if (!etcp) return; + if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] already deleted", etcp->log_name); return; } - if (etcp->io_pool) { - memory_pool_destroy(etcp->io_pool); - etcp->io_pool = NULL; + if (etcp->callbacks_running) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "[%s] FATAL: etcp_connection_close called from inside callback chain — SEGFAULTING to show backtrace", + etcp->log_name); + *(volatile int*)0 = 0; } - u_free(etcp->name); + // === PHASE 1: detach from external world === - { struct UTUN_INSTANCE* inst = etcp->instance; - if (inst && inst->connections) { - struct ETCP_CONN** pp = &inst->connections; - while (*pp && *pp != etcp) pp = &(*pp)->next; - if (*pp) { *pp = etcp->next; if (inst->connections_count) inst->connections_count--; } - } + if (etcp->links_up != 0) { etcp->links_up = 0; etcp_on_down(etcp); } + + // Cancel active timers + if (etcp->retrans_timer) { uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer); etcp->retrans_timer = NULL; } + if (etcp->ack_resp_timer) { uasync_cancel_timeout(etcp->instance->ua, etcp->ack_resp_timer); etcp->ack_resp_timer = NULL; } + etcp_metrics_stop_timer(etcp); + + routing_del_conn(etcp); + etcp_router_transit_queues_destroy(etcp); + + if (etcp->normalizer) { pn_deinit((struct PKTNORM*)etcp->normalizer); etcp->normalizer = NULL; } + + // Close links (detach from socket scan queues) + if (etcp->links) { + struct ETCP_LINK* link = etcp->links; + while (link) { struct ETCP_LINK* next = link->next; etcp_link_close(link); link = next; } + etcp->links = NULL; } - // Clear next pointer to prevent dangling references - etcp->next = NULL; + // Remove from instance queue + if (etcp->conn_queue && etcp->conn_queue_entry) { + queue_remove_data(etcp->conn_queue, etcp->conn_queue_entry); + queue_entry_free(etcp->conn_queue_entry); + etcp->conn_queue_entry = NULL; + etcp->conn_queue = NULL; + } - // TODO: Free rx_list, etc. - u_free(etcp); + etcp->state = 2; // deleted + + // === PHASE 2: deferred resource cleanup (only if no outstanding refs) === + if (etcp->ref_count == 0) + uasync_call_soon(etcp->instance->ua, etcp, etcp_connection_free_deferred); + // else: cleanup deferred until last etcp_conn_ref_free() +} + +// Take a reference on connection. Returns -1 if already deleted (state==2). +int etcp_conn_ref_take(struct ETCP_CONN* conn) { + if (!conn) return -1; + if (conn->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ref_take on deleted conn", conn->log_name); return -1; } + conn->ref_count++; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ref_take ref_count=%d", conn->log_name, conn->ref_count); + return 0; +} + +// Release a reference. If last ref and conn is deleted, schedule deferred free. +void etcp_conn_ref_free(struct ETCP_CONN* conn) { + if (!conn) return; + if (conn->ref_count <= 0) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ref_free underflow ref_count=%d", conn->log_name, conn->ref_count); return; } + conn->ref_count--; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ref_free ref_count=%d state=%d", conn->log_name, conn->ref_count, conn->state); + if (conn->ref_count == 0 && conn->state == 2) + uasync_call_soon(conn->instance->ua, conn, etcp_connection_free_deferred); } // Reset connection @@ -543,24 +566,77 @@ void etcp_conn_reinit(struct ETCP_CONN* etcp) {// Если сбой в обме // внутренняя функция. Вызывается один раз когда первый линк готов. void etcp_conn_ready(struct ETCP_CONN* conn) { if (!conn) return; + if (conn->initialized) return; // already ready conn->initialized = 1; conn->reset_done = 1; - // Если tx_state не установлен (первое подключение без reinit), установим его - if (conn->tx_state == 0) { - conn->tx_state = ETCP_TX_STATE_DATA_WAIT; - } + if (conn->tx_state == 0) { conn->tx_state = ETCP_TX_STATE_DATA_WAIT; } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Connection ready", conn->log_name); + etcp_conn_queue_set_ready(conn); +} + +// Move from pending to indexed connections, fire ready callbacks +void etcp_conn_queue_set_ready(struct ETCP_CONN* conn) { + if (!conn || conn->state != 0) return; + + // Reindex: remove old entry (key=0), add new entry with real peer_node_id + if (conn->conn_queue && conn->conn_queue_entry) { + queue_remove_data(conn->conn_queue, conn->conn_queue_entry); + queue_entry_free(conn->conn_queue_entry); + } + + struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry)); + if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to alloc queue entry for ready", conn->log_name); return; } + struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; + ce->peer_node_id = conn->peer_node_id; + ce->conn = conn; + conn->conn_queue_entry = qe; + conn->conn_queue = conn->instance->connections; + conn->state = 1; + queue_data_put_with_index(conn->instance->connections, qe); + + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] moved to ready queue, state=%d peer_node_id=0x%llx", + conn->log_name, conn->state, (unsigned long long)conn->peer_node_id); + etcp_metrics_start_timer(conn); - // Вызываем callback если установлен conn->callbacks_running = 1; { struct etcp_cbk_entry* cbe = conn->ready_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(conn, cbe->arg); cbe = n; } } if (conn->links_up) etcp_on_up(conn); conn->callbacks_running = 0; } +// Set peer_node_id, reindex if already in connections queue +void etcp_conn_set_peer_node_id(struct ETCP_CONN* conn, uint64_t peer_node_id) { + if (!conn) return; + if (conn->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] set_peer_node_id on deleted conn", conn->log_name); return; } + + if (conn->peer_node_id == peer_node_id) return; + + if (conn->state == 1 && conn->peer_node_id != 0) { + // Already in indexed queue with old ID — reindex + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] peer_node_id changed 0x%llx -> 0x%llx, reindexing", + conn->log_name, (unsigned long long)conn->peer_node_id, (unsigned long long)peer_node_id); + if (conn->conn_queue && conn->conn_queue_entry) { + queue_remove_data(conn->conn_queue, conn->conn_queue_entry); + queue_entry_free(conn->conn_queue_entry); + } + struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry)); + if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] reindex alloc failed", conn->log_name); return; } + struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; + ce->peer_node_id = peer_node_id; + ce->conn = conn; + conn->conn_queue_entry = qe; + conn->conn_queue = conn->instance->connections; + queue_data_put_with_index(conn->instance->connections, qe); + } + + conn->peer_node_id = peer_node_id; + etcp_update_log_name(conn); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] peer_node_id set to 0x%llx", conn->log_name, (unsigned long long)peer_node_id); +} + // Update log_name when peer_node_id becomes known void etcp_update_log_name(struct ETCP_CONN* etcp) { diff --git a/src/etcp.h b/src/etcp.h index 0b3b7faa..640a13eb 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -138,11 +138,19 @@ struct ACK_PACKET { // ETCP connection structure (refactored) struct ETCP_CONN { - struct ETCP_CONN* next; + // State: 0=not ready, 1=ready (indexed in instance->connections), 2=deleted + int state; // 0=pending, 1=ready, 2=deleted (phase 1 of close done) + int ref_count; // External reference count. >0 blocks deferred resource free. + // Take/free via etcp_conn_ref_take()/etcp_conn_ref_free(). + int mtu; struct UTUN_INSTANCE* instance; + // Queue entries in instance->connections or instance->pending_connections + struct ll_entry* conn_queue_entry; // entry в очереди instance + struct ll_queue* conn_queue; // указатель на очередь где лежим (pending или connections) + // Links (channels) - linked list struct ETCP_LINK* links; struct ETCP_LINK* last_rr_link; // последний линк, выбранный round-robin @@ -252,6 +260,29 @@ struct ETCP_CONN { // Functions struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name); void etcp_connection_close(struct ETCP_CONN* etcp); +void etcp_conn_set_peer_node_id(struct ETCP_CONN* conn, uint64_t peer_node_id); +void etcp_conn_queue_set_ready(struct ETCP_CONN* conn); // move pending->connections, fire ready cbks + +/** + * @brief Take a reference on the connection (blocks deferred resource free). + * @param conn connection + * @return 0 on success, -1 if conn is NULL or already in state 2 (deleted). + * + * Increments ref_count. While ref_count > 0, etcp_connection_close() + * will detach but defer resource cleanup until all references are released. + * Always pair with etcp_conn_ref_free(). + */ +int etcp_conn_ref_take(struct ETCP_CONN* conn); + +/** + * @brief Release a reference. Triggers deferred resource free if last ref and state==2. + * @param conn connection + * + * Decrements ref_count. If ref_count reaches 0 and the connection has been closed + * (state == 2), schedules deferred resource cleanup via uasync_call_soon. + * Safe to call on NULL. + */ +void etcp_conn_ref_free(struct ETCP_CONN* conn); void etcp_conn_reset(struct ETCP_CONN* etcp); void etcp_links_reset(struct ETCP_CONN* etcp); diff --git a/src/etcp_connect.c b/src/etcp_connect.c index f18d4281..3f909126 100644 --- a/src/etcp_connect.c +++ b/src/etcp_connect.c @@ -210,22 +210,26 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node, if (!inst || !node || !cb) return -1; uint64_t node_id = node->node->node_id; - struct ETCP_CONN* existing = inst->connections; - while (existing) { - if (existing->peer_node_id == node_id) { - struct ETCP_LINK* l = existing->links; - while (l) { if (l->link_state == 3) break; l = l->next; } - if (l) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx already connected, delivering immediately", - (unsigned long long)node_id); - if (flags & ETCP_CONNECT_EARLY) cb(arg, existing, ETCP_CONNECT_EARLY); - if (flags & ETCP_CONNECT_LATE) cb(arg, existing, ETCP_CONNECT_LATE); - if ((flags & ETCP_CONNECT_BGP_READY) && existing->routing_exchange_active >= 3) - cb(arg, existing, ETCP_CONNECT_BGP_READY); - return 0; + // Check if already connected via indexed queue + { + struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id); + if (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* existing = ce->conn; + if (existing && existing->peer_node_id == node_id) { + struct ETCP_LINK* l = existing->links; + while (l) { if (l->link_state == 3) break; l = l->next; } + if (l) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx already connected, delivering immediately", + (unsigned long long)node_id); + if (flags & ETCP_CONNECT_EARLY) cb(arg, existing, ETCP_CONNECT_EARLY); + if (flags & ETCP_CONNECT_LATE) cb(arg, existing, ETCP_CONNECT_LATE); + if ((flags & ETCP_CONNECT_BGP_READY) && existing->routing_exchange_active >= 3) + cb(arg, existing, ETCP_CONNECT_BGP_READY); + return 0; + } } } - existing = existing->next; } struct ETCP_CONNECT* ctx = connect_find(inst, node_id); @@ -241,7 +245,6 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node, struct ETCP_CONN* conn = etcp_connection_create(inst, NULL); if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] failed to create ETCP_CONN for node 0x%016llx", (unsigned long long)node_id); cb(arg, NULL, 0); return -1; } - conn->next = inst->connections; inst->connections = conn; inst->connections_count++; if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] sc_init_ctx failed for node 0x%016llx", (unsigned long long)node_id); etcp_connection_close(conn); cb(arg, NULL, 0); return -1; diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 22d862f9..dd977e97 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -32,18 +32,17 @@ // TCP server: on new incoming connection → create minimal ETCP_CONN static void tcp_server_on_link(struct stcp_link *link, void *arg) { struct UTUN_INSTANCE *inst = (struct UTUN_INSTANCE *)arg; - struct ETCP_CONN *conn = u_calloc(1, sizeof(struct ETCP_CONN)); + struct ETCP_CONN *conn = etcp_connection_create(inst, NULL); if (!conn) return; - conn->instance = inst; conn->transport_link = link; snprintf(conn->log_name, sizeof(conn->log_name), "tcp-[%p]", (void*)link); - conn->next = inst->connections; - inst->connections = conn; - inst->connections_count++; struct etcp_cbk_entry* cbe = inst->new_conn_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(conn, cbe->arg); cbe = n; } { struct etcp_cbk_entry* rcb = conn->ready_cbks; while (rcb) { struct etcp_cbk_entry* n = rcb->next; rcb->fn(conn, rcb->arg); rcb = n; } } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server new conn=%p total=%d", (void*)conn, inst->connections_count); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server new conn=%p total=%d pending=%d", + (void*)conn, + queue_entry_count(inst->connections), + queue_entry_count(inst->connections)); } // Forward declaration @@ -1464,8 +1463,7 @@ static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_D // DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Received INIT_RESPONSE from server_node_id=%llu, mtu=%d", (unsigned long long)server_node_id, link->mtu); - link->etcp->peer_node_id = server_node_id; // If not set - etcp_update_log_name(link->etcp); // Update log_name with peer_node_id + etcp_conn_set_peer_node_id(link->etcp, server_node_id); link->initialized = 1;// получен init response (client) link->link_state = 3; // connected @@ -1678,11 +1676,11 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { req->src_ipv4[0], req->src_ipv4[1], req->src_ipv4[2], req->src_ipv4[3], src_port, session_id, sockaddr_storage_to_str(&addr).str); - struct ETCP_CONN* conn=e_sock->instance->connections; - while (conn) {// ищем есть ли подключение к этому пиру - if (conn->peer_node_id==peer_id) break; - conn=conn->next; - } + struct ETCP_CONN* conn = NULL; + { + struct ll_entry* e = queue_find_data_by_index(e_sock->instance->connections, (const uint8_t*)&peer_id); + if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; conn = ce->conn; } + } int new_conn=0; if (!conn || conn->peer_node_id!=peer_id) {// создаём новое подключение [new etcp] @@ -1690,14 +1688,11 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { conn=etcp_connection_create(e_sock->instance,""); if (!conn) { errorcode=55; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "failed to create connection"); goto ec_fr; } memcpy(&conn->crypto_ctx, &sc, sizeof(sc)); - conn->peer_node_id=peer_id; - etcp_update_log_name(conn); + etcp_conn_set_peer_node_id(conn, peer_id); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "New connection received on socket %s: log_name=%s peer_id=%lu peer:%s", e_sock->name, conn->log_name, (unsigned long)peer_id, sockaddr_storage_to_str(&addr).str); - DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "New connection from %s peer_id=%ld etcp=%p", ip_to_str(&addr, addr.ss_family).str, peer_id, conn); - conn->next = e_sock->instance->connections; - e_sock->instance->connections = conn; - e_sock->instance->connections_count++; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Added incoming connection %p to instance, total count: %d", conn, e_sock->instance->connections_count); + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "New connection from %s peer_id=%ld etcp=%p total=%d", + ip_to_str(&addr, addr.ss_family).str, peer_id, conn, + queue_entry_count(e_sock->instance->connections)); } else {// check keys если существующее подключение if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) { errorcode=5; DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "peer key mismatch for node %016llx", (unsigned long long)peer_id); goto ec_fr; }// коллизия - peer id совпал а ключи разные. @@ -2065,8 +2060,8 @@ int init_connections(struct UTUN_INSTANCE* instance) { // Initialize clients - create outgoing connections struct CFG_CLIENT* client = config->clients; - DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, connections_count=%d", - instance, config, config ? config->clients : NULL, instance ? instance->connections_count : -1); + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, total_conns=%d", + instance, config, config ? config->clients : NULL, instance ? queue_entry_count(instance->connections) : -1); while (client) { // Check if client has required configuration DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Client %s - keepalive=%d, links=%p, peer_key_len=%zu", @@ -2196,22 +2191,18 @@ int init_connections(struct UTUN_INSTANCE* instance) { client_link = client_link->next; } - etcp_conn->next = instance->connections; - instance->connections = etcp_conn; - instance->connections_count++; - - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Added connection %p to instance, total count: %d", etcp_conn, instance->connections_count); client = client->next; } // If there are clients configured but no connections created, that's an error // If there are no clients (server-only mode), 0 connections is OK (server will accept incoming) - if (instance->connections_count == 0 && config->clients != NULL) { + int total_conns = queue_entry_count(instance->connections); + if (total_conns == 0 && config->clients != NULL) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Clients configured but no connections initialized"); return -1; } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized %d connections", instance->connections_count); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized %d connections", total_conns); // Return 1 if there was a partial socket initialization error if (socket_result == 1) return 1; return 0; diff --git a/src/etcp_dump.c b/src/etcp_dump.c index 9cb372a6..d3f1249d 100644 --- a/src/etcp_dump.c +++ b/src/etcp_dump.c @@ -233,13 +233,16 @@ void etcp_dump_all_conns(struct UTUN_INSTANCE* instance) { } DUMP_HDR("=== DUMP ALL CONNS for instance node=0x%llx ===", (unsigned long long)instance->node_id); - struct ETCP_CONN* conn = instance->connections; + struct ll_entry* entry = instance->connections->head; int idx = 0; - while (conn) { - DUMP_HDR("--- CONN %d ---", idx); - etcp_dump_conn_state(conn); - conn = conn->next; - idx++; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce && ce->conn) { + DUMP_HDR("--- CONN %d ---", idx); + etcp_dump_conn_state(ce->conn); + idx++; + } + entry = entry->next; } if (idx == 0) DUMP_HDR("--- NO CONNECTIONS ---"); DUMP_HDR("=== END DUMP ALL ==="); diff --git a/src/ntp_node_time.c b/src/ntp_node_time.c index d1945cdc..2080f2f0 100644 --- a/src/ntp_node_time.c +++ b/src/ntp_node_time.c @@ -152,9 +152,10 @@ void ntp_node_sync_peers(struct UTUN_INSTANCE* inst) { if (!inst || !ntp_time_is_synced(inst)) return; int sent = 0; - for (struct ETCP_CONN* conn = inst->connections; conn; conn = conn->next) { - if (!conn->initialized || !conn->links_up) continue; - if (send_time_sync(inst, conn->peer_node_id) == 0) sent++; + for (struct ll_entry* entry = inst->connections->head; entry; entry = entry->next) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (!ce->conn->initialized || !ce->conn->links_up) continue; + if (send_time_sync(inst, ce->conn->peer_node_id) == 0) sent++; } DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: broadcast TIME_SYNC to %d peers", sent); diff --git a/src/topo_group.c b/src/topo_group.c index fa788172..737a9cab 100644 --- a/src/topo_group.c +++ b/src/topo_group.c @@ -366,8 +366,8 @@ void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { topo_group_add_to_senders(group, conn); if (peer_nq && peer_nq->alien) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "peer 0x%016llx is alien, skipping route exchange", (unsigned long long)conn->peer_node_id); return; } - struct ETCP_CONN* c = conn->instance->connections; - while (c) { struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->conn && l->nat_check_status < NAT_CHECK_IN_PROGRESS) topo_group_start_link_nat_check(group, l); l = l->next; } c = c->next; } + struct ll_entry* entry = conn->instance->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized && l->conn && l->nat_check_status < NAT_CHECK_IN_PROGRESS) topo_group_start_link_nat_check(group, l); l = l->next; } entry = entry->next; } topo_group_send_table_request(group, conn); } @@ -839,7 +839,7 @@ void topo_group_send_nat_check_req(struct ETCP_CONN* conn, uint8_t socket_id) { void topo_group_request_nat_check_all(struct TOPO_GROUP* group) { if (!group || !group->instance) return; int count = 0; - struct ETCP_CONN* conn = group->instance->connections; while (conn) { struct ETCP_LINK* link = conn->links; while (link) { if (link->nat_check_status < NAT_CHECK_IN_PROGRESS) topo_group_start_link_nat_check(group, link); count++; link = link->next; } conn = conn->next; } + struct ll_entry* entry = group->instance->connections->head; while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; struct ETCP_LINK* link = ce->conn->links; while (link) { if (link->nat_check_status < NAT_CHECK_IN_PROGRESS) topo_group_start_link_nat_check(group, link); count++; link = link->next; } entry = entry->next; } } void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id) { diff --git a/src/utun_instance.c b/src/utun_instance.c index f41417de..54025db6 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -106,6 +106,9 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Failed to create networks queue"); return -1; } + // Initialize connection queue + instance->connections = queue_new(ua, 256, 0, 8, "connections"); + if (!instance->connections) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create connections queue"); return -1; } struct CFG_NETWORK *net = config->networks; while (net) { struct ll_entry *entry = queue_entry_new(sizeof(struct NETWORK_ENTRY)); @@ -397,17 +400,17 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { instance->etcp_sockets = NULL; DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] ETCP sockets cleanup complete"); - // Safe cleanup of ETCP connections with NULL pointer checks - struct ETCP_CONN* conn = instance->connections; - while (conn) { - struct ETCP_CONN* next = conn->next; - DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Closing connection %p", conn); - if (conn) { // Дополнительная проверка на случай поврежденного списка - etcp_connection_close(conn); // Закрыть соединение (с проверкой NULL внутри) + // Cleanup ETCP connections (phase 1 detach + deferred phase 2 via call_soon) + { + struct ll_entry* entry = instance->connections->head; + while (entry) { + struct ll_entry* next = entry->next; + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce && ce->conn) etcp_connection_close(ce->conn); + entry = next; } - conn = next; } - instance->connections = NULL; + queue_free(instance->connections); instance->connections = NULL; DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] ETCP connections cleanup complete"); struct PING_CONTEXT* p = instance->pending_pings; @@ -578,7 +581,7 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) { return -1; } - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully, count=%d", instance->connections_count); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully, count=%d", queue_entry_count(instance->connections)); // Initialize control server if configured if (instance->config->global.control_sock.ss_family != 0) { @@ -654,16 +657,15 @@ void utun_instance_diagnose_leaks(struct UTUN_INSTANCE *instance, const char *ph } // Подсчёт ETCP соединений - struct ETCP_CONN *conn = instance->connections; - while (conn) { - report.etcp_connections_count++; - // Подсчёт линков в соединениях - struct ETCP_LINK *link = conn->links; - while (link) { - report.etcp_links_count++; - link = link->next; + { + struct ll_entry* entry = instance->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + report.etcp_connections_count++; + struct ETCP_LINK *link = ce->conn->links; + while (link) { report.etcp_links_count++; link = link->next; } + entry = entry->next; } - conn = conn->next; } DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[LEAK DIAGNOSIS] Phase: %s", phase); @@ -788,9 +790,12 @@ struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struc // Clients [client:] and links (link=) for (struct CFG_CLIENT *nc = new_config->clients; nc; nc = nc->next) { struct ETCP_CONN *conn = NULL; - for (struct ETCP_CONN *c = instance->connections; c; c = c->next) { - if (strcmp(nc->name, c->name) == 0) { - conn = c; break; + { + struct ll_entry* entry = instance->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (strcmp(nc->name, ce->conn->name) == 0) { conn = ce->conn; break; } + entry = entry->next; } } if (!conn) { diff --git a/src/utun_instance.h b/src/utun_instance.h index 1e705669..f59dd506 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -52,6 +52,12 @@ struct NETWORK_ENTRY { char name[64]; }; +// Queue entry for instance->connections (indexed by peer_node_id) +struct conn_queue_entry { + uint64_t peer_node_id; + struct ETCP_CONN* conn; +}; + // uTun instance configuration struct UTUN_INSTANCE { // Identification @@ -87,9 +93,8 @@ struct UTUN_INSTANCE { // State int running; - // Connections (список всех подключений для instance) - struct ETCP_CONN* connections;// linked-list - int connections_count; // Number of connections + // Connections (очередь всех подключений для instance) + struct ll_queue* connections; // indexed by peer_node_id (0=pending), data=conn_queue_entry // Callback chain for new ETCP connections struct etcp_cbk_entry* new_conn_cbks; @@ -135,7 +140,6 @@ 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/tests/bbr_integration/test_bbr_integration.c b/tests/bbr_integration/test_bbr_integration.c index 375db676..f546b107 100644 --- a/tests/bbr_integration/test_bbr_integration.c +++ b/tests/bbr_integration/test_bbr_integration.c @@ -187,7 +187,9 @@ static void send_wtr_cb(struct ll_queue* q, void* arg) { static void send_burst(struct test_ctx* ctx) { if (ctx->test_done) return; - struct ETCP_CONN* conn = ctx->sender->connections; + struct ETCP_CONN* conn = NULL; + if (ctx->sender->connections && ctx->sender->connections->head) + conn = ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn; if (!conn || !conn->initialized) return; if (conn->normalizer && conn->normalizer->input) { @@ -241,9 +243,11 @@ static void print_metrics(struct test_ctx* ctx) { double pace_kbps = 0; double rtt_ms = 0; - if (ctx->sender->connections && ctx->sender->connections->links) { - struct ETCP_LINK* l = ctx->sender->connections->links; - struct ETCP_CONN* c = ctx->sender->connections; + struct ETCP_CONN* c = NULL; + if (ctx->sender->connections && ctx->sender->connections->head) + c = ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn; + if (c && c->links) { + struct ETCP_LINK* l = c->links; lk_status = l->link_status ? " UP" : " DN"; lk_reinit = c->reinit_count; lk_rtrns = l->total_retransmissions; @@ -376,11 +380,14 @@ int main(void) { /* Wait for connection */ printf("Waiting for connection...\n"); + struct ETCP_CONN* conn = NULL; uint64_t t0 = now_us(); while ((now_us() - t0) < 10000000ULL) { uasync_poll(ctx.ua, 1); - if (ctx.sender->connections && ctx.sender->connections->links) { - struct ETCP_LINK* l = ctx.sender->connections->links; + if (ctx.sender->connections && ctx.sender->connections->head) + conn = ((struct conn_queue_entry*)ctx.sender->connections->head->data)->conn; + if (conn && conn->links) { + struct ETCP_LINK* l = conn->links; if (l->initialized && l->link_status == 1) break; } } @@ -388,9 +395,9 @@ int main(void) { t0 = now_us(); while ((now_us() - t0) < 500000ULL) uasync_poll(ctx.ua, 1); - if (ctx.sender->connections && ctx.sender->connections->normalizer && ctx.sender->connections->normalizer->input) { - queue_set_threshold(ctx.sender->connections->normalizer->input, 0, 0); - queue_set_waiter_defer(ctx.sender->connections->normalizer->input, 1); + if (conn && conn->normalizer && conn->normalizer->input) { + queue_set_threshold(conn->normalizer->input, 0, 0); + queue_set_waiter_defer(conn->normalizer->input, 1); } /* Start test */ @@ -416,10 +423,15 @@ int main(void) { const struct dummynet_stats* ds = dummynet_get_stats(ctx.dn, DUMMYNET_FORWARD); const char* lks = "?"; uint32_t rein = 0, rtrns = 0; - if (ctx.sender->connections && ctx.sender->connections->links) { - lks = ctx.sender->connections->links->link_status ? "UP" : "DN"; - rein = ctx.sender->connections->reinit_count; - rtrns = ctx.sender->connections->links->total_retransmissions; + { + struct ETCP_CONN* c = NULL; + if (ctx.sender->connections && ctx.sender->connections->head) + c = ((struct conn_queue_entry*)ctx.sender->connections->head->data)->conn; + if (c && c->links) { + lks = c->links->link_status ? "UP" : "DN"; + rein = c->reinit_count; + rtrns = c->links->total_retransmissions; + } } printf(" t=%lus Lk=%s Rein=%u Rtrns=%u recv=%lu dn_tx=%llu dn_lost=%llu\n", (unsigned long)elapsed, lks, rein, rtrns, (unsigned long)ctx.bytes_received, @@ -455,31 +467,35 @@ int main(void) { (unsigned long long)ds_bk->recv, (unsigned long long)ds_bk->sent, (unsigned long long)ds_bk->lost, (unsigned long long)ds_bk->dropped); - if (ctx.sender->connections) { - struct ETCP_CONN* c = ctx.sender->connections; - printf("ETCP: reinit=%u reset=%u links_up=%u\n", - c->reinit_count, c->reset_count, c->links_up); - if (c->links) { - struct ETCP_LINK* l = c->links; - int lnum = 0; - while (l) { - lnum++; - printf(" Link%d: status=%s state=%u rKeep=%u sKeep=%u retrans=%lu" - " inflight=%u/%u lim=%u kB rtt=%u ms\n", - lnum, l->link_status ? "UP" : "DN", l->link_state, - l->recv_keepalive, l->remote_keepalive, - (unsigned long)l->total_retransmissions, - l->inflight_bytes, l->inflight_packets, l->inflight_lim_bytes / 1024, - l->rtt_last); - if (l->bbr) - printf(" BBR: mode=%d cycle=%d full_bw=%d loss_rnd=%d" - " minrtt=%.1fms bw_hi=%.0fK bw_lo=%.0fK pace=%.0fK infl_hi=%.1fK infl_lo=%.1fK\n", - l->bbr->mode, l->bbr->cycle_idx, l->bbr->full_bw_reached, l->bbr->loss_in_round, - (double)l->bbr->min_rtt_us / 1000.0, - (double)l->bbr->bw_hi[0] / 2097.152, (double)l->bbr->bw_lo / 2097.152, - (double)l->bbr_pacing_rate / 125.0, - (double)l->bbr->inflight_hi / 1024.0, (double)l->bbr->inflight_lo / 1024.0); - l = l->next; + { + struct ETCP_CONN* c = NULL; + if (ctx.sender->connections && ctx.sender->connections->head) + c = ((struct conn_queue_entry*)ctx.sender->connections->head->data)->conn; + if (c) { + printf("ETCP: reinit=%u reset=%u links_up=%u\n", + c->reinit_count, c->reset_count, c->links_up); + if (c->links) { + struct ETCP_LINK* l = c->links; + int lnum = 0; + while (l) { + lnum++; + printf(" Link%d: status=%s state=%u rKeep=%u sKeep=%u retrans=%lu" + " inflight=%u/%u lim=%u kB rtt=%u ms\n", + lnum, l->link_status ? "UP" : "DN", l->link_state, + l->recv_keepalive, l->remote_keepalive, + (unsigned long)l->total_retransmissions, + l->inflight_bytes, l->inflight_packets, l->inflight_lim_bytes / 1024, + l->rtt_last); + if (l->bbr) + printf(" BBR: mode=%d cycle=%d full_bw=%d loss_rnd=%d" + " minrtt=%.1fms bw_hi=%.0fK bw_lo=%.0fK pace=%.0fK infl_hi=%.1fK infl_lo=%.1fK\n", + l->bbr->mode, l->bbr->cycle_idx, l->bbr->full_bw_reached, l->bbr->loss_in_round, + (double)l->bbr->min_rtt_us / 1000.0, + (double)l->bbr->bw_hi[0] / 2097.152, (double)l->bbr->bw_lo / 2097.152, + (double)l->bbr_pacing_rate / 125.0, + (double)l->bbr->inflight_hi / 1024.0, (double)l->bbr->inflight_lo / 1024.0); + l = l->next; + } } } } @@ -505,8 +521,13 @@ int main(void) { } /* Cleanup */ - if (ctx.sender->connections && ctx.sender->connections->normalizer && ctx.sender->connections->normalizer->input) - queue_waiter_cancel(ctx.sender->connections->normalizer->input, &ctx.waiter); + { + struct ETCP_CONN* cc = NULL; + if (ctx.sender->connections && ctx.sender->connections->head) + cc = ((struct conn_queue_entry*)ctx.sender->connections->head->data)->conn; + if (cc && cc->normalizer && cc->normalizer->input) + queue_waiter_cancel(cc->normalizer->input, &ctx.waiter); + } dummynet_destroy(ctx.dn); utun_instance_destroy(ctx.sender); utun_instance_destroy(ctx.receiver); diff --git a/tests/test_bgp_route_exchange.c b/tests/test_bgp_route_exchange.c index 21da299e..0bf4a9db 100644 --- a/tests/test_bgp_route_exchange.c +++ b/tests/test_bgp_route_exchange.c @@ -195,17 +195,17 @@ static void cleanup_temp_configs(void) { } static int has_initialized_link(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { struct ETCP_LINK* l = conn->links; while (l) { if (l->initialized) return 1; l = l->next; } conn = conn->next; } + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized) return 1; l = l->next; } entry = entry->next; } return 0; } static int count_initialized_links(struct UTUN_INSTANCE* inst) { int n = 0; - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { struct ETCP_LINK* l = conn->links; while (l) { if (l->initialized) n++; l = l->next; } conn = conn->next; } + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized) n++; l = l->next; } entry = entry->next; } return n; } @@ -221,10 +221,11 @@ static int check_learned_route(struct UTUN_INSTANCE* inst, uint32_t network, uin } static struct ETCP_LINK* find_client_link(struct UTUN_INSTANCE* inst, int idx) { - if (!inst) return NULL; - struct ETCP_CONN* conn = inst->connections; - while (conn) { - struct ETCP_LINK* l = conn->links; + if (!inst || !inst->connections) return NULL; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->is_server == 0) { if (idx == 0) return l; @@ -232,7 +233,7 @@ static struct ETCP_LINK* find_client_link(struct UTUN_INSTANCE* inst, int idx) { } l = l->next; } - conn = conn->next; + entry = entry->next; } return NULL; } diff --git a/tests/test_bgp_triangle.c b/tests/test_bgp_triangle.c index 5c3ea308..01eceb35 100644 --- a/tests/test_bgp_triangle.c +++ b/tests/test_bgp_triangle.c @@ -162,26 +162,28 @@ static int route_count_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { static int count_initialized_links(struct UTUN_INSTANCE* inst) { int n = 0; - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { - struct ETCP_LINK* l = conn->links; + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized) n++; l = l->next; } - conn = conn->next; + entry = entry->next; } return n; } static struct ETCP_LINK* find_client_link(struct UTUN_INSTANCE* inst, int idx) { - if (!inst) return NULL; - struct ETCP_CONN* conn = inst->connections; - while (conn) { - struct ETCP_LINK* l = conn->links; + if (!inst || !inst->connections) return NULL; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->is_server == 0) { if (idx == 0) return l; idx--; } l = l->next; } - conn = conn->next; + entry = entry->next; } return NULL; } diff --git a/tests/test_conn_mgr.c b/tests/test_conn_mgr.c index aa77239d..070d87d0 100644 --- a/tests/test_conn_mgr.c +++ b/tests/test_conn_mgr.c @@ -56,8 +56,11 @@ static char* gv(const char* p, const char* k) { free_config(c); return r; } static int lks(struct UTUN_INSTANCE* i) { - int n = 0; struct ETCP_CONN* c = i->connections; - while (c) { struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } c = c->next; } + int n = 0; struct ll_entry* entry = i->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* c = ce->conn; + struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } + entry = entry->next; } return n; } static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 2; } diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index 8f4e447b..e60d3845 100644 --- a/tests/test_db_sync.c +++ b/tests/test_db_sync.c @@ -151,8 +151,11 @@ static void sleep_tb(int tb) { // ---- Condition functions ---- static int cond_links_init(void) { if (!inst_a || !inst_b) return 0; - struct ETCP_CONN* ca = inst_a->connections; - while (ca) { struct ETCP_LINK* l = ca->links; while (l) { if (l->initialized) return 1; l = l->next; } ca = ca->next; } + struct ll_entry* entry = inst_a->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* ca = ce->conn; + struct ETCP_LINK* l = ca->links; while (l) { if (l->initialized) return 1; l = l->next; } + entry = entry->next; } return 0; } diff --git a/tests/test_etcp_100_packets.c b/tests/test_etcp_100_packets.c index 506bf9cd..d069d688 100644 --- a/tests/test_etcp_100_packets.c +++ b/tests/test_etcp_100_packets.c @@ -139,15 +139,17 @@ static void generate_packet_data(int seq, uint8_t* buffer, int size) { // Check if connection is established static int is_connection_established(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } @@ -156,7 +158,7 @@ static int is_connection_established(struct UTUN_INSTANCE* inst) { static void send_packets_fwd(void) { if (!client_instance || packets_sent_fwd >= TOTAL_PACKETS) return; - struct ETCP_CONN* conn = client_instance->connections; + struct ETCP_CONN* conn = (client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; if (!conn || !conn->input_queue) return; // Start timing on first packet @@ -193,7 +195,7 @@ static void send_packets_fwd(void) { static void send_packets_back(void) { if (!server_instance || packets_sent_back >= TOTAL_PACKETS) return; - struct ETCP_CONN* conn = server_instance->connections; + struct ETCP_CONN* conn = (server_instance->connections && server_instance->connections->head) ? ((struct conn_queue_entry*)server_instance->connections->head->data)->conn : NULL; if (!conn || !conn->input_queue) return; // Start timing on first packet @@ -230,7 +232,7 @@ static void send_packets_back(void) { static void check_received_packets_fwd(void) { if (!server_instance) return; - struct ETCP_CONN* conn = server_instance->connections; + struct ETCP_CONN* conn = (server_instance->connections && server_instance->connections->head) ? ((struct conn_queue_entry*)server_instance->connections->head->data)->conn : NULL; if (!conn || !conn->output_queue) return; // Disable routing callback to keep packets in output_queue for test verification @@ -263,7 +265,7 @@ static void check_received_packets_fwd(void) { static void check_received_packets_back(void) { if (!client_instance) return; - struct ETCP_CONN* conn = client_instance->connections; + struct ETCP_CONN* conn = (client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; if (!conn || !conn->output_queue) return; // Disable routing callback to keep packets in output_queue for test verification diff --git a/tests/test_etcp_api.c b/tests/test_etcp_api.c index b74ea4b8..56bfc786 100644 --- a/tests/test_etcp_api.c +++ b/tests/test_etcp_api.c @@ -280,8 +280,10 @@ static void client_recv_callback(struct ETCP_CONN* conn, struct ll_entry* entry) // Check if connection is established and crypto session key is ready static int is_connection_established(struct UTUN_INSTANCE* inst) { if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { // Проверяем что линк инициализирован И session_key установлен (не нулевой) @@ -300,15 +302,17 @@ static int is_connection_established(struct UTUN_INSTANCE* inst) { } link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } // Get first connection from instance static struct ETCP_CONN* get_first_connection(struct UTUN_INSTANCE* inst) { - if (!inst) return NULL; - return inst->connections; + if (!inst || !inst->connections) return NULL; + if (!inst->connections->head) return NULL; + struct conn_queue_entry* ce = (struct conn_queue_entry*)inst->connections->head->data; + return ce->conn; } // Send packets from client to server (forward direction) via etcp_send diff --git a/tests/test_etcp_congestion.c b/tests/test_etcp_congestion.c index 36821fad..eea561b4 100644 --- a/tests/test_etcp_congestion.c +++ b/tests/test_etcp_congestion.c @@ -187,8 +187,8 @@ static void send_timer_cb(void* arg) { static void send_burst(struct test_ctx* ctx) { if (ctx->test_done) return; if (!ctx->sender || !ctx->sender->connections) return; - struct ETCP_CONN* conn = ctx->sender->connections; - if (!conn->initialized) return; + struct ETCP_CONN* conn = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; + if (!conn || !conn->initialized) return; int sent = 0; while (sent < 64) { @@ -232,7 +232,8 @@ static void log_metrics(struct test_ctx* ctx) { ctx->last_bytes_recv = ctx->bytes_received; if (!ctx->sender || !ctx->sender->connections) return; - struct ETCP_CONN* sconn = ctx->sender->connections; + struct ETCP_CONN* sconn = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; + if (!sconn) return; struct ETCP_LINK* link = sconn->links; int idx = 0; while (link && idx < 2) { @@ -343,10 +344,11 @@ int main(void) { int links_ready = 0; while ((now_us() - t0) < 10000000ULL) { uasync_poll(ctx.ua, 1); - if (ctx.sender->connections && ctx.sender->connections->links) { - struct ETCP_LINK* l = ctx.sender->connections->links; + { struct ETCP_CONN* c = (ctx.sender->connections && ctx.sender->connections->head) ? ((struct conn_queue_entry*)ctx.sender->connections->head->data)->conn : NULL; + if (c && c->links) { + struct ETCP_LINK* l = c->links; if (l->initialized && l->link_status == 1) { links_ready = 1; break; } - } + } } } printf("Links ready, waiting stabilize...\n"); t0 = now_us(); @@ -383,7 +385,7 @@ int main(void) { if (elapsed >= TEST_DURATION_MS) ctx.test_done = 1; if (elapsed - last_progress >= 1000) { last_progress = elapsed; - int l1 = ctx.sender->connections && ctx.sender->connections->links ? 1 : 0; + int l1 = (ctx.sender->connections && ctx.sender->connections->head) ? 1 : 0; printf(" t=%lus sent=%lu recv=%lu links=%d\n", (unsigned long)elapsed, (unsigned long)ctx.bytes_sent, (unsigned long)ctx.bytes_received, l1); } } diff --git a/tests/test_etcp_connect.c b/tests/test_etcp_connect.c index eb089d88..4b99549d 100644 --- a/tests/test_etcp_connect.c +++ b/tests/test_etcp_connect.c @@ -45,8 +45,11 @@ static char* gv(const char* p, const char* k) { free_config(c); return r; } static int lks(struct UTUN_INSTANCE* i) { - int n = 0; struct ETCP_CONN* c = i->connections; - while (c) { struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } c = c->next; } + int n = 0; struct ll_entry* entry = i->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* c = ce->conn; + struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } + entry = entry->next; } return n; } static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 1; } diff --git a/tests/test_etcp_dummynet.c b/tests/test_etcp_dummynet.c index 77888be5..4bf3ed5d 100644 --- a/tests/test_etcp_dummynet.c +++ b/tests/test_etcp_dummynet.c @@ -81,8 +81,10 @@ static void send_pkt(void* arg) { if (sending_done || (now_us() - start_us)/1000 >= SEND_MS) { sending_done = 1; return; } - if (!client || !client->connections) return; - + if (!client || !client->connections || !client->connections->head) return; + struct ETCP_CONN* conn = ((struct conn_queue_entry*)client->connections->head->data)->conn; + if (!conn) return; + struct ll_entry* e = ll_alloc_lldgram(sizeof(struct test_pkt)); if (!e) return; struct test_pkt* pkt = (struct test_pkt*)e->dgram; @@ -91,7 +93,7 @@ static void send_pkt(void* arg) { pkt->crc = crc32_calc(pkt->data, PAYLOAD_SIZE); e->len = sizeof(struct test_pkt); - if (etcp_send(client->connections, e) == 0) pkts_sent++; + if (etcp_send(conn, e) == 0) pkts_sent++; else queue_entry_free(e); /* Schedule next packet with 5ms delay for rate limiting (~200 pkt/sec) @@ -229,7 +231,7 @@ int main(void) { /* Check if client has an established connection with ready links */ if (client && client->connections) { - struct ETCP_CONN* conn = client->connections; + struct ETCP_CONN* conn = (client->connections && client->connections->head) ? ((struct conn_queue_entry*)client->connections->head->data)->conn : NULL; /* Check if peer_node_id is set (non-zero) and has links */ if (conn->peer_node_id != 0 && conn->links != NULL) { struct ETCP_LINK* link = conn->links; diff --git a/tests/test_etcp_ping.c b/tests/test_etcp_ping.c index 11e28031..03954cdc 100644 --- a/tests/test_etcp_ping.c +++ b/tests/test_etcp_ping.c @@ -97,27 +97,31 @@ static void monitor_connections(void* arg) { int server_links = 0; int client_links = 0; if (server_instance) { - struct ETCP_CONN* conn = server_instance->connections; - while (conn) { + struct ll_entry* e = server_instance->connections->head; + while (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { server_links++; if (link->initialized) test_completed = 1; link = link->next; } - conn = conn->next; + e = e->next; } } if (client_instance) { - struct ETCP_CONN* conn = client_instance->connections; - while (conn) { + struct ll_entry* e = client_instance->connections->head; + while (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { client_links++; if (link->initialized) test_completed = 1; link = link->next; } - conn = conn->next; + e = e->next; } } if (server_links > 0 && client_links > 0) test_completed = 1; @@ -182,7 +186,7 @@ int main(void) { goto cleanup; } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Links established, running ping test"); - struct ETCP_CONN* conn = client_instance->connections; + struct ETCP_CONN* conn = (client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; if (conn && conn->links) { struct ETCP_LINK* link = conn->links; uint8_t peer_pubkey[SC_PUBKEY_SIZE]; diff --git a/tests/test_etcp_reconnect.c b/tests/test_etcp_reconnect.c index a4107d59..5807858c 100644 --- a/tests/test_etcp_reconnect.c +++ b/tests/test_etcp_reconnect.c @@ -114,19 +114,22 @@ static void cleanup_temp_configs(void) { } static int is_connection_established(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { - if (conn->initialized) return 1; - conn = conn->next; + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce->conn->initialized) return 1; + entry = entry->next; } return 0; } static void drain_received(int count_flag) { - if (!server_instance) return; - struct ETCP_CONN* conn = server_instance->connections; - while (conn) { + if (!server_instance || !server_instance->connections) return; + struct ll_entry* entry = server_instance->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; if (conn->output_queue) { queue_set_callback(conn->output_queue, NULL, NULL); struct ETCP_FRAGMENT* pkt; @@ -136,13 +139,13 @@ static void drain_received(int count_flag) { queue_entry_free((struct ll_entry*)pkt); } } - conn = conn->next; + entry = entry->next; } } static void send_packets(void) { if (!client_instance || packets_sent >= TOTAL_PACKETS) return; - struct ETCP_CONN* conn = client_instance->connections; + struct ETCP_CONN* conn = (client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; if (!conn || !conn->input_queue) return; while (packets_sent < TOTAL_PACKETS) { if (queue_entry_count(conn->input_queue) >= MAX_QUEUE_SIZE) break; @@ -160,7 +163,7 @@ static void monitor(void* arg) { int conn_ok = is_connection_established(client_instance); static int last_reinit_count = -1; int cur_reinit = 0; - struct ETCP_CONN* conn = client_instance ? client_instance->connections : NULL; + struct ETCP_CONN* conn = (client_instance && client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; if (conn) cur_reinit = conn->reinit_count; switch (phase) { @@ -208,8 +211,8 @@ static void monitor(void* arg) { } printf("Server recreated (node_id=%llx)\n", (unsigned long long)server_instance->node_id); // Force client to reinit — otherwise conn_ok stays true on old initialized=1 - { struct ETCP_CONN* c = client_instance->connections; - while (c) { etcp_conn_reinit(c); c = c->next; } } + { struct ll_entry* e = client_instance->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; etcp_conn_reinit(ce->conn); e = e->next; } } restart_action_done = 1; } drain_received(0); diff --git a/tests/test_etcp_reinit_inflight.c b/tests/test_etcp_reinit_inflight.c index 574f50d6..ef68d228 100644 --- a/tests/test_etcp_reinit_inflight.c +++ b/tests/test_etcp_reinit_inflight.c @@ -133,17 +133,20 @@ static void on_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { } static int links_initialized(struct UTUN_INSTANCE* inst) { - struct ETCP_CONN* conn = inst->connections; - while (conn) { + if (!inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } static int buffers_full(struct test_ctx* ctx) { - struct ETCP_CONN* conn = ctx->sender->connections; + struct ETCP_CONN* conn = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; if (!conn) return 0; struct ETCP_LINK* link = conn->links; return link && link->inflight_bytes >= link->inflight_lim_bytes @@ -152,7 +155,7 @@ static int buffers_full(struct test_ctx* ctx) { } static void send_one_packet(struct test_ctx* ctx) { - struct ETCP_CONN* conn = ctx->sender->connections; + struct ETCP_CONN* conn = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; if (!conn || !conn->initialized) return; if (ctx->total_packets_sent >= TOTAL_PACKETS) return; if (conn->input_queue && queue_entry_count(conn->input_queue) > 200) return; @@ -192,11 +195,12 @@ static void monitor(void* arg) { break; case 2: // fill buffers if (buffers_full(ctx)) { - printf(" buffers full: inflight=%u/%u wait_ack=%d sent=%d\n", - ctx->sender->connections->links->inflight_bytes, - ctx->sender->connections->links->inflight_lim_bytes, - ctx->sender->connections->input_wait_ack->count, - ctx->total_packets_sent); + { struct ETCP_CONN* c = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; + printf(" buffers full: inflight=%u/%u wait_ack=%d sent=%d\n", + c ? c->links->inflight_bytes : 0, + c ? c->links->inflight_lim_bytes : 0, + c ? c->input_wait_ack->count : 0, + ctx->total_packets_sent); } ctx->phase = 3; } break; @@ -221,7 +225,7 @@ static void monitor(void* arg) { ctx->phase = 5; break; case 5: // wait delivery - conn = ctx->sender->connections; + conn = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; if (conn && conn->reinit_count > 0 && !ctx->sender_reinit_detected) { ctx->sender_reinit_detected = 1; printf(" reinit detected (count=%u)\n", conn->reinit_count); @@ -232,7 +236,7 @@ static void monitor(void* arg) { } break; case 6: // verify - conn = ctx->sender->connections; + conn = (ctx->sender->connections && ctx->sender->connections->head) ? ((struct conn_queue_entry*)ctx->sender->connections->head->data)->conn : NULL; if (conn && conn->reinit_count < 1) { printf("\n[FAIL] Reinit not triggered (reinit_count=%u)\n", conn->reinit_count); ctx->test_done = 2; diff --git a/tests/test_etcp_router.c b/tests/test_etcp_router.c index be97e9b5..3cc78608 100644 --- a/tests/test_etcp_router.c +++ b/tests/test_etcp_router.c @@ -205,8 +205,11 @@ static void write_configs(void) { // ======================== Helpers ======================== static int conn_established(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - for (struct ETCP_CONN* c = inst->connections; c; c = c->next) { + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* c = ce->conn; for (struct ETCP_LINK* l = c->links; l; l = l->next) { if (l->initialized && c->crypto_ctx.initialized) { int ok = 0; @@ -215,6 +218,7 @@ static int conn_established(struct UTUN_INSTANCE* inst) { if (ok) return 1; } } + entry = entry->next; } return 0; } diff --git a/tests/test_etcp_router_reconnect.c b/tests/test_etcp_router_reconnect.c index 9fb998fe..e130029d 100644 --- a/tests/test_etcp_router_reconnect.c +++ b/tests/test_etcp_router_reconnect.c @@ -186,9 +186,10 @@ static void state_step(void) { break; case ST_B_KILL: g_b->running=0; utun_instance_destroy(g_b); g_b=NULL; - { struct ETCP_CONN *c,*n; - for(c=g_a->connections;c;c=n){n=c->next;if(c->peer_node_id==node_b)etcp_connection_close(c);} - for(c=g_c->connections;c;c=n){n=c->next;if(c->peer_node_id==node_b)etcp_connection_close(c);} } + { struct ll_entry* e = queue_find_data_by_index(g_a->connections, (const uint8_t*)&node_b); + if (e) etcp_connection_close(((struct conn_queue_entry*)e->data)->conn); + e = queue_find_data_by_index(g_c->connections, (const uint8_t*)&node_b); + if (e) etcp_connection_close(((struct conn_queue_entry*)e->data)->conn); } g_phase_sent=0; g_bgp_ready_mask=0; g_state=ST_B_RESTART; break; case ST_B_RESTART: diff --git a/tests/test_etcp_simple_traffic.c b/tests/test_etcp_simple_traffic.c index 3495b152..7c9b66d3 100644 --- a/tests/test_etcp_simple_traffic.c +++ b/tests/test_etcp_simple_traffic.c @@ -107,22 +107,24 @@ static void cleanup_temp_configs(void) { } static int is_connection_established(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } static void send_test_packet(void) { if (!client_instance || packet_sent) return; - struct ETCP_CONN* conn = client_instance->connections; + struct ETCP_CONN* conn = (client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; if (!conn || !conn->input_queue) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "send_test_packet: no connection or input_queue"); return; @@ -138,7 +140,7 @@ static void send_test_packet(void) { static void check_packet_received(void) { if (!server_instance || packet_received) return; - struct ETCP_CONN* conn = server_instance->connections; + struct ETCP_CONN* conn = (server_instance->connections && server_instance->connections->head) ? ((struct conn_queue_entry*)server_instance->connections->head->data)->conn : NULL; if (!conn || !conn->output_queue) return; queue_set_callback(conn->output_queue, NULL, NULL); diff --git a/tests/test_etcp_two_instances.c b/tests/test_etcp_two_instances.c index 96304aed..f3b6b80d 100644 --- a/tests/test_etcp_two_instances.c +++ b/tests/test_etcp_two_instances.c @@ -125,8 +125,10 @@ static void monitor_connections(void* arg) { // Check server if (server_instance) { - struct ETCP_CONN* conn = server_instance->connections; - while (conn) { + struct ll_entry* entry = server_instance->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { server_links++; @@ -136,14 +138,16 @@ static void monitor_connections(void* arg) { link->is_server ? "server" : "client"); link = link->next; } - conn = conn->next; + entry = entry->next; } } // Check client if (client_instance) { - struct ETCP_CONN* conn = client_instance->connections; - while (conn) { + struct ll_entry* entry = client_instance->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { client_links++; @@ -203,7 +207,7 @@ static void monitor_connections(void* arg) { } link = link->next; } - conn = conn->next; + entry = entry->next; } } diff --git a/tests/test_icmp_proxy.c b/tests/test_icmp_proxy.c index 5077a207..20334ee0 100644 --- a/tests/test_icmp_proxy.c +++ b/tests/test_icmp_proxy.c @@ -114,12 +114,18 @@ static void monitor(void* arg) { if (g_test_phase == 0) { int cli_ok = 0, exit_ok = 0; - for (struct ETCP_CONN* c = cli->connections; c; c = c->next) + { struct ll_entry* e = cli->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* c = ce->conn; for (struct ETCP_LINK* l = c->links; l; l = l->next) if (l->initialized && c->crypto_ctx.initialized) cli_ok = 1; - for (struct ETCP_CONN* c = exit_node->connections; c; c = c->next) + e = e->next; } } + { struct ll_entry* e = exit_node->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* c = ce->conn; for (struct ETCP_LINK* l = c->links; l; l = l->next) if (l->initialized && c->crypto_ctx.initialized) exit_ok = 1; + e = e->next; } } if (cli_ok && exit_ok) { g_test_phase = 1; // Override client handler for test verification diff --git a/tests/test_ipv6_sockets.c b/tests/test_ipv6_sockets.c index f0b9d161..ded353a7 100644 --- a/tests/test_ipv6_sockets.c +++ b/tests/test_ipv6_sockets.c @@ -121,15 +121,17 @@ static void cleanup_temp_configs(void) { } static int is_connection_established(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } diff --git a/tests/test_nat_detection.c b/tests/test_nat_detection.c index f08a41ab..3d7986e7 100644 --- a/tests/test_nat_detection.c +++ b/tests/test_nat_detection.c @@ -149,27 +149,25 @@ static void cleanup_temp_configs(void) { } static int has_initialized_link(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } static struct ETCP_LINK* find_link_to_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { - if (!inst) return NULL; - struct ETCP_CONN* conn = inst->connections; - while (conn) { - if (conn->peer_node_id == node_id) return conn->links; - conn = conn->next; - } - return NULL; + if (!inst || !inst->connections) return NULL; + struct ll_entry* entry = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id); + return entry ? ((struct conn_queue_entry*)entry->data)->conn->links : NULL; } static void test_timeout_cb(void* arg) { @@ -335,11 +333,13 @@ int main(void) { } // Verify link nat_type on server - struct ETCP_CONN* conn_sc1 = inst_s->connections; - while (conn_sc1) { - if (conn_sc1->peer_node_id == NODE_ID_C1) break; - conn_sc1 = conn_sc1->next; - } + struct ETCP_CONN* conn_sc1 = NULL; + { struct ll_entry* e = inst_s->connections->head; + while (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + if (ce->conn->peer_node_id == NODE_ID_C1) { conn_sc1 = ce->conn; break; } + e = e->next; + } } if (!link_sc1 || link_sc1->nat_type != NAT_TYPE_EIM) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "FAIL: link nat_type not EIM on server"); goto cleanup; @@ -414,10 +414,11 @@ int main(void) { // 6. Separate STUN check via explicit route_ping_send_req_addr DEBUG_INFO(DEBUG_CATEGORY_BGP, "Performing explicit STUN ping from S via C2 to C1..."); struct ETCP_CONN* conn_sc2 = NULL; - struct ETCP_CONN* conn = inst_s->connections; - while (conn) { - if (conn->peer_node_id == NODE_ID_C2) { conn_sc2 = conn; break; } - conn = conn->next; + struct ll_entry* entry = inst_s->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce->conn->peer_node_id == NODE_ID_C2) { conn_sc2 = ce->conn; break; } + entry = entry->next; } if (!conn_sc2) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "FAIL: S has no connection to C2"); diff --git a/tests/test_nat_transport.c b/tests/test_nat_transport.c index 15a57720..3f9e7ffe 100644 --- a/tests/test_nat_transport.c +++ b/tests/test_nat_transport.c @@ -157,11 +157,12 @@ static void test_timeout_cb(void* arg) { static struct ETCP_LINK* first_initialized_link(struct UTUN_INSTANCE* inst) { if (!inst || !inst->connections) return NULL; - struct ETCP_LINK* link = inst->connections->links; - while (link) { - if (link->initialized) return link; - link = link->next; - } + struct ll_entry* entry = inst->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; + struct ETCP_LINK* link = conn->links; + while (link) { if (link->initialized) return link; link = link->next; } + entry = entry->next; } return NULL; } @@ -278,7 +279,7 @@ static int test_provider_egress(void) { DEBUG_INFO(DEBUG_CATEGORY_NAT, "=== test_provider_egress ==="); // Get the connection from client to provider (for debug info only) - struct ETCP_CONN* client_conn = inst_client->connections; + struct ETCP_CONN* client_conn = (inst_client->connections && inst_client->connections->head) ? ((struct conn_queue_entry*)inst_client->connections->head->data)->conn : NULL; if (!client_conn) { DEBUG_ERROR(DEBUG_CATEGORY_NAT, "No client connections"); return 0; diff --git a/tests/test_pkt_normalizer_etcp.c b/tests/test_pkt_normalizer_etcp.c index e473d6fd..5312ef9d 100644 --- a/tests/test_pkt_normalizer_etcp.c +++ b/tests/test_pkt_normalizer_etcp.c @@ -127,9 +127,9 @@ static int verify_packet_data(uint8_t* buffer, int size, int expected_seq) { } static int is_connection_established(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } conn = conn->next; } + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } entry = entry->next; } return 0; } @@ -189,16 +189,18 @@ static void monitor_and_send(void* arg) { if (!connection_checked) { if (is_connection_established(client_instance)) { connection_checked = 1; - if (!client_pn && client_instance->connections) { - client_pn = pn_init(client_instance->connections); + { struct ETCP_CONN* c = (client_instance->connections && client_instance->connections->head) ? ((struct conn_queue_entry*)client_instance->connections->head->data)->conn : NULL; + if (!client_pn && c) { + client_pn = pn_init(c); if (!client_pn) { test_completed = 2; return; } queue_set_callback(client_pn->output, NULL, NULL); - } - if (!server_pn && server_instance->connections) { - server_pn = pn_init(server_instance->connections); + } } + { struct ETCP_CONN* c = (server_instance->connections && server_instance->connections->head) ? ((struct conn_queue_entry*)server_instance->connections->head->data)->conn : NULL; + if (!server_pn && c) { + server_pn = pn_init(c); if (!server_pn) { test_completed = 2; return; } queue_set_callback(server_pn->output, NULL, NULL); - } + } } } } diff --git a/tests/test_route_ping.c b/tests/test_route_ping.c index 9a3d30ef..c3ce344c 100644 --- a/tests/test_route_ping.c +++ b/tests/test_route_ping.c @@ -157,27 +157,25 @@ static void cleanup_temp_configs(void) { } static int has_initialized_link(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - struct ETCP_CONN* conn = inst->connections; - while (conn) { + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* conn = ce->conn; struct ETCP_LINK* link = conn->links; while (link) { if (link->initialized) return 1; link = link->next; } - conn = conn->next; + entry = entry->next; } return 0; } static struct ETCP_CONN* find_conn_to_peer(struct UTUN_INSTANCE* inst, uint64_t peer_node_id) { - if (!inst) return NULL; - struct ETCP_CONN* conn = inst->connections; - while (conn) { - if (conn->peer_node_id == peer_node_id) return conn; - conn = conn->next; - } - return NULL; + if (!inst || !inst->connections) return NULL; + struct ll_entry* entry = queue_find_data_by_index(inst->connections, (const uint8_t*)&peer_node_id); + return entry ? ((struct conn_queue_entry*)entry->data)->conn : NULL; } static void test_timeout_cb(void* arg) { diff --git a/tests/test_socks_http_proxy.c b/tests/test_socks_http_proxy.c index ba468c0d..eb89db00 100644 --- a/tests/test_socks_http_proxy.c +++ b/tests/test_socks_http_proxy.c @@ -215,10 +215,15 @@ static char* make_cfg_client(void) { } static int conn_ready(struct UTUN_INSTANCE* inst) { - if (!inst) return 0; - for (struct ETCP_CONN* c = inst->connections; c; c = c->next) + if (!inst || !inst->connections) return 0; + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + struct ETCP_CONN* c = ce->conn; for (struct ETCP_LINK* l = c->links; l; l = l->next) if (l->initialized && c->crypto_ctx.initialized) return 1; + entry = entry->next; + } return 0; } diff --git a/tests/test_stcp_traffic.c b/tests/test_stcp_traffic.c index a1a72669..d015f0fe 100644 --- a/tests/test_stcp_traffic.c +++ b/tests/test_stcp_traffic.c @@ -178,8 +178,11 @@ int main(void) { TASSERT(cli_send_ready); TASSERT(cli_link); - { struct ETCP_CONN *c = srv_inst->connections; - while (c) { if (c->transport_link) { srv_link = (struct stcp_link *)c->transport_link; break; } c = c->next; } + { struct ll_entry* e = srv_inst->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN *c = ce->conn; + if (c->transport_link) { srv_link = (struct stcp_link *)c->transport_link; break; } + e = e->next; } } TASSERT(srv_link); diff --git a/tests/test_tcp_proxy_remote.c b/tests/test_tcp_proxy_remote.c index 915af76b..a1536abc 100644 --- a/tests/test_tcp_proxy_remote.c +++ b/tests/test_tcp_proxy_remote.c @@ -118,10 +118,16 @@ static void monitor(void* arg) { if (g_done) return; if (g_phase == 0 && !g_conn_up) { int bup=0, eup=0; - if (g_b) for (struct ETCP_CONN*c=g_b->connections;c;c=c->next) + if (g_b) { struct ll_entry* e=g_b->connections->head; + while (e) { struct conn_queue_entry* ce=(struct conn_queue_entry*)e->data; + struct ETCP_CONN*c=ce->conn; for (struct ETCP_LINK*l=c->links;l;l=l->next) if (l->initialized&&c->crypto_ctx.initialized) bup=1; - if (g_exit) for (struct ETCP_CONN*c=g_exit->connections;c;c=c->next) + e=e->next; } } + if (g_exit) { struct ll_entry* e=g_exit->connections->head; + while (e) { struct conn_queue_entry* ce=(struct conn_queue_entry*)e->data; + struct ETCP_CONN*c=ce->conn; for (struct ETCP_LINK*l=c->links;l;l=l->next) if (l->initialized&&c->crypto_ctx.initialized) eup=1; + e=e->next; } } if (bup && eup) { g_conn_up=1; start_test(); } } if (g_conn_up && !g_done) poll_test(); diff --git a/tests/test_udp_proxy.c b/tests/test_udp_proxy.c index 15e19a0c..3dc96613 100644 --- a/tests/test_udp_proxy.c +++ b/tests/test_udp_proxy.c @@ -118,12 +118,18 @@ static void monitor(void* arg) { if (g_test_phase == 0) { // Wait for ETCP connection int cli_ok = 0, exit_ok = 0; - for (struct ETCP_CONN* c = cli->connections; c; c = c->next) + { struct ll_entry* e = cli->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* c = ce->conn; for (struct ETCP_LINK* l = c->links; l; l = l->next) if (l->initialized && c->crypto_ctx.initialized) cli_ok = 1; - for (struct ETCP_CONN* c = exit_node->connections; c; c = c->next) + e = e->next; } } + { struct ll_entry* e = exit_node->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_CONN* c = ce->conn; for (struct ETCP_LINK* l = c->links; l; l = l->next) if (l->initialized && c->crypto_ctx.initialized) exit_ok = 1; + e = e->next; } } if (cli_ok && exit_ok) { g_test_phase = 1; // Bind client handler for UDP replies (overrides what tcp_proxy_client_create set, for test verification) diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 0016561a..362fa204 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -564,10 +564,6 @@ 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); @@ -764,8 +760,7 @@ 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); } } + etcp_conn_set_peer_node_id(conn, inv->node_id); 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", @@ -810,10 +805,6 @@ 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); @@ -962,10 +953,7 @@ void chat_core_connect_auto(uint64_t node_id, } sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys); sc_set_peer_public_key(&conn->crypto_ctx, pubkey, 0); - conn->peer_node_id = node_id; - conn->next = g_cc.inst->connections; - g_cc.inst->connections = conn; - g_cc.inst->connections_count++; + etcp_conn_set_peer_node_id(conn, node_id); struct ca_ctx* pctx = u_calloc(1, sizeof(struct ca_ctx)); if (!pctx) { etcp_connection_close(conn); pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; } @@ -980,7 +968,6 @@ 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 a6a679d8..d91f6124 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -90,15 +90,18 @@ static void ac_gc(struct auto_connect* ac) { if (now - ac->flights[i].created_tb < deadline) continue; uint64_t nid = ac->flights[i].node_id; /* check if link already UP — if so, just free slot (already connected) */ - struct ETCP_CONN* c = ac->inst->connections; int found_up = 0; - while (c) { - if (c->peer_node_id == nid) { - struct ETCP_LINK* l = c->links; - while (l) { if (l->initialized && l->link_status) { found_up = 1; break; } l = l->next; } - if (found_up) break; + { + struct ll_entry* entry = ac->inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + if (ce->conn->peer_node_id == nid) { + struct ETCP_LINK* l = ce->conn->links; + while (l) { if (l->initialized && l->link_status) { found_up = 1; break; } l = l->next; } + if (found_up) break; + } + entry = entry->next; } - c = c->next; } if (found_up) { chat_core_connect_auto_cancel(ac->flights[i].ca_state); @@ -130,12 +133,9 @@ static int ac_find_free_slot(struct auto_connect* ac) { /* ── check if node already has active ETCP connection ── */ static int ac_node_has_conn(struct UTUN_INSTANCE* inst, uint64_t nid) { - struct ETCP_CONN* c = inst->connections; - while (c) { - if (c->peer_node_id == nid) return 1; - c = c->next; - } - return 0; + if (!inst->connections) return 0; + struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&nid); + return e != NULL; } /* ── fill up to AC_MAX_FLIGHTS by advancing cursor ── */ @@ -433,24 +433,22 @@ 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 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; } + if (!inst->connections) return NULL; + struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id); + if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; if (ce->conn->initialized && ce->conn->links_up) return ce->conn; } return NULL; } /* ── helper: check if peer has active ETCP link ── */ static int cs_is_peer_online(struct UTUN_INSTANCE* inst, uint64_t peer_id) { - struct ETCP_CONN* c = inst->connections; - while (c) { - if (c->peer_node_id == peer_id) { - struct ETCP_LINK* l = c->links; - while (l) { if (l->initialized && l->link_status) return 1; l = l->next; } - } - c = c->next; + if (!inst->connections) return 0; + struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&peer_id); + if (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_LINK* l = ce->conn->links; + while (l) { if (l->initialized && l->link_status) return 1; l = l->next; } } return 0; } @@ -509,10 +507,6 @@ 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; @@ -617,13 +611,14 @@ 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); - etcp_conn_add_down_cbk(c, cs_on_conn_down, NULL); - c = c->next; + { + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + etcp_conn_add_up_cbk(ce->conn, cs_on_conn_up, NULL); + etcp_conn_add_down_cbk(ce->conn, cs_on_conn_down, NULL); + entry = entry->next; + } } cs_refresh_channels(cs); @@ -646,17 +641,18 @@ 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); - struct ETCP_CONN* c = inst->connections; - while (c) { - etcp_conn_remove_up_cbk(c, cs_on_conn_up, NULL); - etcp_conn_remove_down_cbk(c, cs_on_conn_down, NULL); - c = c->next; + { + struct ll_entry* entry = inst->connections->head; + while (entry) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; + etcp_conn_remove_up_cbk(ce->conn, cs_on_conn_up, NULL); + etcp_conn_remove_down_cbk(ce->conn, cs_on_conn_down, NULL); + entry = entry->next; + } } for (int i = 0; i < cs->channel_count; i++)