Browse Source

Fix autosocket ownership and exercise interface churn under load

proxy
evgeny 5 days ago
parent
commit
59f4802904
  1. 104
      src/transport_layer/auto_socket.c
  2. 22
      src/transport_layer/etcp_connections.c
  3. 464
      tests/test_auto_socket_dynamic.c

104
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;
}

22
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);
}

464
tests/test_auto_socket_dynamic.c

@ -1,17 +1,20 @@
#define _GNU_SOURCE
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdarg.h>
#include <ftw.h>
#include "../lib/platform_compat.h"
#include "test_utils.h"
#ifndef _WIN32
#include <unistd.h>
#include <net/if.h>
#include <sched.h>
#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;
}

Loading…
Cancel
Save