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