From e55c4ec7cef25eec09bd8cad39011d027c43d2be Mon Sep 17 00:00:00 2001 From: Evgeny Date: Tue, 9 Jun 2026 01:54:06 +0300 Subject: [PATCH 1/3] bugfixes --- src/lwip_tcp/lwip_tcp.c | 17 +++++++---------- src/lwip_tcp/lwip_tcp.h | 13 ------------- src/lwip_tcp/lwip_tcp_in.c | 8 ++------ src/lwip_tcp/lwip_tcp_opts.h | 7 ++++--- src/lwip_tcp/lwip_tcp_priv.h | 4 +--- 5 files changed, 14 insertions(+), 35 deletions(-) diff --git a/src/lwip_tcp/lwip_tcp.c b/src/lwip_tcp/lwip_tcp.c index ef6f9515..54e84bce 100644 --- a/src/lwip_tcp/lwip_tcp.c +++ b/src/lwip_tcp/lwip_tcp.c @@ -103,7 +103,7 @@ struct lwip_tcp_ctx *lwip_tcp_init(struct UASYNC *ua, tcp_output_fn output, void return NULL; } ctx->iss_seed = (uint16_t)(get_time_tb() & 0xFFFF); - ctx->tmr_interval_ms = TCP_TMR_INTERVAL; + ctx->tmr_interval_ms = TCP_TMR_INTERVAL / 4; ctx->rto_min_ms = 3000; ctx->rto_max_ms = 0; ctx->timer = uasync_set_timeout(ua, ctx->tmr_interval_ms * 10, ctx, tcp_tmr_cb, "lwip_tcp_tmr"); @@ -749,7 +749,6 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx) tcp_pcb_purge(pcb); if (prev != NULL) { prev->next = pcb->next; - prev->next_owner = PCB_NEXT_SLOWTMR; } else { ctx->active_pcbs = pcb->next; } @@ -808,18 +807,16 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx) tcp_pcb_purge(pcb); if (prev != NULL) { prev->next = pcb->next; - prev->next_owner = PCB_NEXT_SLOWTMR; } else { ctx->tw_pcbs = pcb->next; } pcb2 = pcb; pcb = pcb->next; int self_loop = (pcb == pcb2); - uint8_t sl_owner = pcb2->next_owner; tcp_free(pcb2); if (self_loop) { - DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_slowtmr: freed pcb=%p state=TIME_WAIT port=%u next_owner=%d, clearing tw_pcbs", - (void*)pcb2, pcb2->local_port, sl_owner); + DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_slowtmr: freed pcb=%p state=TIME_WAIT port=%u, clearing tw_pcbs", + (void*)pcb2, pcb2->local_port); ctx->tw_pcbs = NULL; break; } @@ -1050,8 +1047,8 @@ static void tcp_kill_timewait(struct lwip_tcp_ctx *ctx) } if (inactive != NULL) { if (inactive->next == inactive) { - DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_kill_timewait: aborting pcb=%p next_owner=%d, clearing tw_pcbs", - (void*)inactive, inactive->next_owner); + DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_kill_timewait: aborting pcb=%p, clearing tw_pcbs", + (void*)inactive); ctx->tw_pcbs = NULL; tcp_free(inactive); } else { @@ -1115,8 +1112,8 @@ struct tcp_pcb *tcp_alloc(struct lwip_tcp_ctx *ctx, uint8_t prio) pcb->rcv_wnd = pcb->rcv_ann_wnd = TCPWND16(TCP_WND); pcb->ttl = 64; pcb->mss = INITIAL_MSS; - pcb->rto = (int16_t)(ctx->rto_min_ms / TCP_SLOW_INTERVAL); - pcb->sv = (int16_t)(ctx->rto_min_ms / TCP_SLOW_INTERVAL); + pcb->rto = (int16_t)(TCP_RTO_MIN_MS / TCP_SLOW_INTERVAL); + pcb->sv = (int16_t)(TCP_RTO_MIN_MS / TCP_SLOW_INTERVAL); pcb->rtime = -1; pcb->cwnd = 1; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "tcp_alloc: cwnd=%u snd_buf=%u mss=%u snd_wnd=%u", pcb->cwnd, pcb->snd_buf, pcb->mss, pcb->snd_wnd); diff --git a/src/lwip_tcp/lwip_tcp.h b/src/lwip_tcp/lwip_tcp.h index 0fa4ccae..d3058d77 100644 --- a/src/lwip_tcp/lwip_tcp.h +++ b/src/lwip_tcp/lwip_tcp.h @@ -67,15 +67,6 @@ enum tcp_err_enum { LERR_CLSD = -15, LERR_ARG = -16 }; - -// кто последним записал pcb->next (диагностика зацикливания tw_pcbs) -enum pcb_next_owner { - PCB_NEXT_NONE = 0, - PCB_NEXT_REG = 1, - PCB_NEXT_RMV = 2, - PCB_NEXT_SLOWTMR = 3, - PCB_NEXT_INPUT = 4, -}; typedef int err_t; // Forward declaration for callbacks @@ -175,9 +166,6 @@ struct tcp_pcb { tcp_connected_fn connected; tcp_poll_fn poll; tcp_err_fn errf; - - uint8_t next_owner; // who last wrote pcb->next (enum pcb_next_owner) - uint32_t keep_idle; uint8_t persist_cnt; uint8_t persist_backoff; @@ -193,7 +181,6 @@ struct tcp_pcb_listen { uint8_t prio; uint16_t local_port; uint32_t local_ip; - uint8_t next_owner; tcp_accept_fn accept; }; diff --git a/src/lwip_tcp/lwip_tcp_in.c b/src/lwip_tcp/lwip_tcp_in.c index 9b1cbdf9..d7d214f5 100644 --- a/src/lwip_tcp/lwip_tcp_in.c +++ b/src/lwip_tcp/lwip_tcp_in.c @@ -198,9 +198,7 @@ void lwip_tcp_input(struct lwip_tcp_ctx *ctx, struct pbuf *p, pcb->local_ip == dst_ip) { if (prev != NULL) { prev->next = pcb->next; - prev->next_owner = PCB_NEXT_INPUT; pcb->next = ctx->active_pcbs; - pcb->next_owner = PCB_NEXT_INPUT; ctx->active_pcbs = pcb; } break; @@ -225,8 +223,8 @@ void lwip_tcp_input(struct lwip_tcp_ctx *ctx, struct pbuf *p, pcb->remote_ip == src_ip && pcb->local_ip == dst_ip) { if (pcb->next == pcb) { - DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in lwip_tcp_input: pcb=%p sport=%u dport=%u next_owner=%d, aborting", - (void*)pcb, sport, dport, pcb->next_owner); + DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in lwip_tcp_input: pcb=%p sport=%u dport=%u, aborting", + (void*)pcb, sport, dport); tcp_abort(pcb); } else { tcp_timewait_input(pcb); @@ -249,9 +247,7 @@ void lwip_tcp_input(struct lwip_tcp_ctx *ctx, struct pbuf *p, if (lpcb != NULL) { if (prev != NULL) { ((struct tcp_pcb_listen *)prev)->next = lpcb->next; - ((struct tcp_pcb_listen *)prev)->next_owner = PCB_NEXT_INPUT; lpcb->next = ctx->listen_pcbs; - lpcb->next_owner = PCB_NEXT_INPUT; ctx->listen_pcbs = (struct tcp_pcb *)lpcb; } tcp_listen_input(lpcb); diff --git a/src/lwip_tcp/lwip_tcp_opts.h b/src/lwip_tcp/lwip_tcp_opts.h index 28cfbf11..4c5f5b59 100644 --- a/src/lwip_tcp/lwip_tcp_opts.h +++ b/src/lwip_tcp/lwip_tcp_opts.h @@ -14,9 +14,10 @@ #define TCP_FAST_INTERVAL TCP_TMR_INTERVAL #define TCP_SLOW_INTERVAL (2 * TCP_TMR_INTERVAL) -#define TCP_FIN_WAIT_TIMEOUT 6000 // ms -#define TCP_SYN_RCVD_TIMEOUT 6000 // ms -#define TCP_MSL 25000 // ms (2*MSL = 50s) +#define TCP_RTO_MIN_MS 3000 // ms, initial RTO +#define TCP_FIN_WAIT_TIMEOUT 20000 // ms +#define TCP_SYN_RCVD_TIMEOUT 20000 // ms +#define TCP_MSL 60000 // ms (2*MSL = 120s) #define TCP_OOSEQ_TIMEOUT 6 // x RTO #define TCP_KEEPIDLE_DEFAULT 7200000 // ms (unused, no keepalive) diff --git a/src/lwip_tcp/lwip_tcp_priv.h b/src/lwip_tcp/lwip_tcp_priv.h index 882a030f..6a12f045 100644 --- a/src/lwip_tcp/lwip_tcp_priv.h +++ b/src/lwip_tcp/lwip_tcp_priv.h @@ -236,15 +236,13 @@ void lwip_tcp_stats_clear(struct lwip_tcp_ctx *ctx); // PCB list management #define TCP_REG(pcbs, npcb) do { \ - (npcb)->next_owner = PCB_NEXT_REG; \ (npcb)->next = *(pcbs); *(pcbs) = (npcb); \ } while(0) #define TCP_RMV(pcbs, npcb) do { \ if(*(pcbs) == (npcb)) { *(pcbs) = (*pcbs)->next; } \ - else { struct tcp_pcb *_tmp; for(_tmp = *(pcbs); _tmp != NULL; _tmp = _tmp->next) { if(_tmp->next == (npcb)) { _tmp->next_owner = PCB_NEXT_SLOWTMR; _tmp->next = (npcb)->next; break; } } } \ + else { struct tcp_pcb *_tmp; for(_tmp = *(pcbs); _tmp != NULL; _tmp = _tmp->next) { if(_tmp->next == (npcb)) { _tmp->next = (npcb)->next; break; } } } \ (npcb)->next = NULL; \ - (npcb)->next_owner = PCB_NEXT_RMV; \ } while(0) #define TCP_REG_ACTIVE(ctx, npcb) TCP_REG(&(ctx)->active_pcbs, npcb) From 5982b011b458978e4e786346a83b3df46ba4b600 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Tue, 9 Jun 2026 16:39:41 +0300 Subject: [PATCH 2/3] etcp: restore etcp_link_update_inflight_lim + state machine fixes + diagnostics - etcp_link_update_inflight_lim: self-contained method with clamping, send_blocked_inflight check, loadbalancer_link_ready on cwnd increase - etcp_conn_on_inflight_lim_changed: recalc optimal_inflight + resume input - etcp_ack_recv: use etcp_link_update_inflight_lim instead of manual set - etcp_request_pkt: fix line 928 check old link (inf_pkt->last_link) not new - etcp_link_close: recalc optimal_inflight on link removal (both branches) - etcp_socket_remove: fix dangling pointer loop (nullify after close) - uasync_print_resources: show timer names, remain time, deleted entries - diagnostics in utun_instance_destroy and test_etcp_reconnect --- .gitignore | 1 + lib/u_async.c | 21 ++++++++++++++++----- src/etcp.c | 20 +++++++++++++------- src/etcp.h | 3 +++ src/etcp_connections.c | 28 ++++++++++++++++++++++------ src/etcp_connections.h | 4 ++++ src/utun_instance.c | 5 ++++- tests/test_etcp_reconnect.c | 21 +++++++++++++++++++-- 8 files changed, 82 insertions(+), 21 deletions(-) diff --git a/.gitignore b/.gitignore index 17e9e7c6..0595b084 100644 --- a/.gitignore +++ b/.gitignore @@ -108,3 +108,4 @@ cl build_win.log /wintun.dll src/lwip_orig/ +lib/timeout/ diff --git a/lib/u_async.c b/lib/u_async.c index 67448299..53316c53 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -1480,17 +1480,28 @@ void uasync_print_resources(struct UASYNC* ua, const char* prefix) { // Показать активные таймеры if (ua->timeout_heap) { + struct timeval now_tv; + get_current_time(&now_tv); + uint64_t now_ms = timeval_to_ms(&now_tv); size_t active_timers = 0; - // Безопасное чтение без извлечения - просто итерируем по массиву + size_t deleted_timers = 0; for (size_t i = 0; i < ua->timeout_heap->size; i++) { - if (!ua->timeout_heap->heap[i].deleted) { + uint64_t exp = ua->timeout_heap->heap[i].expiration; + int64_t remain = (int64_t)(exp - now_ms); + if (ua->timeout_heap->heap[i].deleted) { + deleted_timers++; + struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer DELETED: node=%p name='%s' exp=%llu remain=%lldms", + node, node ? node->name : "null", (unsigned long long)exp, (long long)remain); + } else { active_timers++; struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer: node=%p, expires=%llu ms", - node, (unsigned long long)ua->timeout_heap->heap[i].expiration); + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer ACTIVE: node=%p name='%s' exp=%llu remain=%lldms", + node, node ? node->name : "null", (unsigned long long)exp, (long long)remain); } } - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Active timers in heap: %zu", active_timers); + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Active timers in heap: %zu, deleted: %zu, total: %zu", + active_timers, deleted_timers, ua->timeout_heap->size); } // Показать активные сокеты diff --git a/src/etcp.c b/src/etcp.c index 210a56a6..679844de 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -612,6 +612,16 @@ static void input_queue_try_resume(struct ETCP_CONN* etcp) {// при ACK } } +// Called from link-level etcp_link_update_inflight_lim() after inflight_lim_bytes changes. +// Recalculates connection-level optimal_inflight and resumes input_queue if room opened up. +void etcp_conn_on_inflight_lim_changed(struct ETCP_CONN* etcp) { + if (!etcp) return; + uint32_t sum = 0; + for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes; + etcp->optimal_inflight = sum; + input_queue_try_resume(etcp); +} + void etcp_stats(struct ETCP_CONN* etcp) { if (!etcp) return; @@ -924,7 +934,8 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { bbr_note_loss(inf_pkt->last_link->bbr); inf_pkt->last_link->inflight_bytes -= inf_pkt->ll.len; inf_pkt->last_link->inflight_packets--; - if (link->send_blocked_inflight && link->inflight_bytes < link->inflight_lim_bytes) loadbalancer_link_ready(link); + if (inf_pkt->last_link->send_blocked_inflight && inf_pkt->last_link->inflight_bytes < inf_pkt->last_link->inflight_lim_bytes) + loadbalancer_link_ready(inf_pkt->last_link); } // Always add to the CURRENT link (first send or retransmission) @@ -1239,14 +1250,9 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d uint32_t new_pacing = link->bbr_pacing_rate; bbr_main(link->bbr, &rs, &new_cwnd, &new_pacing, link->mtu, link->inflight_bytes, (link->inflight_bytes >= link->inflight_lim_bytes)); - link->inflight_lim_bytes = new_cwnd; + etcp_link_update_inflight_lim(link, new_cwnd); link->bbr_pacing_rate = new_pacing; link->bandwidth = (uint32_t)((uint64_t)new_pacing * 8 / 1000); - - if (link->send_blocked_inflight && link->inflight_bytes < new_cwnd) { - link->send_blocked_inflight = 0; - loadbalancer_link_ready(link); - } link->bbr_loss_since_ack = 0; } diff --git a/src/etcp.h b/src/etcp.h index 8bbc9af2..e8d2945f 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -230,6 +230,9 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp); // Process ACK receipt - remove acknowledged packet from inflight queues void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t dts); +// Recalculate optimal_inflight and resume input queue after link inflight_lim change +void etcp_conn_on_inflight_lim_changed(struct ETCP_CONN* etcp); + // Update log_name when peer_node_id becomes known void etcp_update_log_name(struct ETCP_CONN* etcp); diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 0b8be32a..1b6a437d 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -736,9 +736,9 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn) { size_t i = 0; while (i < conn->num_channels) { struct ETCP_LINK* l = conn->links[i]; - int had_conn = (l->conn != NULL); etcp_link_close(l); - if (!had_conn) i++; + conn->links[i] = NULL; // обнулить после — etcp_link_close обращается к массиву через remove_link + i++; } u_free(conn->links); @@ -832,10 +832,7 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn while (l && l->next) l=l->next; if (l) l->next = link; else etcp->links = link; -// пересчитать connection-level optimal_inflight - { uint32_t sum = 0; - for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes; - etcp->optimal_inflight = sum; } + etcp_link_update_inflight_lim(link, link->inflight_lim_bytes); DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "NEW link initialized on etcp=[%s] link=%p socket=%s id=%d is_server=%d mtu=%d", etcp->log_name, link, conn->name, link->local_link_id, link->is_server, link->mtu); @@ -847,6 +844,23 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn return link; } +void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) { + if (!link || !link->etcp) return; + if (new_lim < INFLIGHT_LIM_MIN) new_lim = INFLIGHT_LIM_MIN; + if (new_lim > INFLIGHT_LIM_MAX) new_lim = INFLIGHT_LIM_MAX; + + uint32_t old = link->inflight_lim_bytes; + link->inflight_lim_bytes = new_lim; + + if (old != new_lim && link->inflight_bytes < new_lim && link->send_blocked_inflight) { + link->send_blocked_inflight = 0; + loadbalancer_link_ready(link); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] unblocked link %p (inflight_lim %u -> %u)", link->etcp->log_name, link, old, new_lim); + } + + etcp_conn_on_inflight_lim_changed(link->etcp); +} + void etcp_link_close(struct ETCP_LINK* link) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); if (!link) return; @@ -863,6 +877,7 @@ void etcp_link_close(struct ETCP_LINK* link) { if (link->init_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); if (link->shaper_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->shaper_timer); if (link->keepalive_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer); + etcp_conn_on_inflight_lim_changed(link->etcp); u_free(link); return; } @@ -898,6 +913,7 @@ void etcp_link_close(struct ETCP_LINK* link) { remove_link(link->conn, link->ip_port_hash); + etcp_conn_on_inflight_lim_changed(link->etcp); u_free(link->bbr); u_free(link); } diff --git a/src/etcp_connections.h b/src/etcp_connections.h index 4bb68cf6..5fb2c0e4 100644 --- a/src/etcp_connections.h +++ b/src/etcp_connections.h @@ -14,6 +14,9 @@ #define UDP_SC_HDR_SIZE (13+8+4 + 5)// 13+8+4 - sc_nonce+tag size+crc, 5 - payload hdr #define ACK_REZERV 100// сколько байт резервировать под ack и прочие заголовки +#define INFLIGHT_LIM_MIN 8192 // 8K +#define INFLIGHT_LIM_MAX 1048576 // 1M + #define PACKET_DATA_SIZE 1600//1536 #define PACKET_DATA_MAX_MTU 1600 #define ETCP_MAX_PAYLOAD_SIZE (PACKET_DATA_SIZE - ETCP_ACK_BASE_SIZE - 5) /* 1587 */ @@ -262,6 +265,7 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn); // connection functions // создает новый канал связи для etcp подключения (ETCP_CONN) struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server); +void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim); void etcp_link_close(struct ETCP_LINK* link); //int etcp_input_cbk(struct packet_buffer* pkt, struct ETCP_SOCKET* conn);// получает расшифрованный пакет int etcp_encrypt_send(struct ETCP_DGRAM* dgram);// зашифровывает и отправляет пакет diff --git a/src/utun_instance.c b/src/utun_instance.c index c16bb59e..74030a30 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -279,6 +279,7 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { // Диагностика ресурсов ДО cleanup utun_instance_diagnose_leaks(instance, "BEFORE_CLEANUP"); + if (instance->ua) uasync_print_resources(instance->ua, "INSTANCE_DESTROY_BEFORE"); // Stop running if not already instance->running = 0; @@ -397,7 +398,9 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { memory_pool_destroy(instance->data_pool); instance->data_pool = NULL; } - + + if (instance->ua) uasync_print_resources(instance->ua, "INSTANCE_DESTROY_AFTER"); + // Note: uasync is NOT destroyed here - caller must destroy it separately // This allows sharing uasync between multiple instances instance->ua = NULL; diff --git a/tests/test_etcp_reconnect.c b/tests/test_etcp_reconnect.c index 5baab2fe..6c56f7ae 100644 --- a/tests/test_etcp_reconnect.c +++ b/tests/test_etcp_reconnect.c @@ -203,6 +203,9 @@ static void monitor(void* arg) { server_instance->running = 0; utun_instance_destroy(server_instance); server_instance = NULL; + printf("--- uasync after server destroy ---\n"); + debug_set_level(DEBUG_LEVEL_INFO); + uasync_print_resources(ua, "AFTER_SERVER_DESTROY"); printf("Recreating server instance...\n"); server_instance = utun_instance_create(ua, server_config_path); if (!server_instance || utun_instance_init(server_instance) < 0) { @@ -265,6 +268,9 @@ static void monitor(void* arg) { client_instance->running = 0; utun_instance_destroy(client_instance); client_instance = NULL; + printf("--- uasync after client destroy ---\n"); + debug_set_level(DEBUG_LEVEL_INFO); + uasync_print_resources(ua, "AFTER_CLIENT_DESTROY"); printf("Recreating client instance...\n"); client_instance = utun_instance_create(ua, client_config_path); if (!client_instance || utun_instance_init(client_instance) < 0) { @@ -320,8 +326,6 @@ int main(void) { if (create_temp_configs() != 0) return 1; debug_config_init(); - debug_set_level(DEBUG_LEVEL_ERROR); - debug_set_categories(DEBUG_CATEGORY_NONE); printf("=== ETCP Reconnect Test ===\n"); printf("Server port: %d, Client port: %d\n", server_port, client_port); @@ -344,6 +348,11 @@ int main(void) { } printf("Client ready (node_id=%llx)\n", (unsigned long long)client_instance->node_id); + // Enable diagnostic output for timer investigation + debug_set_level(DEBUG_LEVEL_INFO); + debug_set_category_level(DEBUG_CATEGORY_UASYNC, DEBUG_LEVEL_INFO); + debug_set_category_level(DEBUG_CATEGORY_TIMERS, DEBUG_LEVEL_INFO); + packet_timeout_id = uasync_set_timeout(ua, 300, NULL, monitor, "test_monitor"); global_timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS * 10, NULL, test_timeout, "test_timeout"); printf("Starting main loop\n"); @@ -363,8 +372,16 @@ int main(void) { if (packet_timeout_id) uasync_cancel_timeout(ua, packet_timeout_id); if (global_timeout_id) uasync_cancel_timeout(ua, global_timeout_id); + printf("--- uasync before final cleanup ---\n"); + debug_set_level(DEBUG_LEVEL_INFO); + uasync_print_resources(ua, "BEFORE_FINAL_CLEANUP"); + if (server_instance) { server_instance->running = 0; utun_instance_destroy(server_instance); server_instance = NULL; } if (client_instance) { client_instance->running = 0; utun_instance_destroy(client_instance); client_instance = NULL; } + + printf("--- uasync before uasync_destroy ---\n"); + uasync_print_resources(ua, "BEFORE_UASYNC_DESTROY"); + if (ua) { uasync_destroy(ua, 0); ua = NULL; } cleanup_temp_configs(); From 9a3eb1ce8023bde621ae9a83fbbe58c4f6545ff0 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Tue, 9 Jun 2026 19:35:06 +0300 Subject: [PATCH 3/3] =?UTF-8?q?fix:=20reconnect=20test=20=E2=80=94=20clear?= =?UTF-8?q?=5Fqueue=20callback=20deadlock=20+=20conn=5Fok=20timing?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - etcp_conn_reset: queue_resume_callback after clear_queue (input_queue, input_send_q, input_wait_ack) to prevent suspended callback deadlock - is_connection_established: check conn->initialized (reliable across reinit) instead of link->initialized (never cleared on reinit) - test: force etcp_conn_reinit after server recreation so conn_ok waits for fresh INIT handshake - test: guard send_packets() with is_connection_established in reconnect phases to avoid sending data while conn is down --- src/etcp.c | 5 +++++ tests/test_etcp_reconnect.c | 34 ++++++++-------------------------- 2 files changed, 13 insertions(+), 26 deletions(-) diff --git a/src/etcp.c b/src/etcp.c index 679844de..1803a0b7 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -391,6 +391,11 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) { clear_queue(etcp->recv_q); clear_queue(etcp->ack_q); + // clear_queue leaves callback_suspended=1, resume to prevent deadlock + queue_resume_callback(etcp->input_queue); + queue_resume_callback(etcp->input_send_q); + queue_resume_callback(etcp->input_wait_ack); + // В etcp_conn_reset(), после очистки очередей добавьте: struct ETCP_LINK* l = etcp->links; while (l) { diff --git a/tests/test_etcp_reconnect.c b/tests/test_etcp_reconnect.c index 6c56f7ae..2f9b8122 100644 --- a/tests/test_etcp_reconnect.c +++ b/tests/test_etcp_reconnect.c @@ -117,11 +117,7 @@ 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; - } + if (conn->initialized) return 1; conn = conn->next; } return 0; @@ -203,9 +199,6 @@ static void monitor(void* arg) { server_instance->running = 0; utun_instance_destroy(server_instance); server_instance = NULL; - printf("--- uasync after server destroy ---\n"); - debug_set_level(DEBUG_LEVEL_INFO); - uasync_print_resources(ua, "AFTER_SERVER_DESTROY"); printf("Recreating server instance...\n"); server_instance = utun_instance_create(ua, server_config_path); if (!server_instance || utun_instance_init(server_instance) < 0) { @@ -214,6 +207,9 @@ static void monitor(void* arg) { return; } 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; } } restart_action_done = 1; } drain_received(0); @@ -234,7 +230,8 @@ static void monitor(void* arg) { packets_received = 0; drain_received(0); } - send_packets(); + if (is_connection_established(client_instance)) + send_packets(); drain_received(1); if (packets_received >= TOTAL_PACKETS) { printf("=== Phase 4 done: sent=%d received=%d ===\n", packets_sent, packets_received); @@ -268,9 +265,6 @@ static void monitor(void* arg) { client_instance->running = 0; utun_instance_destroy(client_instance); client_instance = NULL; - printf("--- uasync after client destroy ---\n"); - debug_set_level(DEBUG_LEVEL_INFO); - uasync_print_resources(ua, "AFTER_CLIENT_DESTROY"); printf("Recreating client instance...\n"); client_instance = utun_instance_create(ua, client_config_path); if (!client_instance || utun_instance_init(client_instance) < 0) { @@ -299,7 +293,8 @@ static void monitor(void* arg) { packets_received = 0; drain_received(0); } - send_packets(); + if (is_connection_established(client_instance)) + send_packets(); drain_received(1); if (packets_received >= TOTAL_PACKETS) { printf("=== Phase 7 done: sent=%d received=%d ===\n", packets_sent, packets_received); @@ -348,11 +343,6 @@ int main(void) { } printf("Client ready (node_id=%llx)\n", (unsigned long long)client_instance->node_id); - // Enable diagnostic output for timer investigation - debug_set_level(DEBUG_LEVEL_INFO); - debug_set_category_level(DEBUG_CATEGORY_UASYNC, DEBUG_LEVEL_INFO); - debug_set_category_level(DEBUG_CATEGORY_TIMERS, DEBUG_LEVEL_INFO); - packet_timeout_id = uasync_set_timeout(ua, 300, NULL, monitor, "test_monitor"); global_timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS * 10, NULL, test_timeout, "test_timeout"); printf("Starting main loop\n"); @@ -372,16 +362,8 @@ int main(void) { if (packet_timeout_id) uasync_cancel_timeout(ua, packet_timeout_id); if (global_timeout_id) uasync_cancel_timeout(ua, global_timeout_id); - printf("--- uasync before final cleanup ---\n"); - debug_set_level(DEBUG_LEVEL_INFO); - uasync_print_resources(ua, "BEFORE_FINAL_CLEANUP"); - if (server_instance) { server_instance->running = 0; utun_instance_destroy(server_instance); server_instance = NULL; } if (client_instance) { client_instance->running = 0; utun_instance_destroy(client_instance); client_instance = NULL; } - - printf("--- uasync before uasync_destroy ---\n"); - uasync_print_resources(ua, "BEFORE_UASYNC_DESTROY"); - if (ua) { uasync_destroy(ua, 0); ua = NULL; } cleanup_temp_configs();