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 3911ffa2..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) { @@ -612,6 +617,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 +939,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,13 +1255,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)); + 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; } @@ -1277,6 +1289,10 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d void etcp_conn_input(struct ETCP_DGRAM* pkt) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); if (!pkt) return; + if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt); + volatile int _halt = 1; while (_halt) {} + } if (!pkt->data_len) { memory_pool_free(pkt->link->etcp->instance->pkt_pool, pkt); return; @@ -1387,24 +1403,22 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { break; } } - if (queue_find_data_by_index(etcp->ack_q, &seq) == NULL) { - struct ACK_PACKET* p = (struct ACK_PACKET*)queue_entry_new_from_pool(etcp->instance->ack_pool); - if (!p) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ACK_PACKET", etcp->log_name); - len = 0; - break; - } - p->seq=seq; - p->pkt_timestamp=pkt->timestamp; - p->recv_timestamp=get_current_timestamp(); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX add to ack_q seq=%d", etcp->log_name, seq); - - queue_data_put_with_index(etcp->ack_q, (struct ll_entry*)p); - if (etcp->ack_resp_timer == NULL) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] set ack_timer for delayed ACK send", etcp->log_name); - etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb, "etcp_ack_resp"); - } - } else DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX ack dedup: seq=%d already in ack_q", etcp->log_name, seq); + struct ACK_PACKET* p = (struct ACK_PACKET*)queue_entry_new_from_pool(etcp->instance->ack_pool); + if (!p) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ACK_PACKET", etcp->log_name); + len = 0; + break; + } + p->seq=seq; + p->pkt_timestamp=pkt->timestamp; + p->recv_timestamp=get_current_timestamp(); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX add to ack_q seq=%d", etcp->log_name, seq); + + queue_data_put_with_index(etcp->ack_q, (struct ll_entry*)p); + if (etcp->ack_resp_timer == NULL) { + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] set ack_timer for delayed ACK send", etcp->log_name); + etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb, "etcp_ack_resp"); + } if (((int32_t)(etcp->last_delivered_id-seq)<0) && (queue_find_data_by_index(etcp->recv_q, &seq)==NULL)) {// проверяем есть ли пакет с этим seq uint32_t pkt_len=len-5; DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] adding packet seq=%u to recv_q (last_delivered_id=%u)", etcp->log_name, seq, etcp->last_delivered_id); @@ -1546,6 +1560,12 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { } + if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt); + volatile int _halt = 1; while (_halt) {} + } + + memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram } 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 c1d02c17..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); @@ -746,8 +746,6 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn) { } -static void bbr_cwnd_updated(void* ctx, uint32_t new_cwnd); - struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); if (!remote_addr) return NULL; @@ -789,6 +787,7 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn link->keepalive_sent_count = 0; link->keepalive_recv_count = 0; link->ka_period_ms = KA_PERIOD_MIN_MS; + link->inflight_lim_bytes = link->mtu * 4; // BBR init_cwnd (~4 packets) link->bandwidth = 10000; // начальная оценка 10 Mbps для шейпера link->burst_id = 0; link->burst_active = 0; @@ -833,9 +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; - etcp_link_update_inflight_lim(link, link->mtu * 4); - link->bbr->on_cwnd_update = bbr_cwnd_updated; - link->bbr->cwnd_update_ctx = link; + 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,16 +844,21 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn return link; } -static void bbr_cwnd_updated(void* ctx, uint32_t new_cwnd) { - etcp_link_update_inflight_lim((struct ETCP_LINK*)ctx, new_cwnd); -} - 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; - struct ETCP_CONN* etcp = link->etcp; - uint32_t sum = 0; - for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes; - etcp->optimal_inflight = sum; + + 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) { @@ -875,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; } @@ -910,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); } @@ -1743,6 +1747,10 @@ process_decrypted: // log_dump("RECV decrypted:", pkt->data, pkt->data_len, link); if (link->link_state == 3) { + if (memory_pool_is_freed(e_sock->instance->pkt_pool, pkt)) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt); + volatile int _halt = 1; while (_halt) {} + } etcp_conn_input(pkt); } else memory_pool_free(e_sock->instance->pkt_pool, pkt); return; diff --git a/src/etcp_connections.h b/src/etcp_connections.h index 98943c1e..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 */ 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) diff --git a/src/utun_instance.c b/src/utun_instance.c index 251cf054..f8588668 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -284,6 +284,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; @@ -402,7 +403,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..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; @@ -211,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); @@ -231,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); @@ -293,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); @@ -320,8 +321,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);