diff --git a/src/Makefile.am b/src/Makefile.am index 8e68eccc..b0af8400 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -21,6 +21,7 @@ utun_CORE_SOURCES = \ etcp_connections.c \ etcp_loadbalancer.c \ etcp_debug.c \ + etcp_dump.c \ secure_channel.c \ crc32.c \ pkt_normalizer.c \ diff --git a/src/etcp.c b/src/etcp.c index 621626fc..0973dd19 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -290,9 +290,15 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) { l = l->next; } - // Reset timers (just clear the pointers - timers will expire naturally) - etcp->retrans_timer = NULL; - etcp->ack_resp_timer = NULL; + // 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->reset_count++; diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 7c211401..11c8158d 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -1536,7 +1536,6 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { 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; - conn->session_id = session_id; etcp_update_log_name(conn); 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); @@ -1698,6 +1697,10 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { memory_pool_free(e_sock->instance->pkt_pool, pkt); link->initialized = 1; link->link_state = 3; + if (link->init_timer) { + uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); + link->init_timer = NULL; + } if (link->etcp->initialized == 0) { etcp_conn_ready(link->etcp); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Connection established: log_name=%s socket=%s link_id=%d status=UP", link->etcp->log_name, e_sock->name, link->local_link_id); @@ -1826,6 +1829,10 @@ process_decrypted: link->initialized = 1;// получен init response (client) link->link_state = 3; // connected + if (link->init_timer) { + uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); + link->init_timer = NULL; + } if (code == ETCP_INIT_RESPONSE) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit 3 %p", link->etcp); diff --git a/src/etcp_dump.c b/src/etcp_dump.c new file mode 100644 index 00000000..2611de77 --- /dev/null +++ b/src/etcp_dump.c @@ -0,0 +1,255 @@ +#include "etcp_dump.h" +#include "etcp.h" +#include "pkt_normalizer.h" +#include "../lib/ll_queue.h" +#include "../lib/memory_pool.h" +#include "../lib/socket_compat.h" +#include "../lib/swm_min.h" +#include +#include + +#define DUMP_FP stdout + +#define DUMP_LINE(fmt, ...) do { \ + fprintf(DUMP_FP, " " fmt "\n", ##__VA_ARGS__); \ +} while(0) + +#define DUMP_HDR(fmt, ...) do { \ + fprintf(DUMP_FP, fmt "\n", ##__VA_ARGS__); \ +} while(0) + +static void dump_queue(const char* name, struct ll_queue* q) { + if (!q) { DUMP_LINE("%-14s NULL", name); return; } + DUMP_LINE("%-14s %d pkts / %zu bytes", name, queue_entry_count(q), queue_total_bytes(q)); +} + +static void dump_pool(const char* name, struct memory_pool* pool) { + if (!pool) return; + size_t allocs, reuse; + memory_pool_get_stats(pool, &allocs, &reuse); + DUMP_LINE(" %-12s alloc=%zu reuse=%zu", name, allocs, reuse); +} + +static void dump_link_timer(const char* name, void* timer) { + DUMP_LINE(" %-14s %s", name, timer ? "ACTIVE" : "free"); +} + +static const char* ip_to_str_buf(uint32_t ip, char* buf, size_t sz) { + uint8_t* b = (uint8_t*)&ip; + snprintf(buf, sz, "%u.%u.%u.%u", b[0], b[1], b[2], b[3]); + return buf; +} + +static const char* sockaddr_to_str(const struct sockaddr_storage* addr, char* buf, size_t sz) { + if (!addr || addr->ss_family != AF_INET) { snprintf(buf, sz, "?"); return buf; } + const struct sockaddr_in* sin = (const struct sockaddr_in*)addr; + uint8_t* b = (uint8_t*)&sin->sin_addr.s_addr; + snprintf(buf, sz, "%u.%u.%u.%u:%u", b[0], b[1], b[2], b[3], ntohs(sin->sin_port)); + return buf; +} + +void etcp_dump_conn_state(struct ETCP_CONN* conn) { + if (!conn) { DUMP_HDR("etcp_dump_conn_state: NULL conn"); return; } + + char addr_buf[32], nat_buf[32]; + + DUMP_HDR("=== ETCP CONN STATE [%s] ===", conn->log_name); + + /* --- GENERAL --- */ + DUMP_LINE("GENERAL: %s peer=0x%llx initialized=%d links_up=%d tx_state=%d session=0x%08x mtu=%d routing_ex=%d", + conn->name ? conn->name : "", + (unsigned long long)conn->peer_node_id, conn->initialized, conn->links_up, + conn->tx_state, conn->session_id, conn->mtu, conn->routing_exchange_active); + + /* --- IDS --- */ + DUMP_LINE("IDS: next_tx=%u last_rx=%u last_del=%u rx_ack_till=%u", + conn->next_tx_id, conn->last_rx_id, conn->last_delivered_id, conn->rx_ack_till); + + /* --- RTT --- */ + DUMP_LINE("RTT: last=%u avg10=%u avg100=%u jitter=%u", + conn->rtt_last, conn->rtt_avg_10, conn->rtt_avg_100, conn->jitter); + + /* --- STATS --- */ + DUMP_LINE("STATS: bytes_sent=%u retrans=%u reinit=%u reset=%u ack_pkts=%u", + conn->bytes_sent_total, conn->retransmissions_count, + conn->reinit_count, conn->reset_count, conn->ack_packets_count); + + /* --- INFLIGHT --- */ + DUMP_LINE("INFLIGHT: unacked=%u optimal=%u", + conn->unacked_bytes, conn->optimal_inflight); + + /* --- ACK DEBUG --- */ + DUMP_LINE("ACK_DEBUG: hit_inf=%u hit_sndq=%u miss=%u link_wait=%u rx_dup=%u tx_dup=%u", + conn->cnt_ack_hit_inf, conn->cnt_ack_hit_sndq, conn->cnt_ack_miss, + conn->cnt_link_wait, conn->rx_dup_count, conn->tx_dup_count); + + /* --- DEBUG[8] --- */ + DUMP_LINE("DEBUG: [0]=%u [1]=%u [2]=%u [3]=%u [4]=%u [5]=%u [6]=%u [7]=%u", + conn->debug[0], conn->debug[1], conn->debug[2], conn->debug[3], + conn->debug[4], conn->debug[5], conn->debug[6], conn->debug[7]); + + /* --- TIMERS --- */ + DUMP_LINE("TIMERS: retrans=%s ack_resp=%s", + conn->retrans_timer ? "ACTIVE" : "free", + conn->ack_resp_timer ? "ACTIVE" : "free"); + + /* --- QUEUES --- */ + DUMP_HDR(" --- QUEUES ---"); + dump_queue("input", conn->input_queue); + dump_queue("input_send_q", conn->input_send_q); + dump_queue("input_wait_ack", conn->input_wait_ack); + dump_queue("ack_q", conn->ack_q); + dump_queue("recv_q", conn->recv_q); + dump_queue("output", conn->output_queue); + + /* --- NORMALIZER --- */ + if (conn->normalizer) { + struct PKTNORM* pn = (struct PKTNORM*)conn->normalizer; + DUMP_HDR(" --- NORMALIZER ---"); + DUMP_LINE("frag=%u data_ptr=%u/%u flush=%s pending=%s recvpart=%s", + pn->frag_size, pn->data_ptr, pn->data_size, + pn->flush_timer ? "ACTIVE" : "free", + pn->pending ? "ACTIVE" : "free", + pn->recvpart ? "ACTIVE" : "free"); + dump_queue("pn-input", pn->input); + dump_queue("pn-output", pn->output); + DUMP_LINE("alloc_err=%u logic_err=%u in=%llu/%llu out=%llu/%llu", + pn->alloc_errors, pn->logic_errors, + (unsigned long long)pn->in_total_pkts, (unsigned long long)pn->in_total_bytes, + (unsigned long long)pn->out_total_pkts, (unsigned long long)pn->out_total_bytes); + } else { + DUMP_HDR(" --- NORMALIZER --- NULL"); + } + + /* --- POOLS --- */ + DUMP_HDR(" --- POOLS ---"); + dump_pool("inflight", conn->inflight_pool); + dump_pool("io", conn->io_pool); + if (conn->instance) { + dump_pool("data", conn->instance->data_pool); + dump_pool("ack", conn->instance->ack_pool); + dump_pool("pkt", conn->instance->pkt_pool); + } + + /* --- LINKS --- */ + int link_idx = 0; + struct ETCP_LINK* link = conn->links; + while (link) { + DUMP_HDR(" --- LINK %d ---", link_idx); + + sockaddr_to_str(&link->remote_addr, addr_buf, sizeof(addr_buf)); + + /* basic link state */ + DUMP_LINE("BASIC: id=%d/%d is_server=%d link_state=%d link_status=%d initialized=%d mtu=%d", + link->local_link_id, link->remote_link_id, link->is_server, + link->link_state, link->link_status, link->initialized, link->mtu); + DUMP_LINE("ADDR: remote=%s remote_sock=%d remote_type=%d remote_only_local=%d", + addr_buf, link->remote_socket_id, link->remote_type, link->remote_only_local); + + /* keepalive */ + DUMP_LINE("KA: recv=%d remote=%d sent=%u recv=%u interval=%u timeout=%u", + link->recv_keepalive, link->remote_keepalive, + link->keepalive_sent_count, link->keepalive_recv_count, + link->keepalive_interval, link->keepalive_timeout); + + /* timers */ + DUMP_LINE("TIMERS: init=%s(%u/%u) ka=%s shaper=%s stats=%s burst=%s keepalive=%s", + link->init_timer ? "ACTIVE" : "free", + link->init_timeout, link->init_retry_count, + link->keepalive_timer ? "ACTIVE" : "free", + link->shaper_timer ? "ACTIVE" : "free", + link->stats_timer ? "ACTIVE" : "free", + link->burst_resp_timer ? "ACTIVE" : "free", + link->keepalive_sent_count > 0 ? "has_sent" : "idle"); + + /* inflight */ + DUMP_LINE("INFLIGHT: bytes=%u pkts=%u lim=%u blocked=%d phase=%d sst=%u last_win_tb=%llu", + link->inflight_bytes, link->inflight_packets, link->inflight_lim_bytes, + link->send_blocked_inflight, link->inflight_phase, + link->slow_start_threshold, (unsigned long long)link->last_window_update_tb); + + /* rtt */ + char hist_str[128] = ""; + int hist_pos = 0; + for (int i = 0; i < 10 && i < (int)link->rtt_history_count; i++) { + hist_pos += snprintf(hist_str + hist_pos, sizeof(hist_str) - hist_pos, + "%s%u", i == 0 ? "" : ",", link->rtt_history[i]); + } + DUMP_LINE("RTT: last=%u avg10=%u min=%u max=%u jitter=%u hist=[%s] cnt=%d swm=%u", + link->rtt_last, link->rtt_avg10, link->rtt_min, link->rtt_max_val, + link->jitter, hist_str, link->rtt_history_count, + link->rtt_swm ? swm_get_min(link->rtt_swm) : 0); + + /* tt/rt/bandwidth */ + DUMP_LINE("TT/RT: tt=%u rt=%u recv_dt_tx=%u recv_dt_rx=%u bw=%u", + link->tt_last, link->rt_last, + link->recv_dt_avg_tx, link->recv_dt_avg_rx, + link->bandwidth); + + /* errors */ + DUMP_LINE("MTU: %u local=%u remote=%u", link->mtu, link->mtu_local, link->mtu_remote); + DUMP_LINE("ERRORS: enc=%zu dec=%zu snd=%zu rcv=%zu total_enc=%zu total_dec=%zu retrans=%u", + link->encrypt_errors, link->decrypt_errors, + link->send_errors, link->recv_errors, + link->total_encrypted, link->total_decrypted, + link->total_retransmissions); + + /* nat */ + ip_to_str_buf(link->nat_ip, nat_buf, sizeof(nat_buf)); + DUMP_LINE("NAT: ip=%s:%u chg=%u hits=%u check=%d type=%d", + nat_buf, link->nat_port, + link->nat_changes_count, link->nat_hits_count, + link->nat_check_status, link->nat_type); + + /* stats win */ + DUMP_LINE("WIN: ptr=%u tb=%u tx=%u retrans=%u", + link->win_ptr, link->win_timebase, + link->window_pkt_transmitted, link->window_retransmissions); + + /* last recv */ + DUMP_LINE("LAST_RECV: time=%llu ts=%u updated=%d", + (unsigned long long)link->last_recv_local_time, + link->last_recv_timestamp, link->last_recv_updated); + + /* handshake */ + DUMP_LINE("HANDSHAKE: min=%u max=%u", link->handshake_minsize, link->handshake_maxsize); + + /* burst */ + DUMP_LINE("BURST: active=%d id=%u seq=%u/%u last_tb=%llu resp=%s", + link->burst_active, link->burst_id, + link->burst_seq, link->burst_count, + (unsigned long long)link->burst_last_time_tb, + link->burst_resp_timer ? "ACTIVE" : "free"); + + /* shaper */ + DUMP_LINE("SHAPER: load_tb=%llu sub=%llu state=%d", + (unsigned long long)link->shaper_load_time_tb, + (unsigned long long)link->shaper_sub_nanotime, + link->shaper_state); + + link = link->next; + link_idx++; + } + if (link_idx == 0) DUMP_HDR(" --- NO LINKS ---"); + + DUMP_HDR("=== END [%s] ===", conn->log_name); +} + +void etcp_dump_all_conns(struct UTUN_INSTANCE* instance) { + if (!instance) { + DUMP_HDR("etcp_dump_all_conns: NULL instance"); + return; + } + DUMP_HDR("=== DUMP ALL CONNS for instance node=0x%llx ===", + (unsigned long long)instance->node_id); + struct ETCP_CONN* conn = instance->connections; + int idx = 0; + while (conn) { + DUMP_HDR("--- CONN %d ---", idx); + etcp_dump_conn_state(conn); + conn = conn->next; + idx++; + } + if (idx == 0) DUMP_HDR("--- NO CONNECTIONS ---"); + DUMP_HDR("=== END DUMP ALL ==="); +} diff --git a/src/etcp_dump.h b/src/etcp_dump.h new file mode 100644 index 00000000..0d06b217 --- /dev/null +++ b/src/etcp_dump.h @@ -0,0 +1,10 @@ +#ifndef ETCP_DUMP_H +#define ETCP_DUMP_H + +#include "etcp.h" +#include "utun_instance.h" + +void etcp_dump_conn_state(struct ETCP_CONN* conn); +void etcp_dump_all_conns(struct UTUN_INSTANCE* instance); + +#endif diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index f7df370d..59e04c9d 100644 --- a/src/pkt_normalizer.c +++ b/src/pkt_normalizer.c @@ -180,6 +180,7 @@ void pn_reset(struct PKTNORM* pn) { queue_entry_free(entry); } queue_resume_callback(pn->input); + queue_resume_callback(pn->output); } // Send data to packer (copies and adds to input queue or pending, triggering callback) используется только в юниттесте diff --git a/tests/Makefile.am b/tests/Makefile.am index 771e36aa..12e97ad9 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -8,6 +8,7 @@ check_PROGRAMS = \ test_ipv6_sockets \ test_etcp_minimal \ test_etcp_100_packets \ + test_etcp_reconnect \ test_pkt_normalizer_etcp \ test_pkt_normalizer_standalone \ test_etcp_api \ @@ -72,7 +73,8 @@ ETCP_CORE_OBJS = \ $(top_builddir)/src/utun-etcp_loadbalancer.o \ $(top_builddir)/src/utun-pkt_normalizer.o \ $(top_builddir)/src/utun-etcp_api.o \ - $(top_builddir)/src/utun-etcp_debug.o + $(top_builddir)/src/utun-etcp_debug.o \ + $(top_builddir)/src/utun-etcp_dump.o # Platform-specific TUN objects if OS_WINDOWS @@ -178,6 +180,10 @@ test_etcp_100_packets_SOURCES = test_etcp_100_packets.c test_etcp_100_packets_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source test_etcp_100_packets_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_reconnect_SOURCES = test_etcp_reconnect.c +test_etcp_reconnect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source +test_etcp_reconnect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_pkt_normalizer_etcp_SOURCES = test_pkt_normalizer_etcp.c test_pkt_normalizer_etcp_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source test_pkt_normalizer_etcp_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_reconnect.c b/tests/test_etcp_reconnect.c new file mode 100644 index 00000000..5baab2fe --- /dev/null +++ b/tests/test_etcp_reconnect.c @@ -0,0 +1,377 @@ +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifdef _WIN32 +#include +#include +#include +#define getpid _getpid +#else +#include +#endif +#include +#include + +#include "../src/etcp.h" +#include "../src/etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/utun_instance.h" +#include "../src/routing.h" +#include "../src/tun_if.h" +#include "../src/secure_channel.h" +#include "../src/etcp_dump.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" + +#define TEST_TIMEOUT_MS 120000 +#define PACKET_SIZE 100 +#define TOTAL_PACKETS 500 +#define MAX_QUEUE_SIZE 5 + +static struct UTUN_INSTANCE* server_instance = NULL; +static struct UTUN_INSTANCE* client_instance = NULL; +static struct UASYNC* ua = NULL; + +static char temp_dir[] = "/tmp/utun_test_XXXXXX"; +static char server_config_path[256]; +static char client_config_path[256]; +static int server_port = 0; +static int client_port = 0; + +static int test_completed = 0; +static void* packet_timeout_id = NULL; +static void* global_timeout_id = NULL; + +static int phase = 0; +static int packets_sent = 0; +static int packets_received = 0; +static int restart_action_done = 0; +static int _dump_interval = 50; +static int _dump_timer = 0; +static uint8_t packet_buffer[PACKET_SIZE]; + +static int create_temp_configs(void) { + if (test_mkdtemp(temp_dir) != 0) { + fprintf(stderr, "Failed to create temp directory\n"); + return -1; + } + int base_port = 40000 + (getpid() % 20000); + server_port = base_port; + client_port = base_port + 1; + + snprintf(server_config_path, sizeof(server_config_path), "%s/server.conf", temp_dir); + snprintf(client_config_path, sizeof(client_config_path), "%s/client.conf", temp_dir); + + FILE* f = fopen(server_config_path, "w"); + if (!f) { fprintf(stderr, "Failed to create server config file\n"); return -1; } + fprintf(f, + "[global]\n" + "my_node_id=0x1111111111111111\n" + "my_private_key=67b705a92b41bcaae105af2d6a17743faa7b26ccebba8b3b9b0af05e9cd1d5fb\n" + "my_public_key=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9\n" + "tun_ip=10.99.0.1/24\n" + "tun_ifname=tun99\n" + "\n" + "[server: test]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "\n" + "[allowed_keys]\n" + "allow_all=1\n", + server_port); + fclose(f); + + f = fopen(client_config_path, "w"); + if (!f) { fprintf(stderr, "Failed to create client config file\n"); test_unlink(server_config_path); return -1; } + fprintf(f, + "[global]\n" + "my_node_id=0x2222222222222222\n" + "my_private_key=4813d31d28b7e9829247f488c6be7672f2bdf61b2508333128e386d1759afed2\n" + "my_public_key=c594f33c91f3a2222795c2c110c527bf214ad1009197ce14556cb13df3c461b3c373bed8f205a8dd1fc0c364f90bf471d7c6f5db49564c33e4235d268569ac71\n" + "tun_ip=10.99.0.2/24\n" + "tun_ifname=tun98\n" + "\n" + "[server: test]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "\n" + "[client: test_client]\n" + "keepalive=1\n" + "peer_public_key=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9\n" + "link=test:127.0.0.1:%d\n", + client_port, server_port); + fclose(f); + return 0; +} + +static void cleanup_temp_configs(void) { + if (server_config_path[0]) test_unlink(server_config_path); + if (client_config_path[0]) test_unlink(client_config_path); + if (temp_dir[0]) test_rmdir(temp_dir); +} + +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; + } + return 0; +} + +static void drain_received(int count_flag) { + if (!server_instance) return; + struct ETCP_CONN* conn = server_instance->connections; + while (conn) { + if (conn->output_queue) { + queue_set_callback(conn->output_queue, NULL, NULL); + struct ETCP_FRAGMENT* pkt; + while ((pkt = (struct ETCP_FRAGMENT*)queue_data_get(conn->output_queue)) != NULL) { + if (count_flag) packets_received++; + if (pkt->ll.dgram) memory_pool_free(conn->instance->data_pool, pkt->ll.dgram); + queue_entry_free((struct ll_entry*)pkt); + } + } + conn = conn->next; + } +} + +static void send_packets(void) { + if (!client_instance || packets_sent >= TOTAL_PACKETS) return; + struct ETCP_CONN* conn = client_instance->connections; + if (!conn || !conn->input_queue) return; + while (packets_sent < TOTAL_PACKETS) { + if (queue_entry_count(conn->input_queue) >= MAX_QUEUE_SIZE) break; + packet_buffer[0] = (uint8_t)(packets_sent & 0xFF); + for (int i = 1; i < PACKET_SIZE; i++) packet_buffer[i] = (uint8_t)((packets_sent + i) % 256); + if (etcp_int_send(conn, packet_buffer, PACKET_SIZE) != 0) break; + packets_sent++; + } +} + +static void monitor(void* arg) { + (void)arg; + if (test_completed) { packet_timeout_id = NULL; return; } + + 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; + if (conn) cur_reinit = conn->reinit_count; + + switch (phase) { + case 0: // Wait for initial connection + drain_received(0); + if (conn_ok) { + printf("=== Phase 0: Connection established ===\n"); + packets_sent = 0; + packets_received = 0; + phase = 1; + } + break; + + case 1: // Send 500 packets client→server + send_packets(); + drain_received(1); + if (packets_received >= TOTAL_PACKETS) { + printf("=== Phase 1 done: sent=%d received=%d ===\n", packets_sent, packets_received); + packets_received = 0; + phase = 2; + } + break; + + case 2: // Drain and prepare for server restart + drain_received(0); + packets_received = 0; + packets_sent = 0; + restart_action_done = 0; + printf("=== Phase 2: Starting server restart ===\n"); + phase = 3; + break; + + case 3: // Server restart + wait for reconnect + if (!restart_action_done) { + printf("Destroying server instance...\n"); + server_instance->running = 0; + utun_instance_destroy(server_instance); + server_instance = NULL; + printf("Recreating server instance...\n"); + server_instance = utun_instance_create(ua, server_config_path); + if (!server_instance || utun_instance_init(server_instance) < 0) { + fprintf(stderr, "Failed to recreate server instance\n"); + test_completed = 2; + return; + } + printf("Server recreated (node_id=%llx)\n", (unsigned long long)server_instance->node_id); + restart_action_done = 1; + } + drain_received(0); + if (conn_ok) { + printf("=== Phase 3: Reconnected after server restart ===\n"); + packets_sent = 0; + packets_received = 0; + last_reinit_count = cur_reinit; + phase = 4; + } + break; + + case 4: // Send 500 packets after server restart + if (cur_reinit != last_reinit_count) { + printf("=== Phase 4: Reinit detected (%d -> %d), restarting send ===\n", last_reinit_count, cur_reinit); + last_reinit_count = cur_reinit; + packets_sent = 0; + packets_received = 0; + drain_received(0); + } + send_packets(); + drain_received(1); + if (packets_received >= TOTAL_PACKETS) { + printf("=== Phase 4 done: sent=%d received=%d ===\n", packets_sent, packets_received); + packets_received = 0; + _dump_interval = 50; + phase = 5; + } + if (packets_received > 0 && packets_received < TOTAL_PACKETS) { + if (++_dump_timer >= _dump_interval) { + printf("--- DUMP at phase=4 (stuck, recv=%d/%d) ---\n", packets_received, TOTAL_PACKETS); + etcp_dump_all_conns(client_instance); + etcp_dump_all_conns(server_instance); + _dump_timer = 0; + if (_dump_interval < 800) _dump_interval *= 2; + } + } + break; + + case 5: // Drain and prepare for client restart + drain_received(0); + packets_received = 0; + packets_sent = 0; + restart_action_done = 0; + printf("=== Phase 5: Starting client restart ===\n"); + phase = 6; + break; + + case 6: // Client restart + wait for reconnect + if (!restart_action_done) { + printf("Destroying client instance...\n"); + client_instance->running = 0; + utun_instance_destroy(client_instance); + client_instance = NULL; + printf("Recreating client instance...\n"); + client_instance = utun_instance_create(ua, client_config_path); + if (!client_instance || utun_instance_init(client_instance) < 0) { + fprintf(stderr, "Failed to recreate client instance\n"); + test_completed = 2; + return; + } + printf("Client recreated (node_id=%llx)\n", (unsigned long long)client_instance->node_id); + restart_action_done = 1; + } + drain_received(0); + conn_ok = is_connection_established(client_instance); + if (conn_ok) { + printf("=== Phase 6: Reconnected after client restart ===\n"); + packets_sent = 0; + packets_received = 0; + phase = 7; + } + break; + + case 7: // Send 500 packets after client restart + if (cur_reinit != last_reinit_count) { + printf("=== Phase 7: Reinit detected (%d -> %d), restarting send ===\n", last_reinit_count, cur_reinit); + last_reinit_count = cur_reinit; + packets_sent = 0; + packets_received = 0; + drain_received(0); + } + send_packets(); + drain_received(1); + if (packets_received >= TOTAL_PACKETS) { + printf("=== Phase 7 done: sent=%d received=%d ===\n", packets_sent, packets_received); + test_completed = 1; + return; + } + break; + } + + if (!test_completed) + packet_timeout_id = uasync_set_timeout(ua, 10, NULL, monitor, "test_monitor"); +} + +static void test_timeout(void* arg) { + (void)arg; + if (!test_completed) { + printf("\n=== TEST TIMEOUT at phase %d: sent=%d recv=%d ===\n", phase, packets_sent, packets_received); + test_completed = 2; + if (packet_timeout_id) { uasync_cancel_timeout(ua, packet_timeout_id); packet_timeout_id = NULL; } + } +} + +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); + + utun_instance_set_tun_init_enabled(0); + + ua = uasync_create(); + server_instance = utun_instance_create(ua, server_config_path); + if (!server_instance || utun_instance_init(server_instance) < 0) { + fprintf(stderr, "Failed to create server\n"); + return 1; + } + printf("Server ready (node_id=%llx)\n", (unsigned long long)server_instance->node_id); + + client_instance = utun_instance_create(ua, client_config_path); + if (!client_instance || utun_instance_init(client_instance) < 0) { + fprintf(stderr, "Failed to create client\n"); + utun_instance_destroy(server_instance); + return 1; + } + printf("Client ready (node_id=%llx)\n", (unsigned long long)client_instance->node_id); + + 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"); + + int elapsed = 0; + int poll_interval = 50; + int iter = 0; + while (!test_completed && elapsed < TEST_TIMEOUT_MS * 10 + 50000) { + uasync_poll(ua, poll_interval); + elapsed += poll_interval; + if (++iter % 200 == 1) printf("main_loop: iter=%d elapsed=%d phase=%d sent=%d recv=%d\n", + iter, elapsed, phase, packets_sent, packets_received); + } + printf("Main loop exit: test_completed=%d phase=%d elapsed=%d iter=%d\n", + test_completed, phase, elapsed, iter); + + if (packet_timeout_id) uasync_cancel_timeout(ua, packet_timeout_id); + if (global_timeout_id) uasync_cancel_timeout(ua, global_timeout_id); + + 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; } + if (ua) { uasync_destroy(ua, 0); ua = NULL; } + cleanup_temp_configs(); + + if (test_completed == 1) { + printf("=== TEST PASSED: Reconnect works after server and client restart ===\n"); + return 0; + } + printf("=== TEST FAILED: timeout or error at phase %d ===\n", phase); + return 1; +}