#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 "../src/config_parser.h" #include "../src/utun_instance.h" #include "routing.h" #include "../src/tun_if.h" #include "secure_channel.h" #include "../lib/u_async.h" #include "../lib/ll_queue.h" #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" 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 udp_pool_base = 0; static int tcp_pool_base = 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 TCP_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; }; 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; 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=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68\n" "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" "tun_ip=10.99.0.1/24\n" "tun_ifname=tun99\n" "\n" "[server: %s]\n" "addr=127.0.0.1:%d\n" "type=public\n" "\n" "[allowed_keys]\n" "allow_all=1\n", INITIAL_SOCKET_NAME, 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=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" "tun_ip=10.99.0.2/24\n" "tun_ifname=tun98\n" "\n" "[server: %s]\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); 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 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) { 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; } 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 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 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 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 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 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 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 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 struct ETCP_SOCKET* get_client_socket(void) { if (client_instance && client_instance->etcp_sockets) return client_instance->etcp_sockets; return NULL; } 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; } 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); 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 TCP_SOCKET* ts = tcp_socket_add(server_instance, &cfg); if (!ts) { release_tcp_port(port); return -1; } struct stcp_link_config scfg = {.ua = ua, .my_keys = &server_instance->my_keys, .inst = server_instance, .listen_family = AF_INET}; struct stcp_server* srv = stcp_server_listen(&scfg, port, tcp_server_on_link, server_instance); if (!srv) { tcp_socket_remove(ts); release_tcp_port(port); return -1; } stcp_server_list_add(server_instance, srv); 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; struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); etcp_tcp_link_start_connect(tlink, &sa, port); 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); 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 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 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; } 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); } } extra_count = 0; } 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 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); 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(); 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"); return 1; }