Browse Source

test_etcp_link_stress: переписан — 5 вариантов добавления линков (V1/V2, коллизии V3a/b/c), kill/failover и инвариант «нет реинита при живом линке»; фикс stack buffer overflow в etcp_encrypt_send (enc_buf 1600→2048)

v2
evgeny 3 weeks ago
parent
commit
13e2697f80
  1. 2
      src/transport_layer/etcp_connections.c
  2. 710
      tests/test_etcp_link_stress.c

2
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 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; 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; } 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; size_t enc_buf_len=0;
dgram->timestamp=get_current_timestamp(); dgram->timestamp=get_current_timestamp();

710
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 <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include "../lib/platform_compat.h" #include "../lib/platform_compat.h"
#include "test_utils.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 "etcp.h" #include "etcp.h"
#include "etcp_connections.h" #include "etcp_connections.h"
#include "stcp_link.h" #include "etcp_dump.h"
#include "../src/config_parser.h" #include "../src/config_parser.h"
#include "../src/utun_instance.h" #include "../src/utun_instance.h"
#include "routing.h" #include "routing.h"
@ -27,15 +19,14 @@
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
#include "../lib/mem.h" #include "../lib/mem.h"
#define CONN_TIMEOUT_TB 50000 #define PACKET_SIZE 64
#define STRESS_DURATION_MS 2000 #define TRAFFIC_PACKETS 200
#define CYCLE_INTERVAL_MS 100 #define LINK_UP_TIMEOUT_MS 8000
#define PACKET_INTERVAL_TB 10 #define TRAFFIC_TIMEOUT_MS 8000
#define PACKET_SIZE 64 #define LINK_DOWN_TIMEOUT_MS 8000
#define MAX_EXTRA 8 #define COLLISION_TIMEOUT_MS 15000
#define UDP_POOL_SIZE 12
#define TCP_POOL_SIZE 4 enum { COL_A = 0, COL_B, COL_C };
#define INITIAL_SOCKET_NAME "udp0"
static struct UTUN_INSTANCE* server_instance = NULL; static struct UTUN_INSTANCE* server_instance = NULL;
static struct UTUN_INSTANCE* client_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 char client_config_path[256];
static int server_port = 0; static int server_port = 0;
static int client_port = 0; static int client_port = 0;
static int udp_pool_base = 0; static int next_port = 0;
static int tcp_pool_base = 0;
static uint32_t packets_sent = 0; static uint32_t packets_sent = 0;
static uint32_t packets_received = 0; static uint32_t packets_received = 0;
static uint64_t last_received_bytes = 0;
static int cycle_count = 0; // globals for generic poll_until conditions
static int hang_detected = 0; static struct ETCP_LINK* g_link = NULL;
static uint64_t stress_start_tb = 0; static struct ETCP_SOCKET* g_sock = NULL;
static struct ETCP_CONN* g_conn = NULL;
struct extra_info { static uint16_t g_port = 0;
int type; static uint32_t g_reinit_cli0 = 0;
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;
};
static int create_temp_configs(void) { static int create_temp_configs(void) {
if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "Failed to create temp directory\n"); return -1; } if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "Failed to create temp directory\n"); return -1; }
int base_port = 50000 + (getpid() % 15000); int base_port = 50000 + (getpid() % 15000);
server_port = base_port; server_port = base_port;
client_port = base_port + 1; client_port = base_port + 1;
udp_pool_base = base_port + 10; next_port = base_port + 100;
tcp_pool_base = base_port + 30;
snprintf(server_config_path, sizeof(server_config_path), "%s/server.conf", temp_dir); 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); 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" "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n"
"tun_ip=10.99.0.1/24\n" "tun_ip=10.99.0.1/24\n"
"tun_ifname=tun99\n" "tun_ifname=tun99\n"
"keepalive_interval=100\n"
"keepalive_adaptive=0\n"
"\n" "\n"
"[server: %s]\n" "[server: udp0]\n"
"addr=127.0.0.1:%d\n" "addr=127.0.0.1:%d\n"
"type=public\n" "type=public\n"
"\n" "\n"
"[allowed_keys]\n" "[allowed_keys]\n"
"allow_all=1\n", "allow_all=1\n",
INITIAL_SOCKET_NAME, server_port); server_port);
fclose(f); fclose(f);
f = fopen(client_config_path, "w"); f = fopen(client_config_path, "w");
@ -125,16 +89,21 @@ static int create_temp_configs(void) {
"my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n"
"tun_ip=10.99.0.2/24\n" "tun_ip=10.99.0.2/24\n"
"tun_ifname=tun98\n" "tun_ifname=tun98\n"
"keepalive_interval=100\n"
"keepalive_adaptive=0\n"
"\n" "\n"
"[server: %s]\n" "[server: udp0]\n"
"addr=127.0.0.1:%d\n" "addr=127.0.0.1:%d\n"
"type=public\n" "type=public\n"
"\n" "\n"
"[client: test_client]\n" "[client: test_client]\n"
"keepalive=1\n" "keepalive=1\n"
"peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" "peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n"
"link=%s:127.0.0.1:%d\n", "link=udp0:127.0.0.1:%d\n"
INITIAL_SOCKET_NAME, client_port, INITIAL_SOCKET_NAME, server_port); "\n"
"[allowed_keys]\n"
"allow_all=1\n",
client_port, server_port);
fclose(f); fclose(f);
return 0; return 0;
} }
@ -145,205 +114,356 @@ static void cleanup_temp_configs(void) {
if (temp_dir[0]) test_rmdir(temp_dir); if (temp_dir[0]) test_rmdir(temp_dir);
} }
// ===== helpers =====
static struct ETCP_CONN* get_conn(struct UTUN_INSTANCE* inst) { static struct ETCP_CONN* get_conn(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->connections || !inst->connections->head) return NULL; if (!inst || !inst->connections || !inst->connections->head) return NULL;
struct conn_queue_entry* ce = (struct conn_queue_entry*)inst->connections->head->data; struct conn_queue_entry* ce = (struct conn_queue_entry*)inst->connections->head->data;
return ce ? ce->conn : NULL; 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); struct ETCP_CONN* conn = get_conn(inst);
if (!conn) return 0; return conn && conn->initialized && conn->links_up;
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) { static int count_links(struct ETCP_CONN* conn) {
if (!server_instance) return; int n = 0;
struct ETCP_CONN* conn = get_conn(server_instance); for (struct ETCP_LINK* l = conn ? conn->links : NULL; l; l = l->next) n++;
if (!conn || !conn->output_queue) return; return n;
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) { static int count_links_up(struct ETCP_CONN* conn) {
for (int i = 0; i < UDP_POOL_SIZE; i++) { int n = 0;
if (!used_udp_ports[i]) { used_udp_ports[i] = 1; return udp_pool_base + i; } for (struct ETCP_LINK* l = conn ? conn->links : NULL; l; l = l->next) if (l->link_status) n++;
} return n;
return -1;
} }
static void release_udp_port(int port) { static void mk_addr(struct sockaddr_storage* out, uint16_t port) {
for (int i = 0; i < UDP_POOL_SIZE; i++) { if (udp_pool_base + i == port) { used_udp_ports[i] = 0; return; } } 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) { static int alloc_port(void) { return next_port++; }
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) { static struct ETCP_SOCKET* add_udp_socket(struct UTUN_INSTANCE* inst, uint16_t port, const char* name) {
for (int i = 0; i < TCP_POOL_SIZE; i++) { if (tcp_pool_base + i == port) { used_tcp_ports[i] = 0; return; } } 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) { static struct ETCP_LINK* find_link_by_remote_port(struct ETCP_CONN* conn, uint16_t port) {
if (!link) return; for (struct ETCP_LINK* l = conn ? conn->links : NULL; l; l = l->next) {
info->bytes_decrypted = link->total_decrypted; if (l->remote_addr.ss_family == AF_INET &&
info->bytes_encrypted = link->total_encrypted; ntohs(((struct sockaddr_in*)&l->remote_addr)->sin_port) == port)
info->acked_bytes = link->acked_bytes; return l;
info->retrans = link->total_retransmissions; }
info->rtt = link->rtt_last; return NULL;
info->inflight = link->inflight_bytes;
} }
static void remove_udp_sock_from_instance(struct ETCP_SOCKET* sock) { static struct ETCP_LINK* first_link_on_socket(struct ETCP_SOCKET* sock) {
if (!server_instance || !sock) return; if (!sock || !sock->links_queue || !sock->links_queue->head) return NULL;
struct ETCP_SOCKET** pp = &server_instance->etcp_sockets; struct link_queue_entry* lqe = (struct link_queue_entry*)sock->links_queue->head->data;
while (*pp) { if (*pp == sock) { *pp = sock->next; etcp_socket_remove(sock); return; } pp = &(*pp)->next; } return lqe ? lqe->link : NULL;
} }
static void collect_and_remove_oldest(void) { static uint64_t client_link_acked(uint16_t srv_port) {
if (extra_count == 0) return; struct ETCP_LINK* l = find_link_by_remote_port(get_conn(client_instance), srv_port);
int oldest_idx = -1; return l ? l->acked_packets : 0;
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) { static void dump_state(void) {
if (client_instance && client_instance->etcp_sockets) return client_instance->etcp_sockets; printf("--- STATE DUMP ---\n");
return NULL; 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) { #define CHECK(cond, ...) do { if (!(cond)) { \
if (extra_count >= MAX_EXTRA) { collect_and_remove_oldest(); if (extra_count >= MAX_EXTRA) return -1; } printf("FAIL %s:%d: ", __func__, __LINE__); \
int port = get_free_udp_port(); printf(__VA_ARGS__); \
if (port < 0) { fprintf(stderr, "No free UDP ports\n"); return -1; } printf("\n"); \
dump_state(); \
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); goto fail; } } while (0)
sin.sin_family = AF_INET; sin.sin_addr.s_addr = htonl(INADDR_LOOPBACK); sin.sin_port = htons(port);
// ===== poll conditions =====
struct CFG_SERVER cfg; memset(&cfg, 0, sizeof(cfg));
snprintf(cfg.name, sizeof(cfg.name), "dyn_udp_%d", port); static int c_conn_up_both(void) { return conn_up(server_instance) && conn_up(client_instance); }
cfg.type = CFG_SERVER_TYPE_PUBLIC; cfg.mtu = PACKET_DATA_MAX_MTU; static int c_link_up(void) { return g_link && g_link->initialized && g_link->link_status == 1; }
memcpy(&cfg.ip, &sin, sizeof(sin)); 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; }
struct ETCP_SOCKET* es = etcp_socket_add(server_instance, &cfg); 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; }
if (!es) { release_udp_port(port); fprintf(stderr, "Failed to create UDP socket on port %d\n", port); return -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; }
struct ETCP_SOCKET* client_sock = get_client_socket(); static int c_client_reinit_happened(void) { struct ETCP_CONN* c = get_conn(client_instance); return c && c->reinit_count > g_reinit_cli0; }
if (!client_sock) { remove_udp_sock_from_instance(es); release_udp_port(port); fprintf(stderr, "No client socket\n"); return -1; }
static int poll_until(int (*cond)(void), uint64_t timeout_ms, const char* what) {
struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); uint64_t deadline = get_time_tb() + timeout_ms * 10;
struct ETCP_LINK* link = etcp_link_new(conn, client_sock, &sa, 0); while (get_time_tb() < deadline) {
if (!link) { remove_udp_sock_from_instance(es); release_udp_port(port); return -1; } if (cond()) return 0;
uasync_poll(ua, 5);
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; printf("TIMEOUT waiting for: %s (%llums)\n", what, (unsigned long long)timeout_ms);
info->create_tb = get_time_tb(); return -1;
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) { // ===== traffic =====
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)); static void drain_received(int count_flag) {
snprintf(cfg.name, sizeof(cfg.name), "dyn_tcp_%d", port); if (!server_instance) return;
cfg.type = CFG_SERVER_TYPE_PUBLIC; cfg.transport = 1; struct ETCP_CONN* conn = get_conn(server_instance);
memcpy(&cfg.ip, &sin, sizeof(sin)); if (!conn || !conn->output_queue) return;
queue_set_callback(conn->output_queue, NULL, NULL);
struct ETCP_SOCKET* ts = tcp_socket_add(server_instance, &cfg); struct ETCP_FRAGMENT* pkt;
if (!ts) { release_tcp_port(port); return -1; } while ((pkt = (struct ETCP_FRAGMENT*)queue_data_get(conn->output_queue)) != NULL) {
if (count_flag) packets_received++;
struct stcp_link_config scfg = {.inst = server_instance, .listen_family = AF_INET}; if (pkt->ll.dgram) memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
struct stcp_server* srv = stcp_server_listen(&scfg, port, (stcp_server_on_link_cb)tcp_server_on_link, ts); queue_entry_free((struct ll_entry*)pkt);
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); static int run_traffic(int n, uint64_t timeout_ms) {
if (!tlink) { tcp_socket_remove(ts); release_tcp_port(port); return -1; } struct ETCP_CONN* cc = get_conn(client_instance);
tlink->is_tcp = 1; tlink->is_server = 0; 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)); // ===== instance lifecycle =====
etcp_tcp_link_start_connect(tlink, &sa, port);
struct extra_info* info = &extras[extra_count++]; memset(info, 0, sizeof(*info)); static int start_instances(void) {
info->type = 1; info->port = port; info->tcp_sock = ts; info->tcp_srv = srv; info->client_link = tlink; server_instance = utun_instance_create(ua, server_config_path);
info->create_tb = get_time_tb(); if (!server_instance || utun_instance_init(server_instance) < 0) {
printf("[CYCLE %d] ADD TCP: port=%d link=%p sock=%p srv=%p\n", cycle_count, port, (void*)tlink, (void*)ts, (void*)srv); 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; return 0;
} }
static const char* fmt_bytes(uint64_t b) { static void stop_instances(void) {
static char buf[6][32]; static int n = 0; if (server_instance) { server_instance->running = 0; utun_instance_destroy(server_instance); server_instance = NULL; }
char* s = buf[n]; n = (n + 1) % 6; if (client_instance) { client_instance->running = 0; utun_instance_destroy(client_instance); client_instance = NULL; }
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; static void kill_server_socket(struct ETCP_SOCKET* sock, uint16_t srv_port) {
char* s = buf[n]; n = (n + 1) % 6; struct ETCP_LINK* cl = find_link_by_remote_port(get_conn(client_instance), srv_port);
if (b >= 1024) snprintf(s, 32, "+%lluk", (unsigned long long)(b / 1024)); else snprintf(s, 32, "+%llu", (unsigned long long)b); if (cl) cl->last_recv_local_time = 0;
return s; etcp_socket_remove(sock);
} }
static const char* fmt_uint(uint32_t v) { // ===== Session A: one-way adds + kill failover =====
static char buf[6][16]; static int n = 0;
char* s = buf[n]; n = (n + 1) % 6; static int session_a(void) {
snprintf(s, 16, "%u", v); printf("\n=== Session A: one-way adds (V1/V2) + link kill failover ===\n");
return s;
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) { // ===== Session: коллизии =====
for (int i = extra_count - 1; i >= 0; i--) {
struct extra_info* info = &extras[i]; static int session_collision(int which) {
if (info->client_link) { etcp_link_close(info->client_link); info->client_link = NULL; } const char* names[] = {
if (info->udp_sock) { remove_udp_sock_from_instance(info->udp_sock); info->udp_sock = NULL; release_udp_port(info->port); } "V3a: both correct addr (true collision)",
if (info->tcp_sock) { tcp_socket_remove(info->tcp_sock); info->tcp_sock = NULL; release_tcp_port(info->port); } "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) { int main(void) {
@ -351,148 +471,30 @@ int main(void) {
debug_config_init(); debug_config_init();
debug_set_level(DEBUG_LEVEL_ERROR); 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("=== ETCP Multilink Test ===\n");
printf("Server port: %d, Client port: %d\n", server_port, client_port); printf("server_port=%d client_port=%d dyn_ports=%d..\n", server_port, client_port, next_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); utun_instance_set_tun_init_enabled(0);
ua = uasync_create(); ua = uasync_create();
server_instance = utun_instance_create(ua, server_config_path); int rc = 0;
if (!server_instance || utun_instance_init(server_instance) < 0) { rc |= session_a() != 0;
fprintf(stderr, "Failed to create server instance\n"); goto fail; rc |= session_collision(COL_A) != 0;
} rc |= session_collision(COL_B) != 0;
printf("Server ready (node_id=%llx)\n", (unsigned long long)server_instance->node_id); rc |= session_collision(COL_C) != 0;
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; } if (ua) { uasync_destroy(ua, 0); ua = NULL; }
cleanup_temp_configs(); cleanup_temp_configs();
printf("\n=== TEST PASSED ===\n"); if (rc == 0) { printf("\n=== ALL TESTS PASSED ===\n"); return 0; }
return 0; printf("\n=== TESTS FAILED ===\n");
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; return 1;
} }

Loading…
Cancel
Save