Browse Source

fix: reconnect test — cancel timers in conn_reset, add init_timer cancel on link_state=3, add reinit resilience to test

- etcp_conn_reset: cancel retrans_timer/ack_resp_timer via uasync_cancel_timeout
  before nulling (was leaking timer handles, confirmed by 'Timer leak: diff')
- etcp_connections: cancel link->init_timer explicitly when link_state becomes 3
  on both client and server INIT_RESPONSE paths (prevents spurious INIT after
  connection is established that could trigger cascading reinits)
- etcp_connections: remove duplicate conn->session_id assignment in new conn path
  (session_id is set later via conn_reinit flow)
- pkt_normalizer: resume output queue callback after pn_reset
- Add test_etcp_reconnect with reinit detection — monitors conn->reinit_count
  and restarts data transfer if reinit occurs during active phase
- Add etcp_dump.c/h — ETCP connection state dump utility for debugging
congestion
Evgeny 5 months ago
parent
commit
dea3a39598
  1. 1
      src/Makefile.am
  2. 8
      src/etcp.c
  3. 9
      src/etcp_connections.c
  4. 255
      src/etcp_dump.c
  5. 10
      src/etcp_dump.h
  6. 1
      src/pkt_normalizer.c
  7. 8
      tests/Makefile.am
  8. 377
      tests/test_etcp_reconnect.c

1
src/Makefile.am

@ -21,6 +21,7 @@ utun_CORE_SOURCES = \
etcp_connections.c \ etcp_connections.c \
etcp_loadbalancer.c \ etcp_loadbalancer.c \
etcp_debug.c \ etcp_debug.c \
etcp_dump.c \
secure_channel.c \ secure_channel.c \
crc32.c \ crc32.c \
pkt_normalizer.c \ pkt_normalizer.c \

8
src/etcp.c

@ -290,9 +290,15 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) {
l = l->next; l = l->next;
} }
// Reset timers (just clear the pointers - timers will expire naturally) // Cancel active timers to prevent memory leaks
if (etcp->retrans_timer) {
uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer);
etcp->retrans_timer = NULL; 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->ack_resp_timer = NULL;
}
etcp->reset_count++; etcp->reset_count++;

9
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; } if (!conn) { errorcode=55; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "failed to create connection"); goto ec_fr; }
memcpy(&conn->crypto_ctx, &sc, sizeof(sc)); memcpy(&conn->crypto_ctx, &sc, sizeof(sc));
conn->peer_node_id=peer_id; conn->peer_node_id=peer_id;
conn->session_id = session_id;
etcp_update_log_name(conn); 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_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); 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); memory_pool_free(e_sock->instance->pkt_pool, pkt);
link->initialized = 1; link->initialized = 1;
link->link_state = 3; 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) { if (link->etcp->initialized == 0) {
etcp_conn_ready(link->etcp); 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); 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->initialized = 1;// получен init response (client)
link->link_state = 3; // connected 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) { if (code == ETCP_INIT_RESPONSE) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit 3 %p", link->etcp); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit 3 %p", link->etcp);

255
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 <stdio.h>
#include <string.h>
#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 ===");
}

10
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

1
src/pkt_normalizer.c

@ -180,6 +180,7 @@ void pn_reset(struct PKTNORM* pn) {
queue_entry_free(entry); queue_entry_free(entry);
} }
queue_resume_callback(pn->input); queue_resume_callback(pn->input);
queue_resume_callback(pn->output);
} }
// Send data to packer (copies and adds to input queue or pending, triggering callback) используется только в юниттесте // Send data to packer (copies and adds to input queue or pending, triggering callback) используется только в юниттесте

8
tests/Makefile.am

@ -8,6 +8,7 @@ check_PROGRAMS = \
test_ipv6_sockets \ test_ipv6_sockets \
test_etcp_minimal \ test_etcp_minimal \
test_etcp_100_packets \ test_etcp_100_packets \
test_etcp_reconnect \
test_pkt_normalizer_etcp \ test_pkt_normalizer_etcp \
test_pkt_normalizer_standalone \ test_pkt_normalizer_standalone \
test_etcp_api \ test_etcp_api \
@ -72,7 +73,8 @@ ETCP_CORE_OBJS = \
$(top_builddir)/src/utun-etcp_loadbalancer.o \ $(top_builddir)/src/utun-etcp_loadbalancer.o \
$(top_builddir)/src/utun-pkt_normalizer.o \ $(top_builddir)/src/utun-pkt_normalizer.o \
$(top_builddir)/src/utun-etcp_api.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 # Platform-specific TUN objects
if OS_WINDOWS 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_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_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_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_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) test_pkt_normalizer_etcp_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)

377
tests/test_etcp_reconnect.c

@ -0,0 +1,377 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "../lib/platform_compat.h"
#include "test_utils.h"
#ifdef _WIN32
#include <windows.h>
#include <direct.h>
#include <process.h>
#define getpid _getpid
#else
#include <unistd.h>
#endif
#include <time.h>
#include <sys/stat.h>
#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;
}
Loading…
Cancel
Save