diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 00d441b1..d2223564 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -1255,7 +1255,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) { int udp_ov = (dgram->link->remote_addr.ss_family == AF_INET6) ? IP_UDP_OVERHEAD_V6 : IP_UDP_OVERHEAD_V4; int wire_ov = udp_ov + SC_ENCRYPT_OVERHEAD; if (len<0 || len + (int)dgram->noencrypt_len > (int)(dgram->link->mtu - wire_ov)) { dgram->link->send_errors++; errcode=1; goto es_err; } - uint8_t enc_buf[1600]; + uint8_t enc_buf[PACKET_DATA_MAX_MTU];// должно вмещать mtu - udp_ov (до 2020 байт) size_t enc_buf_len=0; dgram->timestamp=get_current_timestamp(); diff --git a/tests/test_etcp_link_stress.c b/tests/test_etcp_link_stress.c index 173493f6..3e00e6c5 100644 --- a/tests/test_etcp_link_stress.c +++ b/tests/test_etcp_link_stress.c @@ -1,22 +1,14 @@ +// test_etcp_link_stress.c - Multilink test: добавляет линки в 5 вариантах (V1/V2 one-way, +// V3a/b/c коллизия), проверяет per-link трафик и инвариант «нет реинита пока жив хотя бы один линк». #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 "etcp.h" #include "etcp_connections.h" -#include "stcp_link.h" +#include "etcp_dump.h" #include "../src/config_parser.h" #include "../src/utun_instance.h" #include "routing.h" @@ -27,15 +19,14 @@ #include "../lib/debug_config.h" #include "../lib/mem.h" -#define CONN_TIMEOUT_TB 50000 -#define STRESS_DURATION_MS 2000 -#define CYCLE_INTERVAL_MS 100 -#define PACKET_INTERVAL_TB 10 -#define PACKET_SIZE 64 -#define MAX_EXTRA 8 -#define UDP_POOL_SIZE 12 -#define TCP_POOL_SIZE 4 -#define INITIAL_SOCKET_NAME "udp0" +#define PACKET_SIZE 64 +#define TRAFFIC_PACKETS 200 +#define LINK_UP_TIMEOUT_MS 8000 +#define TRAFFIC_TIMEOUT_MS 8000 +#define LINK_DOWN_TIMEOUT_MS 8000 +#define COLLISION_TIMEOUT_MS 15000 + +enum { COL_A = 0, COL_B, COL_C }; static struct UTUN_INSTANCE* server_instance = NULL; static struct UTUN_INSTANCE* client_instance = NULL; @@ -46,53 +37,24 @@ static char server_config_path[256]; static char client_config_path[256]; static int server_port = 0; static int client_port = 0; -static int udp_pool_base = 0; -static int tcp_pool_base = 0; +static int next_port = 0; static uint32_t packets_sent = 0; static uint32_t packets_received = 0; -static uint64_t last_received_bytes = 0; -static int cycle_count = 0; -static int hang_detected = 0; -static uint64_t stress_start_tb = 0; - -struct extra_info { - int type; - int port; - struct ETCP_SOCKET* udp_sock; - struct ETCP_SOCKET* tcp_sock; - struct stcp_server* tcp_srv; - struct ETCP_LINK* client_link; - uint64_t create_tb; - uint64_t bytes_decrypted; - uint64_t bytes_encrypted; - uint64_t acked_bytes; - uint64_t prev_decrypted; - uint32_t retrans; - uint16_t rtt; - uint32_t inflight; -}; -static struct extra_info extras[MAX_EXTRA]; -static int extra_count = 0; - -static int used_udp_ports[UDP_POOL_SIZE]; -static int used_tcp_ports[TCP_POOL_SIZE]; -static int udp_port_idx = 0; -static int tcp_port_idx = 0; - -/* struct stcp_server layout from stcp_link.c (opaque in public header) */ -struct stcp_link_server_local { - struct stcp_link_server_local *next; - struct stcp_server *srv; -}; + +// globals for generic poll_until conditions +static struct ETCP_LINK* g_link = NULL; +static struct ETCP_SOCKET* g_sock = NULL; +static struct ETCP_CONN* g_conn = NULL; +static uint16_t g_port = 0; +static uint32_t g_reinit_cli0 = 0; 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 = 50000 + (getpid() % 15000); server_port = base_port; client_port = base_port + 1; - udp_pool_base = base_port + 10; - tcp_pool_base = base_port + 30; + next_port = base_port + 100; 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); @@ -106,14 +68,16 @@ static int create_temp_configs(void) { "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" "tun_ip=10.99.0.1/24\n" "tun_ifname=tun99\n" + "keepalive_interval=100\n" + "keepalive_adaptive=0\n" "\n" - "[server: %s]\n" + "[server: udp0]\n" "addr=127.0.0.1:%d\n" "type=public\n" "\n" "[allowed_keys]\n" "allow_all=1\n", - INITIAL_SOCKET_NAME, server_port); + server_port); fclose(f); f = fopen(client_config_path, "w"); @@ -125,16 +89,21 @@ static int create_temp_configs(void) { "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" "tun_ip=10.99.0.2/24\n" "tun_ifname=tun98\n" + "keepalive_interval=100\n" + "keepalive_adaptive=0\n" "\n" - "[server: %s]\n" + "[server: udp0]\n" "addr=127.0.0.1:%d\n" "type=public\n" "\n" "[client: test_client]\n" "keepalive=1\n" "peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" - "link=%s:127.0.0.1:%d\n", - INITIAL_SOCKET_NAME, client_port, INITIAL_SOCKET_NAME, server_port); + "link=udp0:127.0.0.1:%d\n" + "\n" + "[allowed_keys]\n" + "allow_all=1\n", + client_port, server_port); fclose(f); return 0; } @@ -145,205 +114,356 @@ static void cleanup_temp_configs(void) { if (temp_dir[0]) test_rmdir(temp_dir); } +// ===== helpers ===== + static struct ETCP_CONN* get_conn(struct UTUN_INSTANCE* inst) { if (!inst || !inst->connections || !inst->connections->head) return NULL; struct conn_queue_entry* ce = (struct conn_queue_entry*)inst->connections->head->data; return ce ? ce->conn : NULL; } -static int is_connected(struct UTUN_INSTANCE* inst) { +static int conn_up(struct UTUN_INSTANCE* inst) { struct ETCP_CONN* conn = get_conn(inst); - if (!conn) return 0; - struct ETCP_LINK* l = conn->links; - while (l) { if (l->initialized && l->link_state == 3) return 1; l = l->next; } - return 0; + return conn && conn->initialized && conn->links_up; } -static void drain_received(int count_flag) { - if (!server_instance) return; - struct ETCP_CONN* conn = get_conn(server_instance); - if (!conn || !conn->output_queue) return; - 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); - } +static int count_links(struct ETCP_CONN* conn) { + int n = 0; + for (struct ETCP_LINK* l = conn ? conn->links : NULL; l; l = l->next) n++; + return n; } -static int get_free_udp_port(void) { - for (int i = 0; i < UDP_POOL_SIZE; i++) { - if (!used_udp_ports[i]) { used_udp_ports[i] = 1; return udp_pool_base + i; } - } - return -1; +static int count_links_up(struct ETCP_CONN* conn) { + int n = 0; + for (struct ETCP_LINK* l = conn ? conn->links : NULL; l; l = l->next) if (l->link_status) n++; + return n; } -static void release_udp_port(int port) { - for (int i = 0; i < UDP_POOL_SIZE; i++) { if (udp_pool_base + i == port) { used_udp_ports[i] = 0; return; } } +static void mk_addr(struct sockaddr_storage* out, uint16_t port) { + memset(out, 0, sizeof(*out)); + struct sockaddr_in* sin = (struct sockaddr_in*)out; + sin->sin_family = AF_INET; + sin->sin_addr.s_addr = htonl(INADDR_LOOPBACK); + sin->sin_port = htons(port); } -static int get_free_tcp_port(void) { - for (int i = 0; i < TCP_POOL_SIZE; i++) { - if (!used_tcp_ports[i]) { used_tcp_ports[i] = 1; return tcp_pool_base + i; } - } - return -1; -} +static int alloc_port(void) { return next_port++; } -static void release_tcp_port(int port) { - for (int i = 0; i < TCP_POOL_SIZE; i++) { if (tcp_pool_base + i == port) { used_tcp_ports[i] = 0; return; } } +static struct ETCP_SOCKET* add_udp_socket(struct UTUN_INSTANCE* inst, uint16_t port, const char* name) { + struct CFG_SERVER cfg; memset(&cfg, 0, sizeof(cfg)); + snprintf(cfg.name, sizeof(cfg.name), "%s", name); + cfg.type = CFG_SERVER_TYPE_PUBLIC; + cfg.mtu = PACKET_DATA_MAX_MTU; + mk_addr(&cfg.ip, port); + return etcp_socket_add(inst, &cfg); } -static void collect_link_stats(struct ETCP_LINK* link, struct extra_info* info) { - if (!link) return; - info->bytes_decrypted = link->total_decrypted; - info->bytes_encrypted = link->total_encrypted; - info->acked_bytes = link->acked_bytes; - info->retrans = link->total_retransmissions; - info->rtt = link->rtt_last; - info->inflight = link->inflight_bytes; +static struct ETCP_LINK* find_link_by_remote_port(struct ETCP_CONN* conn, uint16_t port) { + for (struct ETCP_LINK* l = conn ? conn->links : NULL; l; l = l->next) { + if (l->remote_addr.ss_family == AF_INET && + ntohs(((struct sockaddr_in*)&l->remote_addr)->sin_port) == port) + return l; + } + return NULL; } -static void remove_udp_sock_from_instance(struct ETCP_SOCKET* sock) { - if (!server_instance || !sock) return; - struct ETCP_SOCKET** pp = &server_instance->etcp_sockets; - while (*pp) { if (*pp == sock) { *pp = sock->next; etcp_socket_remove(sock); return; } pp = &(*pp)->next; } +static struct ETCP_LINK* first_link_on_socket(struct ETCP_SOCKET* sock) { + if (!sock || !sock->links_queue || !sock->links_queue->head) return NULL; + struct link_queue_entry* lqe = (struct link_queue_entry*)sock->links_queue->head->data; + return lqe ? lqe->link : NULL; } -static void collect_and_remove_oldest(void) { - if (extra_count == 0) return; - int oldest_idx = -1; - uint64_t oldest_tb = 0; - for (int i = 0; i < extra_count; i++) { - if (oldest_idx < 0 || extras[i].create_tb < oldest_tb) { oldest_idx = i; oldest_tb = extras[i].create_tb; } - } - if (oldest_idx < 0) return; - - struct extra_info* info = &extras[oldest_idx]; - uint64_t lifetime_tb = get_time_tb() - info->create_tb; - printf("[CYCLE %d] REMOVE %s: port=%d lifetime=%llums decr=%llu encr=%llu ack=%llu retrans=%u rtt=%u infl=%u\n", - cycle_count, info->type ? "TCP" : "UDP", info->port, - (unsigned long long)(lifetime_tb / 10), - (unsigned long long)info->bytes_decrypted, - (unsigned long long)info->bytes_encrypted, - (unsigned long long)info->acked_bytes, - info->retrans, info->rtt, info->inflight); - - if (info->client_link) { etcp_link_close(info->client_link); info->client_link = NULL; } - if (info->udp_sock) { remove_udp_sock_from_instance(info->udp_sock); info->udp_sock = NULL; release_udp_port(info->port); } - if (info->tcp_srv) { - struct stcp_link_server_local** pp = (struct stcp_link_server_local**)&server_instance->stcp_servers; - while (*pp) { if (*pp == (struct stcp_link_server_local*)info->tcp_srv) { *pp = (*pp)->next; break; } pp = &(*pp)->next; } - stcp_link_server_destroy(info->tcp_srv); - info->tcp_srv = NULL; - } - if (info->tcp_sock) { tcp_socket_remove(info->tcp_sock); info->tcp_sock = NULL; release_tcp_port(info->port); } - - if (oldest_idx < extra_count - 1) - memmove(&extras[oldest_idx], &extras[oldest_idx + 1], (extra_count - oldest_idx - 1) * sizeof(struct extra_info)); - extra_count--; - memset(&extras[extra_count], 0, sizeof(struct extra_info)); +static uint64_t client_link_acked(uint16_t srv_port) { + struct ETCP_LINK* l = find_link_by_remote_port(get_conn(client_instance), srv_port); + return l ? l->acked_packets : 0; } -static struct ETCP_SOCKET* get_client_socket(void) { - if (client_instance && client_instance->etcp_sockets) return client_instance->etcp_sockets; - return NULL; +static void dump_state(void) { + printf("--- STATE DUMP ---\n"); + if (server_instance) etcp_dump_all_conns(server_instance); + if (client_instance) etcp_dump_all_conns(client_instance); + printf("--- END DUMP ---\n"); + fflush(stdout); } -static int add_udp_pair(struct ETCP_CONN* conn) { - if (extra_count >= MAX_EXTRA) { collect_and_remove_oldest(); if (extra_count >= MAX_EXTRA) return -1; } - int port = get_free_udp_port(); - if (port < 0) { fprintf(stderr, "No free UDP ports\n"); return -1; } - - struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); - sin.sin_family = AF_INET; sin.sin_addr.s_addr = htonl(INADDR_LOOPBACK); sin.sin_port = htons(port); - - struct CFG_SERVER cfg; memset(&cfg, 0, sizeof(cfg)); - snprintf(cfg.name, sizeof(cfg.name), "dyn_udp_%d", port); - cfg.type = CFG_SERVER_TYPE_PUBLIC; cfg.mtu = PACKET_DATA_MAX_MTU; - memcpy(&cfg.ip, &sin, sizeof(sin)); - - struct ETCP_SOCKET* es = etcp_socket_add(server_instance, &cfg); - if (!es) { release_udp_port(port); fprintf(stderr, "Failed to create UDP socket on port %d\n", port); return -1; } - - struct ETCP_SOCKET* client_sock = get_client_socket(); - if (!client_sock) { remove_udp_sock_from_instance(es); release_udp_port(port); fprintf(stderr, "No client socket\n"); return -1; } - - struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); - struct ETCP_LINK* link = etcp_link_new(conn, client_sock, &sa, 0); - if (!link) { remove_udp_sock_from_instance(es); release_udp_port(port); return -1; } - - struct extra_info* info = &extras[extra_count++]; memset(info, 0, sizeof(*info)); - info->type = 0; info->port = port; info->udp_sock = es; info->client_link = link; - info->create_tb = get_time_tb(); - printf("[CYCLE %d] ADD UDP: port=%d link=%p sock=%p\n", cycle_count, port, (void*)link, (void*)es); - return 0; +#define CHECK(cond, ...) do { if (!(cond)) { \ + printf("FAIL %s:%d: ", __func__, __LINE__); \ + printf(__VA_ARGS__); \ + printf("\n"); \ + dump_state(); \ + goto fail; } } while (0) + +// ===== poll conditions ===== + +static int c_conn_up_both(void) { return conn_up(server_instance) && conn_up(client_instance); } +static int c_link_up(void) { return g_link && g_link->initialized && g_link->link_status == 1; } +static int c_link_down(void) { return g_link && g_link->link_status == 0; } +static int c_socket_link_up(void) { struct ETCP_LINK* l = first_link_on_socket(g_sock); return l && l->initialized && l->link_status == 1; } +static int c_link_on_port_up(void) { struct ETCP_LINK* l = find_link_by_remote_port(g_conn, g_port); return l && l->initialized && l->link_status == 1; } +static int c_client_all_down(void) { struct ETCP_CONN* c = get_conn(client_instance); return c && c->links_up == 0; } +static int c_server_all_down(void) { struct ETCP_CONN* c = get_conn(server_instance); return c && c->links_up == 0; } +static int c_client_reinit_happened(void) { struct ETCP_CONN* c = get_conn(client_instance); return c && c->reinit_count > g_reinit_cli0; } + +static int poll_until(int (*cond)(void), uint64_t timeout_ms, const char* what) { + uint64_t deadline = get_time_tb() + timeout_ms * 10; + while (get_time_tb() < deadline) { + if (cond()) return 0; + uasync_poll(ua, 5); + } + printf("TIMEOUT waiting for: %s (%llums)\n", what, (unsigned long long)timeout_ms); + return -1; } -static int add_tcp_pair(struct ETCP_CONN* conn) { - if (extra_count >= MAX_EXTRA) { collect_and_remove_oldest(); if (extra_count >= MAX_EXTRA) return -1; } - int port = get_free_tcp_port(); - if (port < 0) return 0; - - struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); - sin.sin_family = AF_INET; sin.sin_addr.s_addr = htonl(INADDR_LOOPBACK); sin.sin_port = htons(port); +// ===== traffic ===== - struct CFG_SERVER cfg; memset(&cfg, 0, sizeof(cfg)); - snprintf(cfg.name, sizeof(cfg.name), "dyn_tcp_%d", port); - cfg.type = CFG_SERVER_TYPE_PUBLIC; cfg.transport = 1; - memcpy(&cfg.ip, &sin, sizeof(sin)); - - struct ETCP_SOCKET* ts = tcp_socket_add(server_instance, &cfg); - if (!ts) { release_tcp_port(port); return -1; } - - struct stcp_link_config scfg = {.inst = server_instance, .listen_family = AF_INET}; - struct stcp_server* srv = stcp_server_listen(&scfg, port, (stcp_server_on_link_cb)tcp_server_on_link, ts); - if (!srv) { tcp_socket_remove(ts); release_tcp_port(port); return -1; } - stcp_server_list_add(server_instance, srv); +static void drain_received(int count_flag) { + if (!server_instance) return; + struct ETCP_CONN* conn = get_conn(server_instance); + if (!conn || !conn->output_queue) return; + 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); + } +} - struct ETCP_LINK* tlink = etcp_link_new(conn, NULL, NULL, 0); - if (!tlink) { tcp_socket_remove(ts); release_tcp_port(port); return -1; } - tlink->is_tcp = 1; tlink->is_server = 0; +static int run_traffic(int n, uint64_t timeout_ms) { + struct ETCP_CONN* cc = get_conn(client_instance); + if (!cc || !cc->input_queue) return -1; + uint32_t target = packets_sent + (uint32_t)n; + uint64_t deadline = get_time_tb() + timeout_ms * 10; + while (1) { + while (packets_sent < target && queue_entry_count(cc->input_queue) < 100) { + uint8_t buf[PACKET_SIZE]; + uint32_t s = packets_sent; + buf[0] = (uint8_t)(s & 0xFF); + buf[1] = (uint8_t)((s >> 8) & 0xFF); + buf[2] = (uint8_t)((s >> 16) & 0xFF); + buf[3] = (uint8_t)((s >> 24) & 0xFF); + for (int i = 4; i < PACKET_SIZE; i++) buf[i] = (uint8_t)((s + i) % 256); + if (etcp_int_send(cc, buf, PACKET_SIZE) != 0) break; + packets_sent++; + } + drain_received(1); + if (packets_received >= target) return 0; + if (get_time_tb() > deadline) { + printf("traffic timeout: sent=%u recv=%u target=%u\n", packets_sent, packets_received, target); + return -1; + } + uasync_poll(ua, 5); + } +} - struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); - etcp_tcp_link_start_connect(tlink, &sa, port); +// ===== instance lifecycle ===== - struct extra_info* info = &extras[extra_count++]; memset(info, 0, sizeof(*info)); - info->type = 1; info->port = port; info->tcp_sock = ts; info->tcp_srv = srv; info->client_link = tlink; - info->create_tb = get_time_tb(); - printf("[CYCLE %d] ADD TCP: port=%d link=%p sock=%p srv=%p\n", cycle_count, port, (void*)tlink, (void*)ts, (void*)srv); +static int start_instances(void) { + server_instance = utun_instance_create(ua, server_config_path); + if (!server_instance || utun_instance_init(server_instance) < 0) { + printf("failed to create server\n"); return -1; + } + client_instance = utun_instance_create(ua, client_config_path); + if (!client_instance || utun_instance_init(client_instance) < 0) { + printf("failed to create client\n"); return -1; + } + if (poll_until(c_conn_up_both, LINK_UP_TIMEOUT_MS, "initial connection") != 0) { + printf("initial connection timeout\n"); dump_state(); return -1; + } + printf("connection established: server links_up=%d/%d client links_up=%d/%d\n", + count_links_up(get_conn(server_instance)), count_links(get_conn(server_instance)), + count_links_up(get_conn(client_instance)), count_links(get_conn(client_instance))); return 0; } -static const char* fmt_bytes(uint64_t b) { - static char buf[6][32]; static int n = 0; - char* s = buf[n]; n = (n + 1) % 6; - if (b >= 1024) snprintf(s, 32, "%lluk", (unsigned long long)(b / 1024)); else snprintf(s, 32, "%llu", (unsigned long long)b); - return s; +static void stop_instances(void) { + 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; } } -static const char* fmt_bytes_delta(uint64_t b) { - static char buf[6][32]; static int n = 0; - char* s = buf[n]; n = (n + 1) % 6; - if (b >= 1024) snprintf(s, 32, "+%lluk", (unsigned long long)(b / 1024)); else snprintf(s, 32, "+%llu", (unsigned long long)b); - return s; +// Убить серверный сокет (линк со стороны сервера закрывается сразу), ускорить детект на клиенте. +static void kill_server_socket(struct ETCP_SOCKET* sock, uint16_t srv_port) { + struct ETCP_LINK* cl = find_link_by_remote_port(get_conn(client_instance), srv_port); + if (cl) cl->last_recv_local_time = 0; + etcp_socket_remove(sock); } -static const char* fmt_uint(uint32_t v) { - static char buf[6][16]; static int n = 0; - char* s = buf[n]; n = (n + 1) % 6; - snprintf(s, 16, "%u", v); - return s; +// ===== Session A: one-way adds + kill failover ===== + +static int session_a(void) { + printf("\n=== Session A: one-way adds (V1/V2) + link kill failover ===\n"); + + if (start_instances() != 0) return -1; + + struct ETCP_CONN* sc = get_conn(server_instance); + struct ETCP_CONN* cc = get_conn(client_instance); + CHECK(sc && cc, "no conn after start"); + + uint32_t srv_reinit0 = sc->reinit_count; + uint32_t cli_reinit0 = cc->reinit_count; + struct ETCP_SOCKET* s0 = server_instance->etcp_sockets; // начальный серверный сокет udp0 + struct ETCP_SOCKET* c0 = client_instance->etcp_sockets; // клиентский сокет udp0 (source для V1) + + // -- per-link трафик на исходном линке -- + printf("[A] traffic on initial link...\n"); + uint64_t ack0 = client_link_acked(server_port); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "initial traffic failed"); + CHECK(client_link_acked(server_port) > ack0, "initial link carried no traffic"); + CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit during initial traffic"); + + // -- V1: клиент дозванивается до нового серверного сокета S1 -- + int sp1 = alloc_port(); + struct ETCP_SOCKET* s1 = add_udp_socket(server_instance, sp1, "dyn_s1"); + CHECK(s1, "V1: failed to add server socket S1"); + struct sockaddr_storage s1_addr; mk_addr(&s1_addr, sp1); + struct ETCP_LINK* lc1 = etcp_link_new(cc, c0, &s1_addr, 0); + CHECK(lc1, "V1: failed to create client link"); + printf("[A] V1: client link %d -> server port %d\n", lc1->local_link_id, sp1); + g_link = lc1; + CHECK(poll_until(c_link_up, LINK_UP_TIMEOUT_MS, "V1 link up (client)") == 0, "V1 link did not come up"); + g_sock = s1; + CHECK(poll_until(c_socket_link_up, LINK_UP_TIMEOUT_MS, "V1 link up (server)") == 0, "V1 server link did not come up"); + CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit during V1"); + ack0 = client_link_acked(sp1); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "V1 traffic failed"); + CHECK(client_link_acked(sp1) > ack0, "V1 link carried no traffic"); + + // -- V2: сервер дозванивается до клиента -- + int sp2 = alloc_port(); + struct ETCP_SOCKET* s2 = add_udp_socket(server_instance, sp2, "dyn_s2"); + CHECK(s2, "V2: failed to add server socket S2"); + struct sockaddr_storage c0_addr; mk_addr(&c0_addr, client_port); + struct ETCP_LINK* ls2 = etcp_link_new(sc, s2, &c0_addr, 0); + CHECK(ls2, "V2: failed to create server outbound link"); + printf("[A] V2: server link %d -> client port %d\n", ls2->local_link_id, client_port); + g_link = ls2; + CHECK(poll_until(c_link_up, LINK_UP_TIMEOUT_MS, "V2 link up (server)") == 0, "V2 link did not come up"); + g_conn = cc; g_port = sp2; + CHECK(poll_until(c_link_on_port_up, LINK_UP_TIMEOUT_MS, "V2 link up (client)") == 0, "V2 client link did not come up"); + CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit during V2"); + ack0 = client_link_acked(sp2); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "V2 traffic failed"); + CHECK(client_link_acked(sp2) > ack0, "V2 link carried no traffic"); + + printf("[A] after adds: server links_up=%d/%d client links_up=%d/%d\n", + count_links_up(sc), count_links(sc), count_links_up(cc), count_links(cc)); + + // -- kill V2 link: failover, без реинита -- + printf("[A] killing V2 link (server socket %d)...\n", sp2); + kill_server_socket(s2, sp2); + g_link = find_link_by_remote_port(cc, sp2); + CHECK(poll_until(c_link_down, LINK_DOWN_TIMEOUT_MS, "V2 link down") == 0, "V2 link did not go down"); + CHECK(cc->links_up == 1, "client links_up != 1 after V2 kill (got %d)", cc->links_up); + CHECK(sc->links_up == 1, "server links_up != 1 after V2 kill (got %d)", sc->links_up); + CHECK(cc->initialized == 1, "client initialized lost after V2 kill"); + CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit after V2 kill"); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "failover traffic after V2 kill failed"); + + // -- kill V1 link -- + printf("[A] killing V1 link (server socket %d)...\n", sp1); + kill_server_socket(s1, sp1); + g_link = find_link_by_remote_port(cc, sp1); + CHECK(poll_until(c_link_down, LINK_DOWN_TIMEOUT_MS, "V1 link down") == 0, "V1 link did not go down"); + CHECK(cc->links_up == 1, "client links_up != 1 after V1 kill (got %d)", cc->links_up); + CHECK(cc->initialized == 1, "client initialized lost after V1 kill"); + CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, "reinit after V1 kill"); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "failover traffic after V1 kill failed"); + + // -- kill initial (последний) линк: всё вниз, реинита нет -- + printf("[A] killing initial link (server socket %d)...\n", server_port); + kill_server_socket(s0, server_port); + g_link = find_link_by_remote_port(cc, server_port); + CHECK(poll_until(c_link_down, LINK_DOWN_TIMEOUT_MS, "initial link down") == 0, "initial link did not go down"); + CHECK(poll_until(c_client_all_down, LINK_DOWN_TIMEOUT_MS, "client all links down") == 0, "client links not all down"); + CHECK(poll_until(c_server_all_down, LINK_DOWN_TIMEOUT_MS, "server all links down") == 0, "server links not all down"); + CHECK(cc->reinit_count == cli_reinit0, "reinit when last link died (client %u->%u)", cli_reinit0, cc->reinit_count); + CHECK(sc->reinit_count == srv_reinit0, "reinit when last link died (server %u->%u)", srv_reinit0, sc->reinit_count); + + printf("[A] Session A PASSED\n"); + stop_instances(); + return 0; + +fail: + stop_instances(); + return -1; } -static void cleanup_all_extras(void) { - for (int i = extra_count - 1; i >= 0; i--) { - struct extra_info* info = &extras[i]; - if (info->client_link) { etcp_link_close(info->client_link); info->client_link = NULL; } - if (info->udp_sock) { remove_udp_sock_from_instance(info->udp_sock); info->udp_sock = NULL; release_udp_port(info->port); } - if (info->tcp_sock) { tcp_socket_remove(info->tcp_sock); info->tcp_sock = NULL; release_tcp_port(info->port); } +// ===== Session: коллизии ===== + +static int session_collision(int which) { + const char* names[] = { + "V3a: both correct addr (true collision)", + "V3b: client wrong addr (server-only establish)", + "V3c: server wrong addr (client-only establish)", + }; + printf("\n=== Session %s ===\n", names[which]); + + if (start_instances() != 0) return -1; + + struct ETCP_CONN* sc = get_conn(server_instance); + struct ETCP_CONN* cc = get_conn(client_instance); + CHECK(sc && cc, "no conn"); + + uint32_t srv_reinit0 = sc->reinit_count; + uint32_t cli_reinit0 = cc->reinit_count; + g_reinit_cli0 = cli_reinit0; + + int sp3 = alloc_port(); + int cp1 = alloc_port(); + int dead = alloc_port(); + + struct ETCP_SOCKET* s3 = add_udp_socket(server_instance, sp3, "dyn_s3"); + CHECK(s3, "failed to add server socket S3"); + struct ETCP_SOCKET* c1 = add_udp_socket(client_instance, cp1, "dyn_c1"); + CHECK(c1, "failed to add client socket C1"); + + struct sockaddr_storage s3_addr; mk_addr(&s3_addr, sp3); + struct sockaddr_storage c1_addr; mk_addr(&c1_addr, cp1); + struct sockaddr_storage dead_addr; mk_addr(&dead_addr, dead); + + struct ETCP_LINK* lc3 = NULL; + struct ETCP_LINK* ls3 = NULL; + if (which == COL_A) { + lc3 = etcp_link_new(cc, c1, &s3_addr, 0); // клиент -> S3 (верно) + ls3 = etcp_link_new(sc, s3, &c1_addr, 0); // сервер -> C1 (верно) + } else if (which == COL_B) { + lc3 = etcp_link_new(cc, c1, &dead_addr, 0); // клиент -> мёртвый порт + ls3 = etcp_link_new(sc, s3, &c1_addr, 0); // сервер -> C1 (верно) + } else { + lc3 = etcp_link_new(cc, c1, &s3_addr, 0); // клиент -> S3 (верно) + ls3 = etcp_link_new(sc, s3, &dead_addr, 0); // сервер -> мёртвый порт } - extra_count = 0; + CHECK(lc3 && ls3, "link creation failed (lc3=%p ls3=%p)", (void*)lc3, (void*)ls3); + + if (which == COL_A) { + CHECK(poll_until(c_client_reinit_happened, COLLISION_TIMEOUT_MS, "collision detected (client yield)") == 0, + "collision NOT detected: client reinit did not happen"); + printf("[V3a] collision detected: client reinit %u -> %u (server %u -> %u)\n", + cli_reinit0, cc->reinit_count, srv_reinit0, sc->reinit_count); + CHECK(poll_until(c_conn_up_both, COLLISION_TIMEOUT_MS, "collision recovery") == 0, + "conn did not recover after collision"); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "traffic after collision failed"); + } else { + struct ETCP_LINK* good = (which == COL_B) ? ls3 : lc3; + g_link = good; + CHECK(poll_until(c_link_up, LINK_UP_TIMEOUT_MS, "good link up") == 0, "good link did not come up"); + uint16_t peer_port = (which == COL_B) ? sp3 : cp1; + g_conn = (which == COL_B) ? cc : sc; + g_port = peer_port; + CHECK(poll_until(c_link_on_port_up, LINK_UP_TIMEOUT_MS, "peer incoming link up") == 0, + "peer incoming link did not come up"); + CHECK(sc->reinit_count == srv_reinit0 && cc->reinit_count == cli_reinit0, + "unexpected reinit (srv %u->%u cli %u->%u)", srv_reinit0, sc->reinit_count, cli_reinit0, cc->reinit_count); + CHECK(run_traffic(TRAFFIC_PACKETS, TRAFFIC_TIMEOUT_MS) == 0, "traffic failed"); + } + + printf("[%s] PASSED\n", names[which]); + stop_instances(); + return 0; + +fail: + stop_instances(); + return -1; } int main(void) { @@ -351,148 +471,30 @@ int main(void) { debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); - debug_set_categories(DEBUG_CATEGORY_NONE); + if (getenv("UTUN_TEST_DEBUG")) { + debug_set_level(DEBUG_LEVEL_DEBUG); + debug_set_category_level(DEBUG_CATEGORY_ETCP, DEBUG_LEVEL_DEBUG); + debug_set_category_level(DEBUG_CATEGORY_CONNECTION, DEBUG_LEVEL_DEBUG); + debug_set_category_level(DEBUG_CATEGORY_KEEPALIVE, DEBUG_LEVEL_DEBUG); + debug_set_category_level(DEBUG_CATEGORY_DEBUG, DEBUG_LEVEL_DEBUG); + } - printf("=== ETCP Link Stress Test ===\n"); - printf("Server port: %d, Client port: %d\n", server_port, client_port); - printf("UDP pool: %d-%d, TCP pool: %d-%d\n", - udp_pool_base, udp_pool_base + UDP_POOL_SIZE - 1, - tcp_pool_base, tcp_pool_base + TCP_POOL_SIZE - 1); + printf("=== ETCP Multilink Test ===\n"); + printf("server_port=%d client_port=%d dyn_ports=%d..\n", server_port, client_port, next_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 instance\n"); goto fail; - } - 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 instance\n"); goto fail; - } - printf("Client ready (node_id=%llx)\n", (unsigned long long)client_instance->node_id); - - memset(used_udp_ports, 0, sizeof(used_udp_ports)); - memset(used_tcp_ports, 0, sizeof(used_tcp_ports)); - memset(extras, 0, sizeof(extras)); - - printf("Waiting for connection...\n"); fflush(stdout); - { uint64_t conn_start = get_time_tb(); - while (!(is_connected(server_instance) && is_connected(client_instance))) { - if (get_time_tb() - conn_start > CONN_TIMEOUT_TB) { printf("Connection timeout\n"); goto fail; } - uasync_poll(ua, 10); - } } - printf("=== Connection established ===\n"); fflush(stdout); - - printf("=== Starting stress test (%dms, %d cycles) ===\n", - STRESS_DURATION_MS, STRESS_DURATION_MS / CYCLE_INTERVAL_MS); fflush(stdout); - stress_start_tb = get_time_tb(); + int rc = 0; + rc |= session_a() != 0; + rc |= session_collision(COL_A) != 0; + rc |= session_collision(COL_B) != 0; + rc |= session_collision(COL_C) != 0; - uint64_t last_pkt_tb = stress_start_tb; - uint64_t last_cycle_tb = stress_start_tb; - - while (1) { - uasync_poll(ua, 5); - uint64_t now_tb = get_time_tb(); - uint64_t elapsed_tb = now_tb - stress_start_tb; - if (elapsed_tb >= STRESS_DURATION_MS * 10) break; - - if (now_tb - last_pkt_tb >= PACKET_INTERVAL_TB) { - last_pkt_tb = now_tb; - drain_received(1); - struct ETCP_CONN* cconn = get_conn(client_instance); - if (cconn && cconn->input_queue) { - uint8_t buf[PACKET_SIZE]; - buf[0] = (uint8_t)(packets_sent & 0xFF); - buf[1] = (uint8_t)((packets_sent >> 8) & 0xFF); - buf[2] = (uint8_t)((packets_sent >> 16) & 0xFF); - buf[3] = (uint8_t)((packets_sent >> 24) & 0xFF); - for (int i = 4; i < PACKET_SIZE; i++) buf[i] = (uint8_t)((packets_sent + i) % 256); - if (etcp_int_send(cconn, buf, PACKET_SIZE) == 0) packets_sent++; - } - } - - if (now_tb - last_cycle_tb >= CYCLE_INTERVAL_MS * 10) { - last_cycle_tb = now_tb; - cycle_count++; - - struct ETCP_CONN* sconn = get_conn(server_instance); - uint64_t cur_bytes = 0; - if (sconn) { struct ETCP_LINK* l = sconn->links; while (l) { cur_bytes += l->total_decrypted; l = l->next; } } - int hung = (cycle_count > 1 && cur_bytes == last_received_bytes && packets_sent > 0); - last_received_bytes = cur_bytes; - if (hung) { hang_detected++; printf("[CYCLE %d] *** HANG: no new data (sent=%u recv=%u) ***\n", cycle_count, packets_sent, packets_received); fflush(stdout); } - - printf(" Links report (cycle %d):\n", cycle_count); - if (sconn) { - struct ETCP_LINK* l = sconn->links; int lidx = 0; - while (l) { - uint64_t link_decrypted_delta = l->total_decrypted; - uint64_t link_acked_delta = l->acked_bytes; - int found = 0; - for (int i = 0; i < extra_count; i++) { - if (extras[i].client_link && extras[i].client_link == l) { found = 1; break; } - } - if (found) { - for (int i = 0; i < extra_count; i++) { - if (extras[i].client_link && extras[i].client_link == l) { - link_decrypted_delta -= extras[i].prev_decrypted; - link_acked_delta -= extras[i].acked_bytes; - extras[i].prev_decrypted = l->total_decrypted; - break; - } - } - } - char flags[16] = ""; - snprintf(flags, sizeof(flags), "%s%s%s", l->is_server ? "S" : "C", l->is_tcp ? "T" : "U", l->initialized ? "I" : "."); - printf(" [%d] type=%s id=%d state=%d decr=%s del=%s encr=%s ack=%s retr=%s rtt=%s\n", - lidx++, flags, l->local_link_id, l->link_state, - fmt_bytes(l->total_decrypted), fmt_bytes_delta(link_decrypted_delta), - fmt_bytes(l->total_encrypted), fmt_bytes(l->acked_bytes), - fmt_uint(l->total_retransmissions), fmt_uint(l->rtt_last)); - fflush(stdout); - l = l->next; - } - } - - for (int i = 0; i < extra_count; i++) - if (extras[i].client_link && extras[i].client_link->initialized) - collect_link_stats(extras[i].client_link, &extras[i]); - - struct ETCP_CONN* cconn = get_conn(client_instance); - if (cconn) { - int op = cycle_count % 4; - if (op < 2) add_udp_pair(cconn); - else add_tcp_pair(cconn); - } - - printf("[CYCLE %d] extra=%d sent=%u recv=%u hung=%d elapsed=%llums\n", - cycle_count, extra_count, packets_sent, packets_received, hung, - (unsigned long long)(elapsed_tb / 10)); fflush(stdout); - } - } - - drain_received(1); - printf("\n=== Stress test complete ===\n"); - printf("Packets: sent=%u received=%u\n", packets_sent, packets_received); - printf("Cycles: %d extras remaining: %d hangs: %d\n", cycle_count, extra_count, hang_detected); - - cleanup_all_extras(); - 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(); - printf("\n=== TEST PASSED ===\n"); - return 0; - -fail: - if (server_instance) { server_instance->running = 0; utun_instance_destroy(server_instance); } - if (client_instance) { client_instance->running = 0; utun_instance_destroy(client_instance); } - if (ua) { uasync_destroy(ua, 0); ua = NULL; } - cleanup_temp_configs(); - printf("\n=== TEST FAILED ===\n"); + if (rc == 0) { printf("\n=== ALL TESTS PASSED ===\n"); return 0; } + printf("\n=== TESTS FAILED ===\n"); return 1; }