From 59f48029045ece0b7272c33a1e15bd1297f67227 Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 27 Sep 2026 11:57:29 +0300 Subject: [PATCH] Fix autosocket ownership and exercise interface churn under load --- src/transport_layer/auto_socket.c | 104 +++--- src/transport_layer/etcp_connections.c | 22 +- tests/test_auto_socket_dynamic.c | 464 +++++++++++++++++-------- 3 files changed, 375 insertions(+), 215 deletions(-) diff --git a/src/transport_layer/auto_socket.c b/src/transport_layer/auto_socket.c index 47f0e695..351cf721 100644 --- a/src/transport_layer/auto_socket.c +++ b/src/transport_layer/auto_socket.c @@ -112,10 +112,13 @@ static int v6_classify(const uint8_t addr[16]) { struct auto_sock_iface { struct auto_sock_iface* next; uint32_t netif_index; // индекс интерфейса (if_nametoindex) + char ifname[IF_NAMESIZE]; // нужен и после удаления интерфейса из ОС struct ETCP_SOCKET* v4_udp; // UDP-сокет для IPv4 (NULL если не создан) struct ETCP_SOCKET* v6_udp; // UDP-сокет для IPv6 struct stcp_server* v4_tcp; // TCP STCP-сервер для IPv4 struct stcp_server* v6_tcp; // TCP STCP-сервер для IPv6 + struct ETCP_SOCKET* v4_tcp_sock; // модель сокета принадлежит тому же интерфейсу + struct ETCP_SOCKET* v6_tcp_sock; uint8_t v4_type; // CFG_SERVER_TYPE_PUBLIC или CFG_SERVER_TYPE_NAT uint8_t v6_type; // CFG_SERVER_TYPE_PUBLIC или CFG_SERVER_TYPE_NAT }; @@ -146,7 +149,8 @@ struct AUTO_SOCKET { /* ─── forward declarations ─── */ static int create_iface_udp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname, int family, uint8_t type); -static struct stcp_server* create_iface_tcp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname, int family, uint8_t type); +static struct stcp_server* create_iface_tcp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname, + int family, uint8_t type, struct ETCP_SOCKET** out_sock); static void remove_iface_sockets(struct AUTO_SOCKET* as, uint32_t ifindex); static int scan_and_classify_iface(uint32_t ifindex, const char* ifname, int* out_has_v4, int* out_has_v6, uint8_t* out_v4_type, uint8_t* out_v6_type); static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname); @@ -351,8 +355,10 @@ static int create_iface_udp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, con * (etcp_link_new + is_tcp=1 + etcp_tcp_link_start_connect) к каждому соединению * * Возвращает stcp_server* при успехе, NULL при ошибке. */ -static struct stcp_server* create_iface_tcp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname, int family, uint8_t type) { +static struct stcp_server* create_iface_tcp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname, + int family, uint8_t type, struct ETCP_SOCKET** out_sock) { struct UTUN_INSTANCE* inst = as->instance; + *out_sock = NULL; uint16_t saved_port = load_port_from_db(as, ifname, family, AS_PROTO_TCP); int max_att = saved_port ? 11 : 10; @@ -403,6 +409,7 @@ static struct stcp_server* create_iface_tcp_socket(struct AUTO_SOCKET* as, uint3 } stcp_server_list_add(inst, tsrv); + *out_sock = ts; save_port_to_db(as, ifname, family, AS_PROTO_TCP, port); @@ -435,24 +442,25 @@ int auto_socket_add_interface(struct UTUN_INSTANCE* inst, uint32_t ifindex, ifa = u_calloc(1, sizeof(*ifa)); if (!ifa) { DEBUG_ERROR(DEBUG_CATEGORY_AS, "[as] alloc failed"); return 0; } ifa->netif_index = ifindex; + snprintf(ifa->ifname, sizeof(ifa->ifname), "%s", ifname); int added = 0; if (has_v4) { if (create_iface_udp_socket(as, ifindex, ifname, AF_INET, v4_type) == 0) { added++; ifa->v4_type = v4_type; struct ETCP_SOCKET* s = inst->etcp_sockets; - while (s) { if (s->netif_index == ifindex && s->local_addr.ss_family == AF_INET && !ifa->v4_udp) { ifa->v4_udp = s; break; } s = s->next; } + while (s) { if (!s->is_tcp && s->netif_index == ifindex && s->local_addr.ss_family == AF_INET && !ifa->v4_udp) { ifa->v4_udp = s; break; } s = s->next; } } - { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET, v4_type); + { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET, v4_type, &ifa->v4_tcp_sock); if (tsrv) { added++; ifa->v4_tcp = tsrv; } } } if (has_v6) { if (create_iface_udp_socket(as, ifindex, ifname, AF_INET6, v6_type) == 0) { added++; ifa->v6_type = v6_type; struct ETCP_SOCKET* s = inst->etcp_sockets; - while (s) { if (s->netif_index == ifindex && s->local_addr.ss_family == AF_INET6 && !ifa->v6_udp) { ifa->v6_udp = s; break; } s = s->next; } + while (s) { if (!s->is_tcp && s->netif_index == ifindex && s->local_addr.ss_family == AF_INET6 && !ifa->v6_udp) { ifa->v6_udp = s; break; } s = s->next; } } - { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET6, v6_type); + { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET6, v6_type, &ifa->v6_tcp_sock); if (tsrv) { added++; ifa->v6_tcp = tsrv; } } } @@ -482,14 +490,23 @@ void auto_socket_remove_interface(struct UTUN_INSTANCE* inst, uint32_t ifindex) * 3. Для TCP — дополнительно останавливает stcp_server * 4. Удаляет запись порта из БД * Затем освобождает запись auto_sock_iface и оповещает GUI. */ +static void remove_iface_tcp(struct AUTO_SOCKET* as, struct stcp_server** server, struct ETCP_SOCKET** socket) { + /* Сначала отменяем все callbacks/линки, затем освобождаем их socket context. */ + if (*server) { + stcp_server_list_remove(as->instance, *server); + stcp_link_server_destroy(*server); + *server = NULL; + } + if (*socket) { tcp_socket_remove(*socket); *socket = NULL; } +} + static void remove_iface_sockets(struct AUTO_SOCKET* as, uint32_t ifindex) { struct auto_sock_iface** pp = &as->ifaces; while (*pp) { if ((*pp)->netif_index == ifindex) { struct auto_sock_iface* ifa = *pp; *pp = ifa->next; - char ifname[IF_NAMESIZE] = ""; - if (!if_indextoname(ifindex, ifname)) snprintf(ifname, sizeof(ifname), "?"); + const char* ifname = ifa->ifname; if (ifa->v4_udp) { DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] removing v4 UDP socket %s ifidx=%u", ifa->v4_udp->name, ifindex); @@ -505,36 +522,12 @@ static void remove_iface_sockets(struct AUTO_SOCKET* as, uint32_t ifindex) { } if (ifa->v4_tcp) { DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] removing v4 TCP socket ifidx=%u", ifindex); - { uint16_t v4_port = load_port_from_db(as, ifname, AF_INET, AS_PROTO_TCP); - struct ETCP_SOCKET** tsp = &as->instance->etcp_sockets; - while (*tsp) { - if (!(*tsp)->is_tcp) { tsp = &(*tsp)->next; continue; } - uint16_t tsp_port = 0; - const struct sockaddr_storage* addr = (*tsp)->interface_addr.ss_family ? &(*tsp)->interface_addr : &(*tsp)->local_addr; - if (addr->ss_family == AF_INET) tsp_port = ntohs(((struct sockaddr_in*)addr)->sin_port); - if (tsp_port == v4_port && v4_port > 0) { struct ETCP_SOCKET* rm = *tsp; *tsp = rm->next; u_free(rm); break; } - tsp = &(*tsp)->next; - } - } - stcp_server_list_remove(as->instance, ifa->v4_tcp); - stcp_link_server_destroy(ifa->v4_tcp); + remove_iface_tcp(as, &ifa->v4_tcp, &ifa->v4_tcp_sock); delete_port_from_db(as, ifname, AF_INET, AS_PROTO_TCP); } if (ifa->v6_tcp) { DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] removing v6 TCP socket ifidx=%u", ifindex); - { uint16_t v6_port = load_port_from_db(as, ifname, AF_INET6, AS_PROTO_TCP); - struct ETCP_SOCKET** tsp = &as->instance->etcp_sockets; - while (*tsp) { - if (!(*tsp)->is_tcp) { tsp = &(*tsp)->next; continue; } - uint16_t tsp_port = 0; - const struct sockaddr_storage* addr = (*tsp)->interface_addr.ss_family ? &(*tsp)->interface_addr : &(*tsp)->local_addr; - if (addr->ss_family == AF_INET6) tsp_port = ntohs(((struct sockaddr_in6*)addr)->sin6_port); - if (tsp_port == v6_port && v6_port > 0) { struct ETCP_SOCKET* rm = *tsp; *tsp = rm->next; u_free(rm); break; } - tsp = &(*tsp)->next; - } - } - stcp_server_list_remove(as->instance, ifa->v6_tcp); - stcp_link_server_destroy(ifa->v6_tcp); + remove_iface_tcp(as, &ifa->v6_tcp, &ifa->v6_tcp_sock); delete_port_from_db(as, ifname, AF_INET6, AS_PROTO_TCP); } @@ -919,12 +912,13 @@ static int iface_has_default_v6(uint32_t ifindex, const char* ifname) { (void)if /* Находит или создаёт запись auto_sock_iface для интерфейса. * Нужна потому что reconcile_iface может вызываться до auto_socket_add_interface * (например при старте через auto_socket_scan_all). */ -static struct auto_sock_iface* reconcile_get_ifa(struct AUTO_SOCKET* as, uint32_t ifindex) { +static struct auto_sock_iface* reconcile_get_ifa(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname) { struct auto_sock_iface* ifa = as->ifaces; while (ifa) { if (ifa->netif_index == ifindex) return ifa; ifa = ifa->next; } ifa = u_calloc(1, sizeof(*ifa)); if (!ifa) return NULL; ifa->netif_index = ifindex; + snprintf(ifa->ifname, sizeof(ifa->ifname), "%s", ifname); ifa->next = as->ifaces; as->ifaces = ifa; return ifa; @@ -984,12 +978,12 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char struct ETCP_SOCKET* sock = ifa ? ifa->v4_udp : NULL; uint8_t cur_type = ifa ? ifa->v4_type : CFG_SERVER_TYPE_NAT; if (has_v4 && !sock) { - if (!ifa) ifa = reconcile_get_ifa(as, ifindex); + if (!ifa) ifa = reconcile_get_ifa(as, ifindex, ifname); if (ifa) { create_iface_udp_socket(as, ifindex, ifname, AF_INET, v4_need_type); ifa->v4_type = v4_need_type; { struct ETCP_SOCKET* s = as->instance->etcp_sockets; - while (s) { if (s->netif_index == ifindex && s->local_addr.ss_family == AF_INET && !ifa->v4_udp) { ifa->v4_udp = s; break; } s = s->next; } } + while (s) { if (!s->is_tcp && s->netif_index == ifindex && s->local_addr.ss_family == AF_INET && !ifa->v4_udp) { ifa->v4_udp = s; break; } s = s->next; } } changed = 1; } } else if (!has_v4 && sock) { @@ -1017,7 +1011,7 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char etcp_socket_remove(sock); ifa->v4_type = new_type; { struct ETCP_SOCKET* s = as->instance->etcp_sockets; - while (s) { if (s->netif_index == ifindex && s->local_addr.ss_family == AF_INET) { ifa->v4_udp = s; break; } s = s->next; } } + while (s) { if (!s->is_tcp && s->netif_index == ifindex && s->local_addr.ss_family == AF_INET) { ifa->v4_udp = s; break; } s = s->next; } } as->v4_addr_changed = 1; changed = 1; } @@ -1029,12 +1023,12 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char struct ETCP_SOCKET* sock = ifa ? ifa->v6_udp : NULL; uint8_t cur_type = ifa ? ifa->v6_type : CFG_SERVER_TYPE_NAT; if (has_v6 && !sock) { - if (!ifa) ifa = reconcile_get_ifa(as, ifindex); + if (!ifa) ifa = reconcile_get_ifa(as, ifindex, ifname); if (ifa) { create_iface_udp_socket(as, ifindex, ifname, AF_INET6, v6_need_type); ifa->v6_type = v6_need_type; { struct ETCP_SOCKET* s = as->instance->etcp_sockets; - while (s) { if (s->netif_index == ifindex && s->local_addr.ss_family == AF_INET6 && !ifa->v6_udp) { ifa->v6_udp = s; break; } s = s->next; } } + while (s) { if (!s->is_tcp && s->netif_index == ifindex && s->local_addr.ss_family == AF_INET6 && !ifa->v6_udp) { ifa->v6_udp = s; break; } s = s->next; } } changed = 1; } } else if (!has_v6 && sock) { @@ -1062,7 +1056,7 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char etcp_socket_remove(sock); ifa->v6_type = new_type; { struct ETCP_SOCKET* s = as->instance->etcp_sockets; - while (s) { if (s->netif_index == ifindex && s->local_addr.ss_family == AF_INET6) { ifa->v6_udp = s; break; } s = s->next; } } + while (s) { if (!s->is_tcp && s->netif_index == ifindex && s->local_addr.ss_family == AF_INET6) { ifa->v6_udp = s; break; } s = s->next; } } as->v6_addr_changed = 1; changed = 1; } @@ -1074,19 +1068,11 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char int has = has_v4 && 1; /* always create both TCP and UDP if v4 present */ struct stcp_server* srv = ifa ? ifa->v4_tcp : NULL; if (has && !srv) { - if (!ifa) ifa = reconcile_get_ifa(as, ifindex); - if (ifa) { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET, v4_need_type); ifa->v4_tcp = tsrv; if (tsrv) changed = 1; } + if (!ifa) ifa = reconcile_get_ifa(as, ifindex, ifname); + if (ifa) { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET, v4_need_type, &ifa->v4_tcp_sock); ifa->v4_tcp = tsrv; if (tsrv) changed = 1; } } else if (!has && srv) { DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v4 TCP removed: %s ifidx=%u", ifname, ifindex); - { char buf[IF_NAMESIZE]; const char* nm = if_indextoname(ifindex, buf) ? buf : ifname; - uint16_t p = load_port_from_db(as, nm, AF_INET, AS_PROTO_TCP); - struct ETCP_SOCKET** tsp = &as->instance->etcp_sockets; - while (*tsp) { if (!(*tsp)->is_tcp) { tsp = &(*tsp)->next; continue; } uint16_t tp = 0; const struct sockaddr_storage* a = (*tsp)->interface_addr.ss_family ? &(*tsp)->interface_addr : &(*tsp)->local_addr; - if (a->ss_family == AF_INET) tp = ntohs(((struct sockaddr_in*)a)->sin_port); - if (tp == p && p > 0) { struct ETCP_SOCKET* rm = *tsp; *tsp = rm->next; u_free(rm); break; } tsp = &(*tsp)->next; } - } - stcp_server_list_remove(as->instance, srv); - stcp_link_server_destroy(srv); + remove_iface_tcp(as, &ifa->v4_tcp, &ifa->v4_tcp_sock); delete_port_from_db(as, ifname, AF_INET, AS_PROTO_TCP); ifa->v4_tcp = NULL; changed = 1; } @@ -1097,19 +1083,11 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char int has = has_v6 && 1; struct stcp_server* srv = ifa ? ifa->v6_tcp : NULL; if (has && !srv) { - if (!ifa) ifa = reconcile_get_ifa(as, ifindex); - if (ifa) { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET6, v6_need_type); ifa->v6_tcp = tsrv; if (tsrv) changed = 1; } + if (!ifa) ifa = reconcile_get_ifa(as, ifindex, ifname); + if (ifa) { struct stcp_server* tsrv = create_iface_tcp_socket(as, ifindex, ifname, AF_INET6, v6_need_type, &ifa->v6_tcp_sock); ifa->v6_tcp = tsrv; if (tsrv) changed = 1; } } else if (!has && srv) { DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v6 TCP removed: %s ifidx=%u", ifname, ifindex); - { char buf[IF_NAMESIZE]; const char* nm = if_indextoname(ifindex, buf) ? buf : ifname; - uint16_t p = load_port_from_db(as, nm, AF_INET6, AS_PROTO_TCP); - struct ETCP_SOCKET** tsp = &as->instance->etcp_sockets; - while (*tsp) { if (!(*tsp)->is_tcp) { tsp = &(*tsp)->next; continue; } uint16_t tp = 0; const struct sockaddr_storage* a = (*tsp)->interface_addr.ss_family ? &(*tsp)->interface_addr : &(*tsp)->local_addr; - if (a->ss_family == AF_INET6) tp = ntohs(((struct sockaddr_in6*)a)->sin6_port); - if (tp == p && p > 0) { struct ETCP_SOCKET* rm = *tsp; *tsp = rm->next; u_free(rm); break; } tsp = &(*tsp)->next; } - } - stcp_server_list_remove(as->instance, srv); - stcp_link_server_destroy(srv); + remove_iface_tcp(as, &ifa->v6_tcp, &ifa->v6_tcp_sock); delete_port_from_db(as, ifname, AF_INET6, AS_PROTO_TCP); ifa->v6_tcp = NULL; changed = 1; } diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 3ce1f8f0..4696f10f 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -903,6 +903,7 @@ struct ETCP_SOCKET* tcp_socket_add(struct UTUN_INSTANCE* instance, struct CFG_SE if (!ts) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_socket_add: alloc failed"); return NULL; } ts->instance = instance; ts->is_tcp = 1; + ts->netif_index = server->netif_index; ts->reality_enabled = server->reality_enabled; ts->fd = SOCKET_INVALID; ts->type = server->type; @@ -979,19 +980,32 @@ struct ETCP_SOCKET* tcp_socket_add(struct UTUN_INSTANCE* instance, struct CFG_SE { char loc_str[64] = "none", if_str[64] = "none"; if (ts->local_addr.ss_family) snprintf(loc_str, sizeof(loc_str), "%s", sockaddr_storage_to_str(&ts->local_addr).str); if (ts->interface_addr.ss_family) snprintf(if_str, sizeof(if_str), "%s", sockaddr_storage_to_str(&ts->interface_addr).str); - DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_socket_add: %s type=%s sock_id=%u local=%s iface=%s", - ts->name, server_type_str(ts->type), ts->sock_id, loc_str, if_str); } + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_socket_add: %s type=%s sock_id=%u ifidx=%u local=%s iface=%s", + ts->name, server_type_str(ts->type), ts->sock_id, ts->netif_index, loc_str, if_str); } return ts; } -// Удаление TCP-сокета из списка instance->etcp_sockets. +// Удаление TCP-сокета и всех его линков, включая созданные вне NCD. void tcp_socket_remove(struct ETCP_SOCKET* sock) { if (!sock || !sock->instance) return; struct UTUN_INSTANCE* inst = sock->instance; struct ETCP_SOCKET** pp = &inst->etcp_sockets; while (*pp && *pp != sock) pp = &(*pp)->next; if (*pp) *pp = sock->next; - DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_socket_remove: %s sock_id=%u", sock->name, sock->sock_id); + unsigned closed = 0; + for (;;) { + struct ETCP_LINK* found = NULL; + for (struct ll_entry* e = inst->connections ? inst->connections->head : NULL; e && !found; e = e->next) { + struct ETCP_CONN* conn = ((struct conn_queue_entry*)e->data)->conn; + for (struct ETCP_LINK* link = conn->links; link; link = link->next) + if (link->conn == sock) { found = link; break; } + } + if (!found) break; + /* DOWN callbacks могут менять список соединений: после close начинаем поиск заново. */ + etcp_link_close(found); + closed++; + } + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_socket_remove: %s sock_id=%u closed_links=%u", sock->name, sock->sock_id, closed); u_free(sock); } diff --git a/tests/test_auto_socket_dynamic.c b/tests/test_auto_socket_dynamic.c index a5b73da9..1c135f21 100644 --- a/tests/test_auto_socket_dynamic.c +++ b/tests/test_auto_socket_dynamic.c @@ -1,17 +1,20 @@ +#define _GNU_SOURCE #include #include #include #include +#include #include "../lib/platform_compat.h" #include "test_utils.h" #ifndef _WIN32 -#include #include +#include #endif #include "etcp.h" #include "etcp_connections.h" #include "etcp_api.h" +#include "pkt_normalizer.h" #include "../src/config_parser.h" #include "../src/config_updater.h" #include "../src/utun_instance.h" @@ -23,46 +26,51 @@ #include "../lib/ll_queue.h" #include "../lib/debug_config.h" #include "../lib/mem.h" +#include "../lib/getmyip.h" +#include "../lib/sqlite3.h" -#define TIMEOUT_TB 300000 +#define TIMEOUT_TB 600000 #define POLL_MS 5 #define STEP_TB 3000 /* 300ms между фазами */ #define TRAF_SEND_TB 50 /* 5ms — отправка */ -#define TRAF_MON_TB 1000 /* 100ms — мониторинг */ +#define TRAF_MON_TB 10000 /* 1s — мониторинг */ #define ETCP_RT_ID_TEST 0xF0 +#define PROBE_SIZE 1024 +#define SEND_BUDGET 64 +#define RX_BUDGET_BYTES (16 * 1024 * 1024) /* бюджет памяти сценария при переупорядочении */ typedef void (*timeout_cb)(void*); +enum timer_id { TIMER_DEADLINE, TIMER_PHASE, TIMER_SEND, TIMER_MONITOR, TIMER_COUNT }; +static struct test_timer { void* handle; timeout_cb callback; } timers[TIMER_COUNT]; +static struct pending_send { + struct ETCP_CONN* conn; + struct queue_waiter_handle waiter; + uint32_t sent, received, snapshot, blocked, peak_bytes; + uint32_t peak_rx_bytes; + uint32_t reinit; + uint32_t outage_sent; + int budget; +} pending[2]; static struct test_ctx { struct UTUN_INSTANCE *server, *client; struct UASYNC* ua; + struct NODE_CONN_DIRECT* handle; + struct ETCP_CONN* tcp_probe; int round; /* 0..7 */ int step; int ip_changes_on_last; char cur_iface[IFNAMSIZ]; char prev_iface[IFNAMSIZ]; - int connected; /* etcp_connect callback fired */ + int connected; uint64_t srv_node_id; uint8_t srv_pubkey[SC_PUBKEY_SIZE]; int result; /* 0=running, 1=fail, 2=pass */ - uint32_t recv_at_ip_change; /* total_recv на момент смены IP */ - - /* traffic */ - uint32_t send_seq, send_count, pong_count, total_recv; - uint32_t expected_rx_seq; /* следующий ожидаемый seq в ответе */ - uint32_t rx_seq_gaps; /* счётчик пропусков в seq */ - uint32_t rx_seq_dups; /* счётчик дубликатов */ - - /* connection health */ - uint32_t conn_reinit_snapshot; /* reinit_count на начало текущей фазы */ - uint32_t max_allowed_reinits; /* максимально допустимых reinits за тест */ - - /* full-disconnect phase */ - int all_deleted; + uint32_t total_recv; uint32_t pre_delete_total_recv; + int draining; } ctx; -static void* timeout_handle; static char tdir[] = "/tmp/utun_as_XXXXXX"; static char scf[256], ccf[256]; @@ -95,7 +103,20 @@ static void fail(const char* msg) { fprintf(stderr, "FAIL r=%d s=%d: %s\n", ctx.round, ctx.step, msg); fflush(stderr); ctx.result = 1; } -static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); ctx.result = 1; } +static void timer_dispatch(void* arg) { + struct test_timer* timer = arg; + timer->handle = NULL; + if (!ctx.result) timer->callback(NULL); +} +static void arm_timer(enum timer_id id, int delay, timeout_cb callback, const char* name) { + struct test_timer* timer = &timers[id]; + if (timer->handle) { fail("test timer already armed"); return; } + timer->callback = callback; + timer->handle = uasync_set_timeout(ctx.ua, delay, timer, timer_dispatch, name); + if (!timer->handle) fail("test timer allocation failed"); +} +static void diag_dump_state(const char* label); +static void to_cb(void* arg) { (void)arg; diag_dump_state("TIMEOUT"); ctx.result = 1; } static int link_on_iface(const struct ETCP_LINK* l, const char* ifname) { if (!l->conn || !l->conn->name || !ifname || !ifname[0]) return 0; @@ -119,8 +140,9 @@ static int count_all_client_tcp(void) { int n = 0; struct ETCP_SOCKET* s = ctx.c static int iface_still_exists(const char* ifname) { return if_nametoindex(ifname) != 0; } static int count_sockets_on_iface(const char* ifname) { - int n = 0; uint32_t idx = if_nametoindex(ifname); if (!idx) return 0; - for (struct ETCP_SOCKET* s = ctx.client->etcp_sockets; s; s = s->next) if (s->netif_index == idx) n++; + int n = 0; size_t len = strlen(ifname); + for (struct ETCP_SOCKET* s = ctx.client->etcp_sockets; s; s = s->next) + if (s->name && !strncmp(s->name, "as_", 3) && !strncmp(s->name + 3, ifname, len) && s->name[3 + len] == '_') n++; return n; } @@ -151,11 +173,10 @@ static struct ETCP_CONN* find_srv_conn(void) { static void check_conn_health(struct ETCP_CONN* conn) { if (!conn) return; if (conn->state == 2) fail("conn state is deleted (2)"); - if (conn->state != 1) { fprintf(stderr, "[warn] conn state=%d\n", conn->state); } - if (conn->reinit_count - ctx.conn_reinit_snapshot > ctx.max_allowed_reinits) fail("too many reinits"); + if (conn->state != 1 || conn->reinit_count) fail("unexpected connection state or reinit"); } -/* ── server node (2 addrs: UDP + TCP) ── */ +/* Сервер теста слушает UDP; TCP autosockets проверяются отдельно по жизненному циклу. */ static struct TOPO_GROUP_NODE* mk_srv_node(void) { struct TOPO_GROUP_NODE* nq = u_calloc(1, sizeof(struct TOPO_GROUP_NODE)); struct TOPO_NODE* ni = u_calloc(1, sizeof(struct TOPO_NODE)); @@ -167,84 +188,159 @@ static struct TOPO_GROUP_NODE* mk_srv_node(void) { struct TOPO_SOCKMETA4* sm = memory_pool_alloc(ctx.client->topo_groups->v4_sock_meta_pool); if (!sm) { u_free(nq); u_free(ni); return NULL; } + memset(sm, 0, sizeof(*sm)); sm->id = 0; sm->config_type = CFG_SERVER_TYPE_PUBLIC; sm->nat_type = NAT_TYPE_UNKNOWN; sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; struct TOPO_ADDR4* a_udp = memory_pool_alloc(ctx.client->topo_groups->v4_addr_pool); - struct TOPO_ADDR4* a_tcp = memory_pool_alloc(ctx.client->topo_groups->v4_addr_pool); - if (!a_udp || !a_tcp) { u_free(nq); u_free(ni); return NULL; } + if (!a_udp) { + memory_pool_free(ctx.client->topo_groups->v4_sock_meta_pool, sm); + u_free(nq); u_free(ni); return NULL; + } + memset(a_udp, 0, sizeof(*a_udp)); a_udp->addr[0] = 10; a_udp->addr[1] = 90; a_udp->addr[2] = 0; a_udp->addr[3] = 1; a_udp->port = 9001; a_udp->protocol = TOPO_PROTO_UDP; a_udp->type = TOPO_ADDR_INTERFACE; a_udp->socket_id = 0; - a_tcp->addr[0] = 10; a_tcp->addr[1] = 90; a_tcp->addr[2] = 0; a_tcp->addr[3] = 1; - a_tcp->port = 9001; a_tcp->protocol = TOPO_PROTO_TCP; - a_tcp->type = TOPO_ADDR_INTERFACE; a_tcp->socket_id = 0; - a_tcp->next = a_udp; ni->v4_addrs = a_tcp; + ni->v4_addrs = a_udp; topo_node_registry_store(ctx.client->topo_groups, ni); return nq; } +static void cancel_send(struct pending_send* p) { + if (!p->conn) return; + queue_waiter_cancel(p->conn->send_input_q, &p->waiter); + struct ETCP_CONN* conn = p->conn; p->conn = NULL; + etcp_conn_ref_free(conn); +} +static void wait_send(struct pending_send* p); +static void check_queues(struct pending_send* p) { + struct ETCP_CONN* c = p->conn; + if (!c || c->state == 2) return; + size_t bytes = queue_total_bytes(c->send_input_q) + queue_total_bytes(c->input_queue) + + queue_total_bytes(c->input_send_q) + queue_total_bytes(c->input_wait_ack); + if (bytes > p->peak_bytes) p->peak_bytes = bytes; + /* Один входной пакет может дать несколько фрагментов плюс flush хвоста. */ + struct PKTNORM* pn = c->normalizer; + size_t rx_bytes = queue_total_bytes(c->recv_q) + queue_total_bytes(c->output_queue) + queue_total_bytes(pn->output); + if (rx_bytes > p->peak_rx_bytes) p->peak_rx_bytes = rx_bytes; + if (rx_bytes > RX_BUDGET_BYTES) { + fprintf(stderr, "receive queue overflow dir=%ld bytes=%zu budget=%u\n", (long)(p - pending), rx_bytes, RX_BUDGET_BYTES); + fail("receive queues exceeded scenario memory budget"); + } + int fragments = (PROBE_SIZE + 2 + pn->frag_size - 1) / pn->frag_size + 1; + unsigned links = 0; + for (struct ETCP_LINK* l = c->links; l; l = l->next) links++; + unsigned probes = 0; + for (struct ll_entry* e = c->send_input_q->head; e; e = e->next) + if (e->len && e->dgram && e->dgram[0] == ETCP_RT_ID_TEST) probes++; + /* У сервера старые адреса остаются до keepalive timeout; лимит задан на каждый линк. */ + size_t limit = (links + 1) * 65536 + (fragments + 2) * c->mtu; + if (c->max_inflight != 65536 || links > 2 * N_ROUNDS + 3 || probes > 1 + || c->input_queue->count > fragments || bytes > limit) { + fprintf(stderr, "queue overflow dir=%ld app=%d input=%d send=%d ack=%d bytes=%zu\n", + (long)(p - pending), c->send_input_q->count, c->input_queue->count, + c->input_send_q->count, c->input_wait_ack->count, bytes); + fail("backpressure did not bound queues"); + } +} +static void send_ready(struct ll_queue* q, void* arg) { + struct pending_send* p = arg; + if (ctx.result || ctx.draining || !p->budget) return; + if (p->conn->state != 1 || p->conn->reinit_count != p->reinit) return; + if (q->count) { wait_send(p); return; } + struct ll_entry* entry = ll_alloc_lldgram(PROBE_SIZE); + if (!entry) { fail("traffic allocation failed"); return; } + uint32_t seq = ++p->sent, direction = p - pending; + entry->len = PROBE_SIZE; entry->dgram[0] = ETCP_RT_ID_TEST; + memcpy(entry->dgram + 1, &seq, 4); memcpy(entry->dgram + 5, &direction, 4); + for (int i = 9; i < PROBE_SIZE; i++) entry->dgram[i] = (uint8_t)(seq + direction + i * 17); + p->budget--; + if (etcp_send(p->conn, entry) != 0) { + queue_dgram_free(entry); queue_entry_free(entry); + fail("traffic send failed"); return; + } + check_queues(p); + if (p->budget && !ctx.result) wait_send(p); +} +static void wait_send(struct pending_send* p) { + int rc = queue_waiter_wait(p->conn->send_input_q, &p->waiter, send_ready, p); + if (rc == 0) p->blocked++; + if (rc < 0) fail("traffic backpressure waiter failed"); +} +static void receive_probe(struct ll_entry* entry, int direction) { + if (!entry || !entry->dgram || entry->len != PROBE_SIZE) { fail("invalid probe size"); goto done; } + uint32_t seq, wire_direction; + memcpy(&seq, entry->dgram + 1, 4); memcpy(&wire_direction, entry->dgram + 5, 4); + struct pending_send* p = &pending[direction]; + if (wire_direction != (uint32_t)direction || seq != p->received + 1 || seq > p->sent) { + fprintf(stderr, "sequence dir=%d wire_dir=%u received=%u seq=%u sent=%u\n", direction, wire_direction, p->received, seq, p->sent); + fail("missing, duplicate or out-of-order packet"); goto done; + } + for (int i = 9; i < PROBE_SIZE; i++) { + if (entry->dgram[i] != (uint8_t)(seq + direction + i * 17)) { fail("corrupt probe payload"); goto done; } + } + p->received++; + ctx.total_recv++; +done: + if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } +} static void srv_traffic_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || entry->len < 5) { if (entry) queue_entry_free(entry); return; } - struct ll_entry* reply = queue_entry_new(0); - if (!reply) { queue_entry_free(entry); return; } - reply->dgram = u_malloc(entry->len); - if (reply->dgram) { memcpy(reply->dgram, entry->dgram, entry->len); reply->len = entry->len; } - queue_entry_free(entry); - if (reply->dgram) etcp_send(conn, reply); - else queue_entry_free(reply); + (void)conn; receive_probe(entry, 0); } static void cli_traffic_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { - (void)conn; - if (!entry || entry->len < 5) { if (entry) queue_entry_free(entry); return; } - uint32_t seq; - memcpy(&seq, (uint8_t*)entry->dgram + 1, 4); - if (ctx.expected_rx_seq == 0) ctx.expected_rx_seq = seq; - if (seq == ctx.expected_rx_seq) { - ctx.expected_rx_seq++; - } else if (seq > ctx.expected_rx_seq) { - ctx.rx_seq_gaps++; - ctx.expected_rx_seq = seq + 1; - } else { - ctx.rx_seq_dups++; - } - ctx.pong_count++; ctx.total_recv++; - queue_entry_free(entry); + (void)conn; receive_probe(entry, 1); } static void traffic_send_timer(void* arg) { (void)arg; - if (ctx.result) return; - struct ETCP_CONN* conn = instance_find_conn(ctx.client, ctx.srv_node_id); - if (conn && conn->state != 2) { - uint8_t buf[5]; buf[0] = ETCP_RT_ID_TEST; - ctx.send_seq++; memcpy(buf + 1, &ctx.send_seq, 4); - struct ll_entry* e = queue_entry_new(0); - if (e) { e->dgram = u_malloc(5); if (e->dgram) { memcpy(e->dgram, buf, 5); e->len = 5; - etcp_send(conn, e); ctx.send_count++; } else queue_entry_free(e); } + if (ctx.result || ctx.draining) return; + for (int i = 0; i < 2; i++) { + struct pending_send* p = &pending[i]; + struct ETCP_CONN* conn = instance_find_conn(i ? ctx.server : ctx.client, + i ? ctx.client->config->global.my_node_id : ctx.srv_node_id); + if (p->conn && (p->conn != conn || p->conn->state == 2)) { + fail("owned traffic session unexpectedly replaced"); return; + } + if (!p->conn && conn && conn->state == 1) { + if (etcp_conn_ref_take(conn) != 0) { fail("cannot retain traffic connection"); return; } + p->conn = conn; + p->reinit = conn->reinit_count; + queue_set_threshold(conn->send_input_q, 0, 0); + } + if (!p->conn) continue; + if (p->conn->reinit_count != p->reinit) { + fail("unexpected traffic session reset"); return; + } + if (conn->state != 1) continue; + check_queues(p); + p->budget = SEND_BUDGET; + if (!p->waiter.internal && !p->waiter.call_soon_id) wait_send(p); } - if (!ctx.result) timeout_handle = uasync_set_timeout(ctx.ua, TRAF_SEND_TB, NULL, traffic_send_timer, "traf_snd"); + if (!ctx.result) arm_timer(TIMER_SEND, TRAF_SEND_TB, traffic_send_timer, "traf_snd"); } static void traffic_monitor_timer(void* arg) { (void)arg; if (ctx.result) return; - uint32_t d = ctx.pong_count; ctx.pong_count = 0; - fprintf(stderr, " [traf] r=%d tx=%u rx=%u+d=%u rate=%u/s gaps=%u dups=%u\n", - ctx.round, ctx.send_count, ctx.total_recv, d, d * 10, ctx.rx_seq_gaps, ctx.rx_seq_dups); + for (int i = 0; i < 2; i++) { + struct pending_send* p = &pending[i]; + check_queues(p); + fprintf(stderr, " [traf] r=%d dir=%d sent=%u received=%u blocked=%u peak_tx=%u peak_rx=%u\n", + ctx.round, i, p->sent, p->received, p->blocked, p->peak_bytes, p->peak_rx_bytes); + } fflush(stderr); - if (ctx.rx_seq_gaps > 100 || ctx.rx_seq_dups > 50) fail("too many seq anomalies"); - if (!ctx.result) timeout_handle = uasync_set_timeout(ctx.ua, TRAF_MON_TB, NULL, traffic_monitor_timer, "traf_mon"); + if (!ctx.result) arm_timer(TIMER_MONITOR, TRAF_MON_TB, traffic_monitor_timer, "traf_mon"); } /* ── etcp_connect callback ── */ -static void connect_cb(void* arg, struct ETCP_CONN* conn, int type) { - (void)arg; (void)conn; - if (type == ETCP_CONNECT_EARLY || type == ETCP_CONNECT_LATE) ctx.connected = 1; +static void connect_cb(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) { + (void)arg; (void)handle; + if (event == NCD_EVENT_UP) ctx.connected = 1; + if (event == NCD_EVENT_CLOSED || event == NCD_EVENT_TIMEOUT) fail("owned test connection closed or timed out"); } static void start_traffic(void) { - timeout_handle = uasync_set_timeout(ctx.ua, TRAF_SEND_TB, NULL, traffic_send_timer, "traf_snd"); - timeout_handle = uasync_set_timeout(ctx.ua, TRAF_MON_TB, NULL, traffic_monitor_timer, "traf_mon"); + arm_timer(TIMER_SEND, TRAF_SEND_TB, traffic_send_timer, "traf_snd"); + arm_timer(TIMER_MONITOR, TRAF_MON_TB, traffic_monitor_timer, "traf_mon"); } /* ── diagnostics ── */ @@ -294,24 +390,29 @@ static void diag_dump_state(const char* label) { /* ── ip addr add/del via system ── */ static int ip_addr_add(const char* ifname, const char* cidr) { - char cmd[256]; snprintf(cmd, sizeof(cmd), "ip addr add %s dev %s 2>/dev/null", cidr, ifname); - return system(cmd); + char cmd[256]; snprintf(cmd, sizeof(cmd), "ip addr add %s dev %s", cidr, ifname); + int rc = system(cmd); if (rc) fail(cmd); return rc; } static int ip_addr_del(const char* ifname, const char* cidr) { - char cmd[256]; snprintf(cmd, sizeof(cmd), "ip addr del %s dev %s 2>/dev/null", cidr, ifname); - return system(cmd); + char cmd[256]; snprintf(cmd, sizeof(cmd), "ip addr del %s dev %s", cidr, ifname); + int rc = system(cmd); if (rc) fail(cmd); return rc; } static int ip_link_add(const char* ifname) { - char cmd[256]; snprintf(cmd, sizeof(cmd), "ip link add %s type dummy 2>/dev/null", ifname); - return system(cmd); + char cmd[256]; snprintf(cmd, sizeof(cmd), "ip link add %s type dummy", ifname); + int rc = system(cmd); if (rc) fail(cmd); return rc; } static int ip_link_del(const char* ifname) { - char cmd[256]; snprintf(cmd, sizeof(cmd), "ip link del %s 2>/dev/null", ifname); - return system(cmd); + char cmd[256]; snprintf(cmd, sizeof(cmd), "ip link del %s", ifname); + int rc = system(cmd); if (rc) fail(cmd); return rc; } static int ip_link_up(const char* ifname) { - char cmd[256]; snprintf(cmd, sizeof(cmd), "ip link set %s up 2>/dev/null", ifname); - return system(cmd); + char cmd[256]; snprintf(cmd, sizeof(cmd), "ip link set %s up", ifname); + int rc = system(cmd); if (rc) fail(cmd); return rc; +} +static void add_default_route(const char* ifname) { + char cmd[256]; + snprintf(cmd, sizeof(cmd), "ip route replace default dev %s metric %u", ifname, if_nametoindex(ifname)); + if (system(cmd)) fail(cmd); } /* ═══════════════════════════════════════════════════════════ @@ -333,9 +434,20 @@ static void phase_done(void* arg); static void phase_done(void* arg) { (void)arg; if (ctx.result) return; - struct ETCP_CONN* conn = find_srv_conn(); - fprintf(stderr, "=== ALL PASSED === gaps=%u dups=%u reinit=%u\n", - ctx.rx_seq_gaps, ctx.rx_seq_dups, conn ? conn->reinit_count : 0); + ctx.draining = 1; + int outstanding = 0; + for (int i = 0; i < 2; i++) { + struct pending_send* p = &pending[i]; + struct ETCP_CONN* c = p->conn; + if (!c || !p->blocked || p->received < 100) { fail("insufficient traffic/backpressure coverage"); return; } + queue_waiter_cancel(c->send_input_q, &p->waiter); + struct PKTNORM* pn = c->normalizer; + outstanding += p->sent != p->received || c->send_input_q->count || c->input_queue->count + || c->input_send_q->count || c->input_wait_ack->count || c->recv_q->count || c->output_queue->count + || pn->output->count || pn->data_ptr || pn->recvpart; + } + if (outstanding) { arm_timer(TIMER_PHASE, STEP_TB / 3, phase_done, "drain"); return; } + fprintf(stderr, "=== TRAFFIC PASSED === delivered=%u/%u queues=empty reinit=0\n", pending[0].received, pending[1].received); fflush(stderr); ctx.result = 2; } @@ -344,9 +456,9 @@ static void phase_del_last(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 11; - diag_dump_state("before del_last"); + if (getenv("UTUN_TEST_DEBUG")) diag_dump_state("before del_last"); ip_link_del(rounds[N_ROUNDS-1].ifname); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_del_last, "chk_del_last"); + arm_timer(TIMER_PHASE, STEP_TB, phase_check_del_last, "chk_del_last"); } /* ══ check wrappers ══ */ @@ -359,6 +471,20 @@ static void do_check(void) { check_conn_health(conn); if (has_stale_links(conn)) fail("stale links detected"); } + unsigned ifindex = if_nametoindex(ctx.cur_iface); + uint32_t ip = get_interface_ip_by_index(ifindex); + int udp = 0, tcp = 0; + for (struct ETCP_SOCKET* s = ctx.client->etcp_sockets; s; s = s->next) { + if (strncmp(s->name + 3, ctx.cur_iface, strlen(ctx.cur_iface)) != 0) continue; + if (s->netif_index != ifindex || s->interface_addr.ss_family != AF_INET + || ((struct sockaddr_in*)&s->interface_addr)->sin_addr.s_addr != ip) { + fprintf(stderr, "socket binding: %s ifidx=%u expected=%u addr=%s\n", + s->name, s->netif_index, ifindex, sockaddr_storage_to_str(&s->interface_addr).str); + fail("socket bound to wrong interface/address"); return; + } + if (s->is_tcp) tcp++; else udp++; + } + if (udp != 1 || tcp != 1) fail("expected exactly one UDP and one TCP socket"); } static void phase_check_del_last(void* arg) { @@ -367,7 +493,7 @@ static void phase_check_del_last(void* arg) { int l = count_links_to_srv(rounds[N_ROUNDS-1].ifname); if (l > 0) { static int cnt = 0; - if (++cnt < 2) { uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_del_last, "chk_del_last"); return; } + if (++cnt < 2) { arm_timer(TIMER_PHASE, STEP_TB, phase_check_del_last, "chk_del_last"); return; } diag_dump_state("chk_del_last fail"); fail("links survived last del"); return; } @@ -375,13 +501,11 @@ static void phase_check_del_last(void* arg) { int ntotal = count_all_client_udp() + count_all_client_tcp(); int ndummy = 0; for (int i = 0; i < N_ROUNDS; i++) ndummy += count_sockets_on_iface(rounds[i].ifname); - if (ndummy > 0) fail("dummy sockets still exist after all deleted"); + if (ndummy > 0 || ntotal > 0) fail("sockets still exist after all client interfaces deleted"); fprintf(stderr, " r=%d s=%d: no links after del_last total_socks=%d dummy_socks=%d (OK)\n", ctx.round, ctx.step, ntotal, ndummy); fflush(stderr); - if (!ctx.all_deleted); - ctx.all_deleted = 1; - uasync_set_timeout(ctx.ua, STEP_TB * 3, NULL, phase_full_disconnect, "full_disconnect"); + arm_timer(TIMER_PHASE, STEP_TB * 3, phase_full_disconnect, "full_disconnect"); } #define MIN_RECOVERY_PKTS 5 @@ -390,14 +514,13 @@ static void phase_check_ip(void* arg) { (void)arg; if (ctx.result) return; static uint64_t wait_start = 0; static int diag_cnt = 0; - static uint32_t recv_ok; ctx.step = 6; - if (count_links_to_srv(ctx.cur_iface) == 0 || ctx.total_recv <= ctx.recv_at_ip_change + MIN_RECOVERY_PKTS) { + if (count_links_to_srv(ctx.cur_iface) == 0 || pending[0].received <= pending[0].snapshot + MIN_RECOVERY_PKTS + || pending[1].received <= pending[1].snapshot + MIN_RECOVERY_PKTS) { if (++diag_cnt <= 3) diag_dump_state("wait chk_ip"); - if (!wait_start) { wait_start = get_time_tb(); recv_ok = 0; } - else if (ctx.total_recv > ctx.recv_at_ip_change) recv_ok++; + if (!wait_start) wait_start = get_time_tb(); if (get_time_tb() - wait_start < (uint64_t)STEP_TB * 10) { - uasync_set_timeout(ctx.ua, STEP_TB/3, NULL, phase_check_ip, "chk_ip"); return; + arm_timer(TIMER_PHASE, STEP_TB/3, phase_check_ip, "chk_ip"); return; } diag_dump_state("timeout chk_ip"); fail("traffic not recovered after IP change"); @@ -410,71 +533,92 @@ static void phase_check_ip(void* arg) { if (ctx.round == N_ROUNDS - 1 && ctx.ip_changes_on_last < 2) { ctx.ip_changes_on_last++; fprintf(stderr, " r=%d s=%d: IP change #%d verified — another\n", ctx.round, ctx.step, ctx.ip_changes_on_last); fflush(stderr); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip2"); + arm_timer(TIMER_PHASE, STEP_TB, phase_change_ip, "chg_ip2"); } else if (ctx.round == N_ROUNDS - 1 && ctx.ip_changes_on_last >= 2) { - fprintf(stderr, " r=%d s=%d: both IP changes on last round — deleting\n", ctx.round, ctx.step); fflush(stderr); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_del_last, "del_last"); + fprintf(stderr, " r=%d s=%d: three IP changes verified — deleting\n", ctx.round, ctx.step); fflush(stderr); + arm_timer(TIMER_PHASE, STEP_TB, phase_del_last, "del_last"); } else { strncpy(ctx.prev_iface, ctx.cur_iface, IFNAMSIZ - 1); ctx.round++; - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_add, "next_add"); + arm_timer(TIMER_PHASE, STEP_TB, phase_add, "next_add"); } } static void phase_change_ip(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 5; - const char* old_ip = (ctx.ip_changes_on_last >= 2) ? rounds[ctx.round].ip2 : rounds[ctx.round].ip1; - const char* new_ip = (ctx.ip_changes_on_last >= 2) ? rounds[ctx.round].ip1 : rounds[ctx.round].ip2; + const char* old_ip = (ctx.ip_changes_on_last % 2) ? rounds[ctx.round].ip2 : rounds[ctx.round].ip1; + const char* new_ip = (ctx.ip_changes_on_last % 2) ? rounds[ctx.round].ip1 : rounds[ctx.round].ip2; ip_addr_del(rounds[ctx.round].ifname, old_ip); ip_addr_add(rounds[ctx.round].ifname, new_ip); - ctx.recv_at_ip_change = ctx.total_recv; - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_ip, "chk_ip"); + add_default_route(rounds[ctx.round].ifname); + for (int i = 0; i < 2; i++) pending[i].snapshot = pending[i].received; + arm_timer(TIMER_PHASE, STEP_TB, phase_check_ip, "chk_ip"); } static void phase_check_del(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 4; do_check(); if (ctx.result) return; + if (ctx.tcp_probe) { + if (ctx.tcp_probe->links) { fail("TCP link outside NCD survived socket deletion"); return; } + etcp_connection_close(ctx.tcp_probe); etcp_conn_ref_free(ctx.tcp_probe); ctx.tcp_probe = NULL; + } /* сокеты удалённого интерфейса должны исчезнуть */ int stale = count_sockets_on_iface(ctx.prev_iface); - if (stale > 0) fail("sockets survived for deleted iface"); + if (stale > 0) { fail("sockets survived for deleted iface"); return; } + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(ctx.client->topo_sqlite_db, "SELECT count(*) FROM auto_socket_ports WHERE if_name=?", + -1, &stmt, NULL) != SQLITE_OK) { fail("cannot inspect autosocket port records"); return; } + sqlite3_bind_text(stmt, 1, ctx.prev_iface, -1, SQLITE_TRANSIENT); + if (sqlite3_step(stmt) != SQLITE_ROW || sqlite3_column_int(stmt, 0) != 0) fail("port records survived interface deletion"); + sqlite3_finalize(stmt); + if (ctx.result) return; fprintf(stderr, " r=%d s=%d: del_prev OK links=%d stale_socks=%d\n", ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface), stale); fflush(stderr); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip"); + arm_timer(TIMER_PHASE, STEP_TB, phase_change_ip, "chg_ip"); } static void phase_del_prev(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 3; + if (ctx.round == 1) { + /* TCP link без NCD-владельца тоже должен закрыться до освобождения сокета. */ + struct ETCP_SOCKET* sock = ctx.client->etcp_sockets; + while (sock && (!sock->is_tcp || sock->netif_index != if_nametoindex(ctx.prev_iface))) sock = sock->next; + if (!sock) { fail("missing TCP socket for lifecycle probe"); return; } + ctx.tcp_probe = etcp_connection_create(ctx.client, "tcp_remove_probe"); + if (!ctx.tcp_probe || etcp_conn_ref_take(ctx.tcp_probe) != 0) { fail("cannot create TCP lifecycle probe"); return; } + struct sockaddr_storage addr = {0}; + struct sockaddr_in* sin = (struct sockaddr_in*)&addr; + sin->sin_family = AF_INET; sin->sin_port = htons(9001); inet_pton(AF_INET, "10.90.0.1", &sin->sin_addr); + if (!etcp_link_new(ctx.tcp_probe, sock, &addr, 0)) { fail("cannot create TCP probe link"); return; } + } ip_link_del(ctx.prev_iface); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_del, "chk_del"); + arm_timer(TIMER_PHASE, STEP_TB, phase_check_del, "chk_del"); } static void phase_check_add(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 2; if (ctx.round == 0 && !ctx.connected) { - uasync_set_timeout(ctx.ua, STEP_TB / 3, NULL, phase_check_add, "chk_add"); return; + arm_timer(TIMER_PHASE, STEP_TB / 3, phase_check_add, "chk_add"); return; } do_check(); if (ctx.result) return; /* на текущем интерфейсе должен быть хотя бы 1 UDP и 1 TCP сокет */ int socks = count_sockets_on_iface(ctx.cur_iface); - if (socks < 2) { fprintf(stderr, "[warn] only %d sockets on %s\n", socks, ctx.cur_iface); } + if (socks != 2) { fail("expected one UDP and one TCP socket on current interface"); return; } + if (count_sockets_on_iface("dummy_srv")) { fail("interface without default route was not excluded"); return; } fprintf(stderr, " r=%d s=%d: add OK links=%d socks=%d udp=%d tcp=%d\n", ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface), socks, count_all_client_udp(), count_all_client_tcp()); fflush(stderr); if (ctx.round == 0) start_traffic(); - /* snapshot reinit на входе в цикл */ - struct ETCP_CONN* conn = find_srv_conn(); - if (conn) { ctx.conn_reinit_snapshot = conn->reinit_count; ctx.max_allowed_reinits = 2; } - if (ctx.prev_iface[0]) { - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_del_prev, "del_prev"); + arm_timer(TIMER_PHASE, STEP_TB, phase_del_prev, "del_prev"); } else { - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip"); + arm_timer(TIMER_PHASE, STEP_TB, phase_change_ip, "chg_ip"); } } @@ -483,27 +627,32 @@ static void phase_add(void* arg) { ctx.step = 1; strncpy(ctx.cur_iface, rounds[ctx.round].ifname, IFNAMSIZ - 1); - ip_link_add(ctx.cur_iface); - ip_addr_add(ctx.cur_iface, rounds[ctx.round].ip1); - ip_link_up(ctx.cur_iface); + if (ctx.round > 0) { + ip_link_add(ctx.cur_iface); + ip_addr_add(ctx.cur_iface, rounds[ctx.round].ip1); + ip_link_up(ctx.cur_iface); + add_default_route(ctx.cur_iface); + } if (ctx.round == 0) { struct TOPO_GROUP_NODE* sn = mk_srv_node(); if (!sn) { fail("mk_srv_node"); return; } queue_data_put_with_index(topo_groups_get_default(ctx.client->topo_groups)->nodes, &sn->ll); - etcp_connect(ctx.client, sn, connect_cb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE); + if (node_conn_direct_open(ctx.client, ctx.srv_node_id, connect_cb, NULL, &ctx.handle, NULL) < 0) + fail("cannot open test connection"); } - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_add, "chk_add"); + arm_timer(TIMER_PHASE, STEP_TB, phase_check_add, "chk_add"); } /* ── full disconnect + reconnect ── */ static void phase_full_disconnect(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 20; - for (int i = 0; i < N_ROUNDS; i++) ip_link_del(rounds[i].ifname); ctx.pre_delete_total_recv = ctx.total_recv; + for (int i = 0; i < 2; i++) pending[i].snapshot = pending[i].received; fprintf(stderr, " r=%d s=%d: all interfaces deleted — awaiting reconnect\n", ctx.round, ctx.step); fflush(stderr); - uasync_set_timeout(ctx.ua, STEP_TB * 5, NULL, phase_reconnect, "reconnect"); + for (int i = 0; i < 2; i++) pending[i].outage_sent = pending[i].sent; + arm_timer(TIMER_PHASE, STEP_TB * 5, phase_reconnect, "reconnect"); } static void phase_reconnect(void* arg) { @@ -513,27 +662,30 @@ static void phase_reconnect(void* arg) { struct ETCP_CONN* conn = find_srv_conn(); int live = conn ? count_live_links_on_conn(conn) : 0; fprintf(stderr, " r=%d s=%d: live_links=%d\n", ctx.round, ctx.step, live); fflush(stderr); + if (live || ctx.total_recv != ctx.pre_delete_total_recv) { fail("traffic survived full network outage"); return; } + for (int i = 0; i < 2; i++) { + if (pending[i].sent != pending[i].outage_sent) { fail("producer did not stop during outage"); return; } + } /* создаём новый интерфейс */ char new_if[] = "dummy_reconn"; ip_link_add(new_if); ip_addr_add(new_if, "10.90.1.100/16"); ip_link_up(new_if); + add_default_route(new_if); strncpy(ctx.cur_iface, new_if, IFNAMSIZ - 1); - ctx.recv_at_ip_change = ctx.total_recv; - conn = find_srv_conn(); - ctx.conn_reinit_snapshot = conn ? conn->reinit_count : 0; - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_reconnect, "chk_reconn"); + arm_timer(TIMER_PHASE, STEP_TB, phase_check_reconnect, "chk_reconn"); } static void phase_check_reconnect(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 22; static uint64_t wait_start = 0; - if (count_links_to_srv(ctx.cur_iface) == 0 || ctx.total_recv <= ctx.pre_delete_total_recv + MIN_RECOVERY_PKTS) { + if (count_links_to_srv(ctx.cur_iface) == 0 || pending[0].received <= pending[0].snapshot + MIN_RECOVERY_PKTS + || pending[1].received <= pending[1].snapshot + MIN_RECOVERY_PKTS) { if (!wait_start) wait_start = get_time_tb(); if (get_time_tb() - wait_start < (uint64_t)STEP_TB * 20) { - uasync_set_timeout(ctx.ua, STEP_TB/3, NULL, phase_check_reconnect, "chk_reconn"); return; + arm_timer(TIMER_PHASE, STEP_TB/3, phase_check_reconnect, "chk_reconn"); return; } diag_dump_state("reconnect timeout"); fail("connection not recovered after full disconnect"); @@ -544,25 +696,19 @@ static void phase_check_reconnect(void* arg) { do_check(); if (ctx.result) return; struct ETCP_CONN* conn = find_srv_conn(); check_conn_health(conn); - ip_link_del("dummy_reconn"); fprintf(stderr, " r=%d s=%d: reconnect OK links=%d reinit=%u\n", ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface), conn ? conn->reinit_count : 0); fflush(stderr); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_done, "done"); + arm_timer(TIMER_PHASE, STEP_TB, phase_done, "done"); } /* ═══════════════════════════════════════════════════════════ * setup / cleanup / main * ═══════════════════════════════════════════════════════════ */ -static void cleanup_ifaces(void) { - (void)system("ip link del dummy_srv 2>/dev/null"); - for (int i = 0; i < N_ROUNDS; i++) ip_link_del(rounds[i].ifname); -} - static void setup(void) { if (geteuid() != 0) { fprintf(stderr, "SKIP: test requires root\n"); exit(77); } - cleanup_ifaces(); - test_mkdtemp(tdir); + if (unshare(CLONE_NEWNET) != 0) { perror("SKIP: test requires a private network namespace"); exit(77); } + if (ip_link_up("lo") != 0 || test_mkdtemp(tdir) != 0) { fail("test setup failed"); return; } snprintf(scf, sizeof(scf), "%s/s.conf", tdir); snprintf(ccf, sizeof(ccf), "%s/c.conf", tdir); @@ -579,13 +725,14 @@ static void setup(void) { ip_link_add(rounds[0].ifname); ip_addr_add(rounds[0].ifname, rounds[0].ip1); ip_link_up(rounds[0].ifname); + add_default_route(rounds[0].ifname); /* Step 1: write minimal configs → generate keys + node_id via config_ensure */ wf(scf, "[global]\ntun_ip=10.99.0.1/24\ntun_ifname=tun_srv\ntun_test_mode=1\n" - "auto_sockets=no\ndb_path=%s/db_srv\n" + "auto_sockets=no\nbbr_max_cwnd=65536\ndb_path=%s/db_srv\n" "[server: fixed]\naddr=10.90.0.1:9001\ntype=public\n[allowed_keys]\nallow_all=1\n", tdir); wf(ccf, "[global]\ntun_ip=10.99.0.2/24\ntun_ifname=tun_cli\ntun_test_mode=1\n" - "auto_sockets=yes\ndb_path=%s/db_cli\n[allowed_keys]\nallow_all=1\n", tdir); + "auto_sockets=yes\nbbr_max_cwnd=65536\ndb_path=%s/db_cli\n[allowed_keys]\nallow_all=1\n", tdir); config_ensure_keys_and_node_id(scf); config_ensure_keys_and_node_id(ccf); @@ -597,12 +744,12 @@ static void setup(void) { /* Step 3: write final configs with all values embedded */ wf(scf, "[global]\nmy_private_key=%s\nmy_public_key=%s\nmy_node_id=0x%016llx\n" "tun_ip=10.99.0.1/24\ntun_ifname=tun_srv\ntun_test_mode=1\n" - "auto_sockets=no\ndb_path=%s/db_srv\n" + "auto_sockets=no\nbbr_max_cwnd=65536\ndb_path=%s/db_srv\n" "[server: fixed]\naddr=10.90.0.1:9001\ntype=public\n[allowed_keys]\nallow_all=1\n", spriv, spub, (unsigned long long)ctx.srv_node_id, tdir); wf(ccf, "[global]\nmy_private_key=%s\nmy_public_key=%s\n" "tun_ip=10.99.0.2/24\ntun_ifname=tun_cli\ntun_test_mode=1\n" - "auto_sockets=yes\ndb_path=%s/db_cli\n[allowed_keys]\nallow_all=1\n", + "auto_sockets=yes\nbbr_max_cwnd=65536\ndb_path=%s/db_cli\n[allowed_keys]\nallow_all=1\n", cpriv, cpub, tdir); /* extract server pubkey binary */ @@ -614,13 +761,25 @@ static void setup(void) { u_free(spub); u_free(spriv); u_free(cpub); u_free(cpriv); } -static void cleanup(void) { test_unlink(scf); test_unlink(ccf); test_rmdir(tdir); } +static int cleanup_path(const char* path, const struct stat* st, int type, struct FTW* ftw) { + (void)st; (void)type; (void)ftw; + if (remove(path)) { perror(path); return -1; } + return 0; +} +static void cleanup(void) { + if (nftw(tdir, cleanup_path, 16, FTW_DEPTH | FTW_PHYS) != 0) fail("temporary directory cleanup failed"); +} int main(void) { debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); + if (getenv("UTUN_TEST_DEBUG")) { + debug_set_category_level(DEBUG_CATEGORY_SOCKET, DEBUG_LEVEL_INFO); + debug_set_category_level(DEBUG_CATEGORY_CONNECTION, DEBUG_LEVEL_INFO); + } utun_instance_set_tun_init_enabled(0); setup(); + if (ctx.result) return 1; ctx.ua = uasync_create(); ctx.server = utun_instance_create(ctx.ua, scf); @@ -638,9 +797,9 @@ int main(void) { strncpy(ctx.prev_iface, rounds[0].ifname, IFNAMSIZ - 1); ctx.prev_iface[0] = '\0'; /* round 0 has no prev */ - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_add, "init"); + arm_timer(TIMER_PHASE, STEP_TB, phase_add, "init"); - timeout_handle = uasync_set_timeout(ctx.ua, TIMEOUT_TB, NULL, to_cb, "to"); + arm_timer(TIMER_DEADLINE, TIMEOUT_TB, to_cb, "to"); { uint64_t start = get_time_tb(); while (!ctx.result && (int)(get_time_tb() - start) < TIMEOUT_TB + 50000) @@ -649,10 +808,19 @@ int main(void) { fprintf(stderr, "final result=%d\n", ctx.result); fflush(stderr); done: - if (timeout_handle) uasync_cancel_timeout(ctx.ua, timeout_handle); + for (int i = 0; i < TIMER_COUNT; i++) if (timers[i].handle) uasync_cancel_timeout(ctx.ua, timers[i].handle); + for (int i = 0; i < 2; i++) cancel_send(&pending[i]); + if (ctx.tcp_probe) { etcp_connection_close(ctx.tcp_probe); etcp_conn_ref_free(ctx.tcp_probe); ctx.tcp_probe = NULL; } + if (ctx.handle) { node_conn_direct_force_close(ctx.handle); ctx.handle = NULL; } if (ctx.server) { ctx.server->running = 0; utun_instance_destroy(ctx.server); } if (ctx.client) { ctx.client->running = 0; utun_instance_destroy(ctx.client); } - if (ctx.ua) { uasync_destroy(ctx.ua, 0); ctx.ua = NULL; } + if (ctx.ua) { + uasync_poll(ctx.ua, 0); + if (ctx.ua->timer_alloc_count != ctx.ua->timer_free_count) fail("timers leaked after cleanup"); + if (ctx.ua->socket_alloc_count != ctx.ua->socket_free_count) fail("sockets leaked after cleanup"); + uasync_destroy(ctx.ua, 0); ctx.ua = NULL; + } cleanup(); + fprintf(stderr, "=== %s === timers/sockets cleaned\n", ctx.result == 2 ? "PASS" : "FAIL"); return (ctx.result == 2) ? 0 : 1; }