diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 238e5be3..be2ceaf3 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -9,6 +9,10 @@ 1. Модуль подключения p2p (ETCP/STCP) Позволяет подключиться двум узлам друг к другу. Весь трафик шифруется. Для подключения нужен публичный ключ удаленного узла. Позволяет использовать несколько линков до узла параллельно и динамически их переключать (агрегация/балансировка нагрузки/failover) +Модуль может работать в автоматическим режиме: мониторит изменения на интерфейсах и при появлении нового сетевого интерфейса автоматически добавлять подключения через новые интерфейсы. +При этом подключения работат параллельно по всем доступным интерфейсам распределяя нагрузку (какой интерфейс быстрее - тот бОльшую нагрузку берет на себя). +Разуммется, бесшовная работа - если есть хотяюы один живой линк то соединение работает. В процессе могут добавляться-удаляться линки. + 2. Логические группы узлов (topo_group) Один сервис может работать с несколькими группами. diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 43be2255..a3702430 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -653,6 +653,20 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { u_free(cs); } +void chat_sync_retry_channels_on_socket_change(struct UTUN_INSTANCE* inst) { + if (!g_cs || !g_cs->initialized || !inst->topo_groups) return; + for (int i = 0; i < g_cs->channel_count; i++) { + uint64_t gid = strtoull(g_cs->channels[i].channel_id, NULL, 10); + struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, gid); + if (!group || !group->connect) continue; + if (topo_group_connect_active_count(group) == 0) { + DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: socket change — restarting connect for ch=%s (no active connections)", + CS_ID, g_cs->channels[i].channel_id); + topo_group_connect_restart(group); + } + } +} + void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { (void)inst; (void)node_id; } diff --git a/src/chat/chat_sync.h b/src/chat/chat_sync.h index 21912a35..de171b72 100644 --- a/src/chat/chat_sync.h +++ b/src/chat/chat_sync.h @@ -108,6 +108,10 @@ void chat_sync_deny_invite(struct UTUN_INSTANCE* inst, uint64_t channel_id, uint Канал (DB + TOPO_GROUP) создаётся локально при получении CHANNEL_INFO_RESP. */ void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t target_node_id); +/* Перезапустить групповые подключения для каналов без активных соединений. + Вызывается из auto_socket при появлении/изменении сетевой связности. */ +void chat_sync_retry_channels_on_socket_change(struct UTUN_INSTANCE* inst); + #ifdef __cplusplus } #endif diff --git a/src/routing_layer/topo_group_connect.c b/src/routing_layer/topo_group_connect.c index b3a71a73..be6bd328 100644 --- a/src/routing_layer/topo_group_connect.c +++ b/src/routing_layer/topo_group_connect.c @@ -109,9 +109,8 @@ void topo_group_connect_destroy(struct TOPO_GROUP* group) { void topo_group_connect_restart(struct TOPO_GROUP* group) { struct TOPO_GROUP_CONNECT* gc = group->connect; if (!gc) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s no gc → init", TGC_ID, group->channel_id); topo_group_connect_init(group); return; } - if (gc->phase != TGC_PHASE_DONE) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s skip phase=%d", TGC_ID, group->channel_id, gc->phase); return; } if (gc->active_conn_count > 0) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s skip active=%d", TGC_ID, group->channel_id, gc->active_conn_count); return; } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s", TGC_ID, group->channel_id); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s phase=%d", TGC_ID, group->channel_id, gc->phase); topo_group_connect_destroy(group); topo_group_connect_init(group); } diff --git a/src/transport_layer/auto_socket.c b/src/transport_layer/auto_socket.c index 7044c61c..6d9dd498 100644 --- a/src/transport_layer/auto_socket.c +++ b/src/transport_layer/auto_socket.c @@ -1,25 +1,45 @@ /* - * auto_socket.c — автоматическое управление UDP-сокетами при изменении сетевых интерфейсов + * auto_socket.c — автоматическое управление сокетами (UDP+TCP) при изменении сетевых интерфейсов * - * При auto_sockets=yes: - * - Для каждого не-loopback интерфейса (UP, с IP-адресами) создаётся UDP-сокет - * с bind на 0.0.0.0/[::] + случайный порт, привязанный к интерфейсу (SO_BINDTODEVICE) - * - Тип: PUBLIC если есть публичный адрес, иначе NAT - * - При добавлении сокета — линки ко всем активным NCD-соединениям - * - При удалении интерфейса — линки удаляются, сокет закрывается + * === Зачем === + * При auto_sockets=yes модуль сам создаёт UDP и TCP сокеты на всех не-loopback + * интерфейсах с IP-адресами. Не нужно вручную описывать серверы в конфиге. * - * Стратегия reconcile: - * При любом изменении адресов/линков сканируем фактические IP интерфейса, - * сравниваем с состоянием сокетов и устраняем расхождения: + * === Что создаёт === + * Для каждого интерфейса с v4/v6 адресами: + * 1. UDP-сокет (etcp_socket_add) — bind на 0.0.0.0/[::] + SO_BINDTODEVICE, случайный порт + * 2. TCP-сокет (tcp_socket_add + stcp_server) — слушает входящие STCP-подключения + * 3. Для каждого сокета сразу добавляются линки (etcp_link_new) ко ВСЕМ активным + * ETCP-соединениям (inst->connections), случайно выбирая один совместимый + * адрес пира отдельно для v4 и v6. + * + * === Тип сокета === + * PUBLIC — если на интерфейсе есть глобальный (не RFC1918) адрес; иначе NAT. + * При смене типа (например VPN включился/выключился) сокет пересоздаётся. + * + * === Persistence портов === + * Порт сохраняется в SQLite-таблицу auto_socket_ports и переиспользуется + * при перезапуске utun. Если сохранённый порт занят — генерируется новый случайный. + * + * === Стратегия reconcile === + * При любом изменении адресов/линков сканируем фактические IP интерфейса, + * сравниваем с состоянием сокетов и устраняем расхождения: * - появился адрес нужного family → создать сокет + линки * - исчезли все адреса family → удалить сокет + линки * - изменился тип (NAT↔PUBLIC) → пересоздать сокет с новым типом + линки - * - без изменений → только обновить interface_addr + * - изменился сам адрес → обновить interface_addr, оповестить подписчиков * - интерфейс стал DOWN → удалить все сокеты * - * Слои: - * Platform-agnostic: auto_socket_add_interface / auto_socket_remove_interface - * Platform monitoring: netlink (Linux), route socket (BSD), IP Helper (Windows) + * === Линки === + * Каждый созданный сокет получает линки ко всем активным соединениям через + * ncd_add_socket_links() — по одному случайному адресу пира на каждый address family. + * На линке сохраняется local_bound_addr (interface_addr в момент создания) для + * последующей проверки принадлежности. + * + * === Platform monitoring === + * Встроенные обработчики: netlink (Linux), route socket (BSD), IP Helper (Windows). + * При отсутствии поддержки платформы мониторинг неактивен — можно вызывать + * auto_socket_add_interface / auto_socket_on_network_change извне (JNI и т.д.). */ #include "auto_socket.h" @@ -31,6 +51,7 @@ #include "utun_instance.h" #include "config_parser.h" #include "../chat/chat_event.h" +#include "../chat/chat_sync.h" #include "../lib/debug_config.h" #include "../lib/mem.h" #include "../lib/u_async.h" @@ -65,7 +86,9 @@ #define AS_PROTO_UDP 0 #define AS_PROTO_TCP 1 -/* ─── IPv4-адрес: публичный (интернет) или приватный/NAT ─── */ +/* Классифицирует IPv4-адрес: 1 = глобальный/публичный (можно использовать как PUBLIC-сокет), + * 0 = приватный, CGNAT, loopback, link-local, multicast или зарезервированный. + * Определяет тип создаваемого сокета по адресам на интерфейсе. */ static int is_ipv4_public(uint32_t addr_be) { uint32_t a = ntohl(addr_be); if ((a & 0xFF000000u) == 0x0A000000u) return 0; /* 10.0.0.0/8 */ @@ -80,7 +103,11 @@ static int is_ipv4_public(uint32_t addr_be) { return 1; } -/* ─── IPv6-адрес: глобальный или локальный (аналог cm_classify_v6_addr) ─── */ +/* Классифицирует IPv6-адрес для определения типа сокета: + * AS_V6_LL = link-local (fe80::) — создаём NAT-сокет + * AS_V6_LOC = ULA (fc00::/fd00::) — локальный, NAT-сокет + * AS_V6_DIR = глобальный — создаём PUBLIC-сокет + * AS_V6_OTH = loopback, multicast, :: — пропускаем */ #define AS_V6_LL 0 #define AS_V6_LOC 1 #define AS_V6_DIR 2 @@ -95,31 +122,40 @@ static int v6_classify(const uint8_t addr[16]) { return AS_V6_DIR; } -/* ─── Внутреннее состояние модуля ─── */ - +/* Один отслеживаемый интерфейс: хранит указатели на созданные сокеты (UDP — ETCP_SOCKET, + * TCP — stcp_server) и их актуальный тип (NAT/PUBLIC). По типу определяется когда + * нужно пересоздать сокет при изменении адресов. */ struct auto_sock_iface { struct auto_sock_iface* next; - uint32_t netif_index; - struct ETCP_SOCKET* v4_udp; - struct ETCP_SOCKET* v6_udp; - struct stcp_server* v4_tcp; - struct stcp_server* v6_tcp; - uint8_t v4_type; - uint8_t v6_type; + uint32_t netif_index; // индекс интерфейса (if_nametoindex) + 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 + uint8_t v4_type; // CFG_SERVER_TYPE_PUBLIC или CFG_SERVER_TYPE_NAT + uint8_t v6_type; // CFG_SERVER_TYPE_PUBLIC или CFG_SERVER_TYPE_NAT }; +/* Внутреннее состояние модуля auto_socket: хранит связный список отслеживаемых + * интерфейсов и платформенные handle'ы для мониторинга изменений сети. + * Создаётся в auto_socket_init, уничтожается в auto_socket_destroy. + * Хранится в inst->auto_socket_state. */ struct AUTO_SOCKET { - struct UTUN_INSTANCE* instance; - struct auto_sock_iface* ifaces; + struct UTUN_INSTANCE* instance; // родительский utun instance + struct auto_sock_iface* ifaces; // связный список интерфейсов + uint8_t v4_usable; // есть ли сейчас рабочий v4-сокет с реальным IP + uint8_t v6_usable; // есть ли сейчас рабочий v6-сокет с реальным IP + uint8_t v4_addr_changed; // v4-адрес или тип изменился с прошлой проверки + uint8_t v6_addr_changed; // v6-адрес или тип изменился с прошлой проверки #ifdef __linux__ - socket_t nl_sock; - void* uasync_handle; + socket_t nl_sock; // netlink socket fd для RTMGRP_LINK/IPV4/IPV6 + void* uasync_handle; // handle в uasync для netlink fd #elif defined(__FreeBSD__) || defined(__FreeBSD_kernel__) - socket_t route_sock; - void* uasync_handle; + socket_t route_sock; // route socket fd (PF_ROUTE) + void* uasync_handle; // handle в uasync для route fd #elif defined(_WIN32) - HANDLE addr_notify_handle; - HANDLE route_notify_handle; + HANDLE addr_notify_handle; // NotifyUnicastIpAddressChange handle + HANDLE route_notify_handle; // NotifyRouteChange2 handle #endif }; @@ -131,8 +167,12 @@ static int scan_and_classify_iface(uint32_t ifindex, const char* ifname, int* o static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname); static void auto_socket_scan_all(struct AUTO_SOCKET* as); -/* ─── DB helpers (port persistence) ─── */ +/* ─── DB helpers (port persistence) ─── + * Сохраняют порт сокета в SQLite чтобы переиспользовать его при перезапуске utun. + * Это сохраняет NAT-привязку (endpoint mapping) на роутере и позволяет пирам + * продолжать достукиваться на тот же порт. */ +/* Возвращает сохранённый порт для указанного интерфейса/family/протокола, или 0 если нет. */ static uint16_t load_port_from_db(struct AUTO_SOCKET* as, const char* ifname, int family, int protocol) { sqlite3* db = as->instance->topo_sqlite_db; if (!db) return 0; @@ -149,6 +189,7 @@ static uint16_t load_port_from_db(struct AUTO_SOCKET* as, const char* ifname, in return port; } +/* Сохраняет порт в БД для переиспользования при следующем запуске. */ static void save_port_to_db(struct AUTO_SOCKET* as, const char* ifname, int family, int protocol, uint16_t port) { sqlite3* db = as->instance->topo_sqlite_db; if (!db) return; @@ -169,6 +210,7 @@ static void save_port_to_db(struct AUTO_SOCKET* as, const char* ifname, int fami sqlite3_finalize(stmt); } +/* Удаляет запись порта из БД после закрытия сокета. */ static void delete_port_from_db(struct AUTO_SOCKET* as, const char* ifname, int family, int protocol) { sqlite3* db = as->instance->topo_sqlite_db; if (!db) return; @@ -183,19 +225,62 @@ static void delete_port_from_db(struct AUTO_SOCKET* as, const char* ifname, int sqlite3_finalize(stmt); } +/* Отправляет событие CHAT_EVT_LOCAL_SOCKETS в GUI/headless — оповещает что + * список локальных сокетов изменился (добавился/удалился интерфейс). */ +static int as_check_usable(struct AUTO_SOCKET* as, int family) { + struct ETCP_SOCKET* s = as->instance->etcp_sockets; + while (s) { + if (s->local_addr.ss_family == family && s->interface_addr.ss_family == family) return 1; + s = s->next; + } + return 0; +} + static void auto_socket_post_sockets_changed(struct AUTO_SOCKET* as) { + int old_v4 = as->v4_usable, old_v6 = as->v6_usable; + as->v4_usable = as_check_usable(as, AF_INET); + as->v6_usable = as_check_usable(as, AF_INET6); + + int need_reconnect = 0; + if ((!old_v4 && as->v4_usable) || (!old_v6 && as->v6_usable)) + need_reconnect = 1; + if (as->v4_addr_changed || as->v6_addr_changed) + need_reconnect = 1; + as->v4_addr_changed = 0; as->v6_addr_changed = 0; + chat_event_post(CHAT_EVT_LOCAL_SOCKETS, NULL, 0); - (void)as; + + if (need_reconnect) { + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] connectivity change: v4 %d→%d v6 %d→%d — triggering reconnect", + old_v4, as->v4_usable, old_v6, as->v6_usable); + chat_sync_retry_channels_on_socket_change(as->instance); + } } /* ═══════════ Platform-agnostic core ═══════════ */ +/* Генерирует случайный порт в диапазоне 1024..65535 для нового сокета. */ static uint16_t random_port(void) { uint32_t r; random_bytes((uint8_t*)&r, sizeof(r)); return (uint16_t)((r % (65535 - 1024)) + 1024); } +/* + * Создаёт UDP-сокет привязанный к конкретному интерфейсу и добавляет к нему + * линки ко всем активным ETCP-соединениям. + * + * Алгоритм: + * 1. Пытается переиспользовать сохранённый порт из БД (на случай перезапуска) + * 2. Если занят — генерирует случайный (до 10 попыток) + * 3. Вызывает etcp_socket_add — создаёт ETCP_SOCKET с bind на интерфейс + * 4. Обновляет interface_addr через socket_monitor_update_if_addr + * 5. Отбрасывает сокет если адрес нулевой (INADDR_ANY / ::) — интерфейс без реального IP + * 6. Сохраняет порт в БД для будущих перезапусков + * 7. Вызывает ncd_add_socket_links — добавляет по одному случайному линку + * к каждому активному соединению (для каждого address family) + * + * Возвращает 0 при успехе, -1 при ошибке. */ static int create_iface_udp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname, int family, uint8_t type) { struct UTUN_INSTANCE* inst = as->instance; struct CFG_SERVER server; @@ -266,6 +351,20 @@ static int create_iface_udp_socket(struct AUTO_SOCKET* as, uint32_t ifindex, con return -1; } +/* + * Создаёт TCP-сокет и STCP-сервер на интерфейсе, добавляет исходящие TCP-линки + * ко всем активным ETCP-соединениям. + * + * Алгоритм аналогичен create_iface_udp_socket: + * 1. Переиспользует сохранённый порт или генерирует случайный + * 2. Создаёт ETCP_SOCKET через tcp_socket_add + * 3. Отбрасывает если адрес нулевой + * 4. Создаёт stcp_server (слушает входящие STCP-подключения) + * 5. Регистрирует сервер через stcp_server_list_add + * 6. Вызывает ncd_add_socket_links — добавляет исходящие TCP-линки + * (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) { struct UTUN_INSTANCE* inst = as->instance; @@ -323,6 +422,11 @@ static struct stcp_server* create_iface_tcp_socket(struct AUTO_SOCKET* as, uint3 save_port_to_db(as, ifname, family, AS_PROTO_TCP, port); + { int tcp_links = ncd_add_socket_links(inst, ts); + if (tcp_links > 0) + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] added %d TCP links on if=%s", tcp_links, ifname); + } + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] TCP socket created: if=%s(%u) type=%d port=%u", ifname, ifindex, type, port); return tsrv; } @@ -386,6 +490,14 @@ void auto_socket_remove_interface(struct UTUN_INSTANCE* inst, uint32_t ifindex) remove_iface_sockets(as, ifindex); } +/* + * Удаляет все сокеты и линки для указанного интерфейса. + * Для каждого из 4 сокетов (v4 UDP/TCP, v6 UDP/TCP): + * 1. Закрывает ETCP-линки через ncd_remove_socket_links + * 2. Удаляет ETCP_SOCKET (etcp_socket_remove / tcp_socket_remove) + * 3. Для TCP — дополнительно останавливает stcp_server + * 4. Удаляет запись порта из БД + * Затем освобождает запись auto_sock_iface и оповещает GUI. */ static void remove_iface_sockets(struct AUTO_SOCKET* as, uint32_t ifindex) { struct auto_sock_iface** pp = &as->ifaces; while (*pp) { @@ -450,12 +562,19 @@ static void remove_iface_sockets(struct AUTO_SOCKET* as, uint32_t ifindex) { } } -/* ─── Сканирование адресов интерфейса и определение типа ─── */ +/* ─── Сканирование адресов интерфейса и определение типа ─── + * + * Эти функции анализируют реальное состояние интерфейса: есть ли v4/v6 адреса, + * являются ли они публичными или приватными. Результат используется reconcile_iface + * для принятия решений о создании/удалении/пересоздании сокетов. */ #if !defined(_WIN32) #include +/* Проверяет через /proc/net/if_inet6: есть ли на интерфейсе постоянные (не temporary) + * IPv6-адреса. Временные адреса (privacy extensions) пропускаются — они нестабильны + * и не подходят для долгоживущих P2P-соединений. */ /* /proc/net/if_inet6 format: hex_addr(32) ifindex(2) prefixlen(2) scope(2) flags(2) ifname flags bit 0x20 = IFA_F_TEMPORARY */ static int iface_has_permanent_v6(const char* ifname) { @@ -476,6 +595,8 @@ static int iface_has_permanent_v6(const char* ifname) { return found; } +/* Проверяет является ли конкретный IPv6-адрес временным (IFA_F_TEMPORARY, + * privacy extensions). Такие адреса меняются и не годятся для P2P. */ static int is_v6_addr_temporary(const uint8_t addr[16], const char* ifname) { char hex[33]; for (int i = 0; i < 16; i++) snprintf(hex + i * 2, 3, "%02x", addr[i]); @@ -495,6 +616,13 @@ static int is_v6_addr_temporary(const uint8_t addr[16], const char* ifname) { return tmp; } +/* + * Сканирует все IP-адреса на интерфейсе и определяет: + * - есть ли IPv4/IPv6 адреса (out_has_v4, out_has_v6) + * - тип адресов: PUBLIC (глобальный) или NAT (приватный/link-local) + * Для IPv6 временные адреса (privacy extensions) пропускаются. + * Если интерфейс не UP — все out-параметры сбрасываются в 0/NAT. + * Используется reconcile_iface для сравнения ожидаемого и фактического состояния. */ 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) { @@ -539,6 +667,10 @@ static int scan_and_classify_iface(uint32_t ifindex, const char* ifname, return 0; } +/* + * Полное сканирование: обходит все UP-интерфейсы (не loopback) и для каждого + * вызывает reconcile_iface. Используется при старте и при ручном запросе + * пересканирования (auto_socket_on_network_change). */ static void auto_socket_scan_all(struct AUTO_SOCKET* as) { struct ifaddrs* ifa_list = NULL; if (getifaddrs(&ifa_list) < 0) { DEBUG_ERROR(DEBUG_CATEGORY_AS, "[as] getifaddrs: %s", strerror(errno)); return; } @@ -585,9 +717,16 @@ static void auto_socket_scan_all(struct AUTO_SOCKET* as) { #endif /* _WIN32 */ -/* ═══════════ reconcile: сверка и исправление состояния сокетов ═══════════ */ - -/* helper: ensure ifa record exists, allocate if needed */ +/* ═══════════ reconcile: сверка и исправление состояния сокетов ═══════════ + * + * reconcile — ключевой механизм модуля. При любом изменении сети сравнивает + * фактическое состояние интерфейса (IP-адреса, их тип) с текущим набором сокетов + * и устраняет расхождения: создаёт недостающие сокеты, удаляет лишние, + * пересоздаёт при смене типа (NAT↔PUBLIC), обновляет interface_addr. */ + +/* Находит или создаёт запись 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) { struct auto_sock_iface* ifa = as->ifaces; while (ifa) { if (ifa->netif_index == ifindex) return ifa; ifa = ifa->next; } @@ -599,7 +738,8 @@ static struct auto_sock_iface* reconcile_get_ifa(struct AUTO_SOCKET* as, uint32_ return ifa; } -/* helper: remove ifa record if all 4 sockets gone */ +/* Удаляет запись auto_sock_iface если все 4 сокета (v4/v6 UDP+TCP) отсутствуют. + * Вызывается после reconcile чтобы подчистить записи для интерфейсов без сокетов. */ static void reconcile_prune_ifa(struct AUTO_SOCKET* as, struct auto_sock_iface* ifa, uint32_t ifindex, const char* ifname) { if (!ifa) return; if (ifa->v4_udp || ifa->v6_udp || ifa->v4_tcp || ifa->v6_tcp) return; @@ -610,6 +750,19 @@ static void reconcile_prune_ifa(struct AUTO_SOCKET* as, struct auto_sock_iface* u_free(ifa); } +/* + * Сверяет состояние интерфейса с набором сокетов и устраняет расхождения. + * Это основная логика модуля — вызывается при ЛЮБОМ изменении сети. + * + * Для каждого из 4 сокетов (v4/v6 UDP/TCP) проверяет 4 сценария: + * 1. Адрес есть а сокета нет → создать сокет + линки + * 2. Адреса нет а сокет есть → удалить сокет + линки + * 3. Тип изменился (NAT↔PUBLIC) → пересоздать сокет с новым типом + * 4. Адрес изменился (тот же семейство/тип, другой IP) → обновить + * interface_addr, оповестить через ETCP_SOCKET_EVENT_ADDR_CHANGED + * + * В конце удаляет запись auto_sock_iface если все сокеты отсутствуют, + * и оповещает GUI если были изменения. */ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char* ifname) { if (ifindex == 0) return; int changed = 0; @@ -639,22 +792,28 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char ncd_remove_socket_links(as->instance, sock); etcp_socket_remove(sock); ifa->v4_udp = NULL; changed = 1; - } else if (has_v4 && sock && cur_type != v4_need_type) { - DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v4 UDP type changed: %s type %d→%d", ifname, cur_type, v4_need_type); - ncd_remove_socket_links(as->instance, sock); - etcp_socket_remove(sock); - 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 = s; break; } s = s->next; } } - changed = 1; } else if (has_v4 && sock) { struct sockaddr_storage old4 = sock->interface_addr; socket_monitor_update_if_addr(sock); - if (memcmp(&old4, &sock->interface_addr, sizeof(old4)) != 0) { - DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v4 UDP (%s) addr changed: %s → %s", - ifname, sockaddr_storage_to_str(&old4).str, sockaddr_storage_to_str(&sock->interface_addr).str); - etcp_socket_cbk_fire(sock, ETCP_SOCKET_EVENT_ADDR_CHANGED); + int addr_changed = (memcmp(&old4, &sock->interface_addr, sizeof(old4)) != 0); + int type_changed = (cur_type != v4_need_type); + if (addr_changed || type_changed) { + uint8_t new_type = type_changed ? v4_need_type : cur_type; + if (addr_changed && type_changed) + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v4 UDP (%s) addr+type changed: %s→%s type %d→%d", ifname, + sockaddr_storage_to_str(&old4).str, sockaddr_storage_to_str(&sock->interface_addr).str, cur_type, new_type); + else if (type_changed) + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v4 UDP (%s) type changed: %d→%d", ifname, cur_type, new_type); + else + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v4 UDP (%s) addr changed: %s→%s", ifname, + sockaddr_storage_to_str(&old4).str, sockaddr_storage_to_str(&sock->interface_addr).str); + ncd_remove_socket_links(as->instance, sock); + etcp_socket_remove(sock); + create_iface_udp_socket(as, ifindex, ifname, AF_INET, new_type); + 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; } } + as->v4_addr_changed = 1; changed = 1; } } @@ -678,22 +837,28 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char ncd_remove_socket_links(as->instance, sock); etcp_socket_remove(sock); ifa->v6_udp = NULL; changed = 1; - } else if (has_v6 && sock && cur_type != v6_need_type) { - DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v6 UDP type changed: %s type %d→%d", ifname, cur_type, v6_need_type); - ncd_remove_socket_links(as->instance, sock); - etcp_socket_remove(sock); - 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 = s; break; } s = s->next; } } - changed = 1; } else if (has_v6 && sock) { struct sockaddr_storage old6 = sock->interface_addr; socket_monitor_update_if_addr(sock); - if (memcmp(&old6, &sock->interface_addr, sizeof(old6)) != 0) { - DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v6 UDP (%s) addr changed: %s → %s", - ifname, sockaddr_storage_to_str(&old6).str, sockaddr_storage_to_str(&sock->interface_addr).str); - etcp_socket_cbk_fire(sock, ETCP_SOCKET_EVENT_ADDR_CHANGED); + int addr_changed = (memcmp(&old6, &sock->interface_addr, sizeof(old6)) != 0); + int type_changed = (cur_type != v6_need_type); + if (addr_changed || type_changed) { + uint8_t new_type = type_changed ? v6_need_type : cur_type; + if (addr_changed && type_changed) + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v6 UDP (%s) addr+type changed: %s→%s type %d→%d", ifname, + sockaddr_storage_to_str(&old6).str, sockaddr_storage_to_str(&sock->interface_addr).str, cur_type, new_type); + else if (type_changed) + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v6 UDP (%s) type changed: %d→%d", ifname, cur_type, new_type); + else + DEBUG_INFO(DEBUG_CATEGORY_AS, "[as] v6 UDP (%s) addr changed: %s→%s", ifname, + sockaddr_storage_to_str(&old6).str, sockaddr_storage_to_str(&sock->interface_addr).str); + ncd_remove_socket_links(as->instance, sock); + etcp_socket_remove(sock); + create_iface_udp_socket(as, ifindex, ifname, AF_INET6, new_type); + 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; } } + as->v6_addr_changed = 1; changed = 1; } } @@ -754,6 +919,7 @@ static void reconcile_iface(struct AUTO_SOCKET* as, uint32_t ifindex, const char /* ─── Linux: netlink ─── */ #ifdef __linux__ +/* При изменении IP-адреса (RTM_NEWADDR / RTM_DELADDR) — запускаем reconcile для интерфейса. */ static void handle_nl_addr(struct AUTO_SOCKET* as, struct ifaddrmsg* ifa, int msg_type) { (void)msg_type; char ifname[IFNAMSIZ]; @@ -761,6 +927,9 @@ static void handle_nl_addr(struct AUTO_SOCKET* as, struct ifaddrmsg* ifa, int ms reconcile_iface(as, ifa->ifa_index, ifname); } +/* При изменении состояния линка (RTM_NEWLINK / RTM_DELLINK): + * - DOWN или удаление → удаляем все сокеты интерфейса + * - UP (новый или поднялся) → запускаем reconcile */ static void handle_nl_link(struct AUTO_SOCKET* as, struct ifinfomsg* ifi, int msg_type) { char ifname[IFNAMSIZ]; if (!if_indextoname(ifi->ifi_index, ifname)) return; @@ -776,6 +945,8 @@ static void handle_nl_link(struct AUTO_SOCKET* as, struct ifinfomsg* ifi, int ms } } +/* Коллбэк uasync для netlink-сокета. Принимает все сообщения netlink, + * разбирает их и направляет в handle_nl_addr / handle_nl_link. */ static void auto_socket_netlink_cb(int fd, void* arg) { struct AUTO_SOCKET* as = (struct AUTO_SOCKET*)arg; char buf[8192]; @@ -822,12 +993,15 @@ static int auto_socket_init_monitor(struct AUTO_SOCKET* as) { return 0; } +/* Отключает netlink-мониторинг: убирает fd из uasync, закрывает сокет. */ static void auto_socket_destroy_monitor(struct AUTO_SOCKET* as) { if (as->uasync_handle) { uasync_remove_socket(as->instance->ua, as->uasync_handle); as->uasync_handle = NULL; } if (as->nl_sock >= 0) { close(as->nl_sock); as->nl_sock = -1; } } -/* ─── BSD: route socket ─── */ +/* ─── BSD: route socket ─── + * На FreeBSD мониторим изменения через PF_ROUTE сокет. При любых изменениях + * адресов/интерфейсов запускаем полное пересканирование. */ #elif defined(__FreeBSD__) || defined(__FreeBSD_kernel__) static void auto_socket_bsd_cb(int fd, void* arg) { @@ -859,7 +1033,9 @@ static void auto_socket_destroy_monitor(struct AUTO_SOCKET* as) { if (as->route_sock >= 0) { close(as->route_sock); as->route_sock = -1; } } -/* ─── Windows: IP Helper ─── */ +/* ─── Windows: IP Helper ─── + * Используем NotifyUnicastIpAddressChange + NotifyRouteChange2. + * Коллбэки приходят из системных потоков → uasync_post для обработки в главном потоке. */ #elif defined(_WIN32) static void process_win_notify(void* arg) { @@ -894,12 +1070,19 @@ static void auto_socket_destroy_monitor(struct AUTO_SOCKET* as) { } #else -static int auto_socket_init_monitor(struct AUTO_SOCKET* as) { (void)as; return 0; } +/* Платформа без встроенного мониторинга — только явные вызовы API (JNI и т.д.). */ +static int auto_socket_init_monitor(struct AUTO_SOCKET* as) { (void)as; return 0; } static void auto_socket_destroy_monitor(struct AUTO_SOCKET* as) { (void)as; } #endif /* ═══════════ Public API ═══════════ */ +/* + * Инициализирует модуль auto_socket. Если auto_sockets=no в конфиге — сразу + * возвращает 0 без действий. Иначе создаёт таблицу портов в БД, запускает + * платформенный мониторинг и сканирует все текущие UP-интерфейсы. + * Состояние хранится в inst->auto_socket_state. + */ int auto_socket_init(struct UTUN_INSTANCE* inst) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_AS, "[as] instance is NULL"); return -1; } if (!inst->config || !inst->config->global.auto_sockets) return 0; @@ -930,6 +1113,9 @@ int auto_socket_init(struct UTUN_INSTANCE* inst) { return 0; } +/* + * Завершает работу модуля: останавливает платформенный мониторинг, удаляет + * все сокеты и линки, освобождает память. Безопасно вызывать повторно. */ void auto_socket_destroy(struct UTUN_INSTANCE* inst) { if (!inst || !inst->auto_socket_state) return; struct AUTO_SOCKET* as = (struct AUTO_SOCKET*)inst->auto_socket_state; diff --git a/src/transport_layer/auto_socket.h b/src/transport_layer/auto_socket.h index d4a6e8b2..cc8c60ff 100644 --- a/src/transport_layer/auto_socket.h +++ b/src/transport_layer/auto_socket.h @@ -6,20 +6,32 @@ extern "C" { #endif /* - * auto_socket.h — автоматическое управление UDP-сокетами при изменении сетевых интерфейсов + * auto_socket.h — автоматическое управление сокетами при изменении сетевых интерфейсов * - * При включённой опции auto_sockets=yes: - * - Ручные [server] сокеты пропускаются - * - Для каждого не-loopback интерфейса (UP, с IP-адресами) создаётся UDP-сокет (v4/v6) - * с bind на 0.0.0.0/[::] + случайный порт, привязанный к интерфейсу - * - Тип сокета: PUBLIC если есть публичный адрес, иначе NAT - * - Для каждого нового сокета добавляются линки ко всем активным NCD-соединениям - * - При пропадании интерфейса линки удаляются, сокет закрывается + * === Зачем нужен модуль === + * При auto_sockets=yes не нужно вручную прописывать секции [server] в конфиге. + * Модуль сам создаёт UDP+TCP сокеты на всех не-loopback интерфейсах с IP-адресами, + * и сразу добавляет линки ко всем активным ETCP-соединениям. При пропадании + * интерфейса линки удаляются, сокеты закрываются. * - * API разделён на два слоя: - * 1. auto_socket_init / auto_socket_destroy — жизненный цикл модуля - * 2. auto_socket_add_interface / auto_socket_remove_interface — platform-agnostic - * вызываются из платформенного мониторинга (netlink, Android JNI, etc.) + * === Что делает === + * - Для каждого интерфейса (UP, с v4/v6 адресами) создаёт пару сокетов: + * UDP (случайный порт) + TCP (stcp_server на случайном порту) + * - Тип сокета: PUBLIC если на интерфейсе есть глобальный адрес, иначе NAT + * - Порт сохраняется в SQLite для переиспользования при перезапуске + * - Для каждого нового сокета: etcp_link_new ко всем соединениям (inst->connections), + * случайно выбирая один совместимый адрес пира (отдельно для v4 и v6) + * - При изменении адреса на интерфейсе — обновляет interface_addr и оповещает + * подписчиков (ETCP_SOCKET_EVENT_ADDR_CHANGED) + * - При смене типа (NAT↔PUBLIC) — пересоздаёт сокет с новым типом + * - При пропадании интерфейса — удаляет линки, закрывает сокет, чистит БД + * + * === Слои === + * 1. Platform-agnostic API: auto_socket_init / auto_socket_destroy + * auto_socket_add_interface / auto_socket_remove_interface + * (можно вызывать из любого платформенного мониторинга: netlink, JNI и т.д.) + * 2. Platform monitoring: встроенные обработчики netlink (Linux), + * route socket (BSD), IP Helper (Windows) — работают автоматически */ #include @@ -27,33 +39,46 @@ extern "C" { struct UTUN_INSTANCE; struct ETCP_SOCKET; +/* Внутреннее состояние модуля — хранится в inst->auto_socket_state. + * Содержит связный список отслеживаемых интерфейсов (ifaces) и платформенные + * handle'ы для мониторинга (netlink fd, route socket fd, IP Helper handles). */ +struct AUTO_SOCKET; + /* ─── Жизненный цикл ─── */ +/* Инициализирует модуль: создаёт таблицу портов в БД, запускает платформенный + * мониторинг интерфейсов, сканирует все текущие UP-интерфейсы и создаёт сокеты. + * Если auto_sockets=no в конфиге — ничего не делает, возвращает 0. */ int auto_socket_init(struct UTUN_INSTANCE* inst); + +/* Останавливает мониторинг, удаляет все созданные сокеты и линки, освобождает память. */ void auto_socket_destroy(struct UTUN_INSTANCE* inst); -/* Platform notification: call from any thread when network interfaces change. - * Posts auto_socket_scan_all() to uasync thread. Safe for JNI/ConnectivityManager callbacks. */ +/* Вызвать при изменении сетевых интерфейсов из любого потока (JNI, ConnectivityManager). + * Откладывает полное сканирование (auto_socket_scan_all) в uasync-поток. */ void auto_socket_on_network_change(struct UTUN_INSTANCE* inst); /* ─── Platform-agnostic API: управление сокетами на конкретном интерфейсе ─── */ /* - * Создать сокет(ы) для интерфейса и добавить линки ко всем активным соединениям. - * ifindex — индекс интерфейса (0 = без привязки, использовать только ifname) + * Создать UDP+TCP сокеты для интерфейса и добавить линки ко всем активным соединениям. + * ifindex — индекс интерфейса (0 = автоопределение по ifname) * ifname — имя интерфейса для SO_BINDTODEVICE - * has_v4 — есть IPv4 адреса на интерфейсе - * has_v6 — есть IPv6 адреса на интерфейсе + * has_v4 — есть ли IPv4 адреса на интерфейсе + * has_v6 — есть ли IPv6 адреса на интерфейсе * v4_type — CFG_SERVER_TYPE_PUBLIC или CFG_SERVER_TYPE_NAT * v6_type — CFG_SERVER_TYPE_PUBLIC или CFG_SERVER_TYPE_NAT - * Возвращает количество созданных сокетов (0, 1 или 2) + * Возвращает количество созданных сокетов (0, 1, 2, 3 или 4 — до 4: v4 UDP+TCP, v6 UDP+TCP). + * Если сокеты для этого интерфейса уже есть — ничего не делает, возвращает 0. */ int auto_socket_add_interface(struct UTUN_INSTANCE* inst, uint32_t ifindex, const char* ifname, int has_v4, int has_v6, uint8_t v4_type, uint8_t v6_type); /* - * Удалить сокет(ы) для интерфейса: закрыть линки, удалить сокет. + * Удалить все сокеты и линки для интерфейса. + * Закрывает UDP/TCP сокеты, удаляет etcp-линки через ncd_remove_socket_links, + * чистит запись порта в БД. * ifindex — индекс интерфейса, сокеты которого удаляем */ void auto_socket_remove_interface(struct UTUN_INSTANCE* inst, uint32_t ifindex); diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 02a692a4..eb966287 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -453,19 +453,19 @@ restart_timer: link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive"); } -static void sockaddr_to_key(struct sockaddr_storage* addr, uint8_t key[LINK_ADDR_KEY_SIZE]) { +static void sockaddr_to_key(struct sockaddr_storage* addr, uint8_t key[LINK_ADDR_KEY_SIZE], int is_tcp) { memset(key, 0, LINK_ADDR_KEY_SIZE); if (!addr) return; if (addr->ss_family == AF_INET) { struct sockaddr_in* sa = (struct sockaddr_in*)addr; memcpy(key, &sa->sin_port, 2); memcpy(key + 2, &sa->sin_addr.s_addr, 4); - key[18] = 4; + key[18] = 4 | (is_tcp ? 0x80 : 0); } else { struct sockaddr_in6* sa = (struct sockaddr_in6*)addr; memcpy(key, &sa->sin6_port, 2); memcpy(key + 2, &sa->sin6_addr, 16); - key[18] = 6; + key[18] = 6 | (is_tcp ? 0x80 : 0); } } @@ -476,7 +476,7 @@ static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) if (!e_sock || !link || !e_sock->links_queue) return -1; uint8_t key[LINK_ADDR_KEY_SIZE]; - sockaddr_to_key(&link->remote_addr, key); + sockaddr_to_key(&link->remote_addr, key, link->is_tcp); struct ll_entry* dup_qe = queue_find_data_by_index(e_sock->links_queue, key); if (dup_qe) { struct link_queue_entry* dup_lqe = (struct link_queue_entry*)dup_qe->data; @@ -529,11 +529,11 @@ static int sockaddr_equal(const struct sockaddr_storage* a, const struct sockadd return 0; } -struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr) { +struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr, int is_tcp) { DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); if (!e_sock || !addr || !e_sock->links_queue) return NULL; uint8_t key[LINK_ADDR_KEY_SIZE]; - sockaddr_to_key(addr, key); + sockaddr_to_key(addr, key, is_tcp); struct ll_entry* e = queue_find_data_by_index(e_sock->links_queue, key); if (!e) return NULL; struct link_queue_entry* lqe = (struct link_queue_entry*)e->data; @@ -1038,6 +1038,9 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn if (remote_addr) memcpy(&link->remote_addr, remote_addr, sizeof(struct sockaddr_storage)); + if (conn && conn->interface_addr.ss_family) memcpy(&link->local_bound_addr, &conn->interface_addr, sizeof(struct sockaddr_storage)); + else memset(&link->local_bound_addr, 0, sizeof(link->local_bound_addr)); + link->last_recv_local_time = get_time_tb(); // Initialize to prevent immediate timeout link->total_retransmissions = 0; @@ -1151,6 +1154,7 @@ void etcp_link_close(struct ETCP_LINK* link) { etcp_fire_link_status_cbk(link, link->link_state, link->link_status); etcp_inflight_nullify_link(link->etcp->input_wait_ack, link); etcp_inflight_nullify_link(link->etcp->input_send_q, link); + etcp_on_link_down(link->etcp, link); u_free(link->bbr); u_free(link); } @@ -1974,8 +1978,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { size_t pkt_len=0; int errorcode=0; - struct ETCP_LINK* link=etcp_link_find_by_addr(e_sock, &addr); - + struct ETCP_LINK* link=etcp_link_find_by_addr(e_sock, &addr, 0); // Try normal decryption first if we have an established link with session keys // This is the common case for data packets and responses // if (link) { @@ -2141,8 +2144,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; conn = ce->conn; } } if (!conn) { - struct ETCP_LINK* ol = etcp_link_find_by_addr(e_sock, &addr); - if (ol && ol->etcp && ol->etcp->conn_queue_entry) { + struct ETCP_LINK* ol = etcp_link_find_by_addr(e_sock, &addr, 0); if (ol && ol->etcp && ol->etcp->conn_queue_entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)ol->etcp->conn_queue_entry->data; if (ce->peer_node_id == 0) { conn = ol->etcp; @@ -2201,7 +2203,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { struct ETCP_LINK* existing_link = etcp_link_find_by_remote_id(conn, req->link_id); if (!existing_link) { - existing_link = etcp_link_find_by_addr(e_sock, &addr); + existing_link = etcp_link_find_by_addr(e_sock, &addr, 0); if (existing_link && existing_link->etcp == conn) { if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] link address match but pubkey mismatch, firing node_changed", conn->log_name); diff --git a/src/transport_layer/etcp_connections.h b/src/transport_layer/etcp_connections.h index 80f46456..38da781d 100644 --- a/src/transport_layer/etcp_connections.h +++ b/src/transport_layer/etcp_connections.h @@ -318,6 +318,7 @@ struct ETCP_LINK { etcp_udp_send_fn_t send_hook; // NULL = socket_sendto напрямую void* send_hook_ctx; + struct sockaddr_storage local_bound_addr; // локальный адрес к которому привязан линк (interface_addr на момент создания) uint8_t is_tcp; // 1 = TCP транспорт (STCP) struct stcp_link *tcp_link; // STCP линк (для TCP) void *tcp_reconnect_timer; // таймер реконнекта (TCP клиент) @@ -355,7 +356,7 @@ void etcp_link_close(struct ETCP_LINK* link); //int etcp_input_cbk(struct packet_buffer* pkt, struct ETCP_SOCKET* conn);// получает расшифрованный пакет int etcp_encrypt_send(struct ETCP_DGRAM* dgram);// зашифровывает и отправляет пакет // find link by address -struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr); +struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr, int is_tcp); // find free local_link_id for connection // scans all links in connection, marks used ids in bit array diff --git a/src/transport_layer/node_conn_direct.c b/src/transport_layer/node_conn_direct.c index adaf7048..3f118653 100644 --- a/src/transport_layer/node_conn_direct.c +++ b/src/transport_layer/node_conn_direct.c @@ -157,7 +157,7 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni, struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); struct ETCP_SOCKET* use_sock = socks[rr++ % sock_count]; { - struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa); + struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa, 0); if (stale && stale->etcp == conn) continue; if (stale && stale->etcp != conn && memcmp(conn->crypto_ctx.peer_public_key, stale->etcp->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE)) @@ -209,7 +209,7 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni, if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = use_sock->netif_index; struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6)); { - struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa); + struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa, 0); if (stale && stale->etcp == conn) continue; if (stale && stale->etcp != conn && memcmp(conn->crypto_ctx.peer_public_key, stale->etcp->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE)) @@ -965,23 +965,132 @@ struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h) { /* ═══════════ Управление линками на конкретном сокете (для auto_socket) ═══════════ */ +/* ncd_add_udp_link_random / ncd_add_tcp_link_random — собирают совместимые адреса пира, + * случайно выбирают один, проверяют что линк на (sock, addr, proto) ещё не существует, + * и создают новый. */ +static int ncd_add_udp_link_random(struct ETCP_CONN* conn, struct TOPO_NODE* ni, struct ETCP_SOCKET* sock) { + int family = sock->local_addr.ss_family; + if (family == AF_INET) { + struct TOPO_ADDR4* compat[64]; int count = 0; + for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) { + if (!a->port || !(a->protocol & TOPO_PROTO_UDP)) continue; + int zero = 1; for (int j = 0; j < 4; j++) if (a->addr[j]) { zero = 0; break; } + if (zero) continue; + if (count < 64) compat[count++] = (struct TOPO_ADDR4*)a; + } + if (count == 0) return 0; + uint32_t r; random_bytes((uint8_t*)&r, sizeof(r)); + struct TOPO_ADDR4* a = compat[r % count]; + struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port); + struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); + if (etcp_link_find_by_addr(sock, &sa, 0)) return 0; + if (etcp_link_new(conn, sock, &sa, 0)) return 1; + return 0; + } + if (family == AF_INET6) { + struct TOPO_ADDR6* compat[64]; int count = 0; + for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) { + if (!a->port || !(a->protocol & TOPO_PROTO_UDP)) continue; + int zero = 1; for (int j = 0; j < 16; j++) if (a->addr[j]) { zero = 0; break; } + if (zero) continue; + if (count < 64) compat[count++] = (struct TOPO_ADDR6*)a; + } + if (count == 0) return 0; + uint32_t r; random_bytes((uint8_t*)&r, sizeof(r)); + struct TOPO_ADDR6* a = compat[r % count]; + struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6; + memcpy(&sin6.sin6_addr, a->addr, 16); sin6.sin6_port = htons(a->port); + if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = sock->netif_index; + struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6)); + if (etcp_link_find_by_addr(sock, &sa, 0)) return 0; + if (etcp_link_new(conn, sock, &sa, 0)) return 1; + return 0; + } + return 0; +} + +static int ncd_add_tcp_link_random(struct ETCP_CONN* conn, struct TOPO_NODE* ni, struct ETCP_SOCKET* sock) { + int family = sock->local_addr.ss_family; + if (family == AF_INET) { + struct TOPO_ADDR4* compat[64]; int count = 0; + for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) { + if (!a->port || !(a->protocol & TOPO_PROTO_TCP)) continue; + int zero = 1; for (int j = 0; j < 4; j++) if (a->addr[j]) { zero = 0; break; } + if (zero) continue; + if (count < 64) compat[count++] = (struct TOPO_ADDR4*)a; + } + if (count == 0) return 0; + uint32_t r; random_bytes((uint8_t*)&r, sizeof(r)); + struct TOPO_ADDR4* a = compat[r % count]; + struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port); + struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); + if (etcp_link_find_by_addr(sock, &sa, 1)) return 0; + struct ETCP_LINK* tlink = etcp_link_new(conn, sock, &sa, 0); + if (!tlink) return 0; + tlink->is_tcp = 1; + etcp_tcp_link_start_connect(tlink, &sa, a->port); + return 1; + } + if (family == AF_INET6) { + struct TOPO_ADDR6* compat[64]; int count = 0; + for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) { + if (!a->port || !(a->protocol & TOPO_PROTO_TCP)) continue; + int zero = 1; for (int j = 0; j < 16; j++) if (a->addr[j]) { zero = 0; break; } + if (zero) continue; + if (count < 64) compat[count++] = (struct TOPO_ADDR6*)a; + } + if (count == 0) return 0; + uint32_t r; random_bytes((uint8_t*)&r, sizeof(r)); + struct TOPO_ADDR6* a = compat[r % count]; + struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6; + memcpy(&sin6.sin6_addr, a->addr, 16); sin6.sin6_port = htons(a->port); + if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = sock->netif_index; + struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6)); + if (etcp_link_find_by_addr(sock, &sa, 1)) return 0; + struct ETCP_LINK* tlink = etcp_link_new(conn, sock, &sa, 0); + if (!tlink) return 0; + tlink->is_tcp = 1; + etcp_tcp_link_start_connect(tlink, &sa, a->port); + return 1; + } + return 0; +} + int ncd_add_socket_links(struct UTUN_INSTANCE* inst, struct ETCP_SOCKET* sock) { if (!inst || !sock) return 0; - int added = 0; - struct ncd_entry* entry = (struct ncd_entry*)inst->ncd_registry; - while (entry) { - struct TOPO_NODE* ni = ncd_lookup_node(inst, entry->node_id); - if (ni) { - int n = ncd_create_links(entry, ni, sock); - if (n > 0) { - DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] add_socket_links socket=%s node=0x%016llx links=%d", - sock->name, (unsigned long long)entry->node_id, n); - added += n; - } - topo_node_registry_unref(inst->topo_groups, ni->node_id); + int added = 0, total_conns = 0, skipped_closed = 0, skipped_no_addr = 0; + + struct ll_entry* qe = inst->connections ? inst->connections->head : NULL; + while (qe) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; + struct ETCP_CONN* conn = ce->conn; + qe = qe->next; + if (!conn) continue; + total_conns++; + + if (conn->state == 2) { skipped_closed++; continue; } + + struct TOPO_NODE* ni = ncd_lookup_node(inst, ce->peer_node_id); + if (!ni) { skipped_no_addr++; continue; } + + int n; + if (sock->is_tcp) n = ncd_add_tcp_link_random(conn, ni, sock); + else n = ncd_add_udp_link_random(conn, ni, sock); + if (n > 0) { + DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] add_socket_links sock=%s fd=%d conn=[%s] node=0x%016llx addr=%s is_tcp=%d", + sock->name, (int)sock->fd, conn->log_name, + (unsigned long long)ce->peer_node_id, + sockaddr_storage_to_str(&sock->local_addr).str, sock->is_tcp); + added += n; } - entry = entry->next; + topo_node_registry_unref(inst->topo_groups, ni->node_id); } + + DEBUG_INFO(DEBUG_CATEGORY_NCD, + "[ncd] add_socket_links done: sock=%s is_tcp=%d total_conns=%d skipped(cl=%d,no_addr=%d) added=%d", + sock->name, sock->is_tcp, total_conns, skipped_closed, skipped_no_addr, added); return added; } diff --git a/tests/Makefile.am b/tests/Makefile.am index bce44f6e..6816eb7c 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -71,6 +71,7 @@ check_PROGRAMS = \ test_media_delivery_integration \ test_media_delivery_full \ test_etcp_link_stress \ + test_auto_socket_dynamic \ bench_timeout_heap \ bench_uasync_timeouts @@ -366,6 +367,10 @@ test_etcp_link_stress_SOURCES = test_etcp_link_stress.c test_etcp_link_stress_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_etcp_link_stress_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_auto_socket_dynamic_SOURCES = test_auto_socket_dynamic.c +test_auto_socket_dynamic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/lib +test_auto_socket_dynamic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_media_delivery_full_SOURCES = test_media_delivery_full.c test_media_delivery_full_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/chat test_media_delivery_full_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) @@ -408,6 +413,11 @@ check-local: $(check_PROGRAMS) elapsed_ms=$$(( (end - start) / 1000000 )); \ printf "$${G}[PASS]$${N} %s (%sms)\n" "$$test" "$$elapsed_ms"; \ passed=$$((passed + 1)); \ + elif test $$? -eq 77; then \ + end=$$(date +%s%N); \ + elapsed_ms=$$(( (end - start) / 1000000 )); \ + printf "$${Y}[SKIP]$${N} %s (requires root)\n" "$$test"; \ + skipped=$$((skipped + 1)); \ else \ end=$$(date +%s%N); \ elapsed_ms=$$(( (end - start) / 1000000 )); \ diff --git a/tests/test_auto_socket_dynamic.c b/tests/test_auto_socket_dynamic.c new file mode 100644 index 00000000..a6354268 --- /dev/null +++ b/tests/test_auto_socket_dynamic.c @@ -0,0 +1,474 @@ +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifndef _WIN32 +#include +#include +#endif + +#include "etcp.h" +#include "etcp_connections.h" +#include "etcp_api.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/utun_instance.h" +#include "topo_group.h" +#include "topo_node.h" +#include "secure_channel.h" +#include "../src/transport_layer/node_conn_direct.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define TIMEOUT_TB 300000 +#define POLL_MS 5 +#define STEP_TB 3000 /* 300ms между фазами */ +#define TRAF_SEND_TB 50 /* 5ms — отправка */ +#define TRAF_MON_TB 1000 /* 100ms — мониторинг */ +#define ETCP_RT_ID_TEST 0xF0 + +typedef void (*timeout_cb)(void*); + +static struct test_ctx { + struct UTUN_INSTANCE *server, *client; + struct UASYNC* ua; + int round; /* 0..7 */ + int step; + int keep_tcp; + int ip_changes_on_last; + char cur_iface[IFNAMSIZ]; + char prev_iface[IFNAMSIZ]; + struct ETCP_CONN* srv_conn; + uint64_t srv_node_id; + uint8_t srv_pubkey[SC_PUBKEY_SIZE]; + int result; /* 0=running, 1=fail, 2=pass */ + + /* traffic */ + uint32_t send_seq, send_count, pong_count, total_recv; +} ctx; + +static void* timeout_handle; +static char tdir[] = "/tmp/utun_as_XXXXXX"; +static char scf[256], ccf[256]; + +/* ── rounds ── */ +static const struct { + const char* ifname, *ip1, *ip2; + int keep_tcp; +} rounds[] = { + {"dummy_cli1", "10.90.0.2/24", "10.90.1.2/24", 1}, + {"dummy_cli2", "10.90.0.3/24", "10.90.1.3/24", 1}, + {"dummy_cli3", "10.90.0.4/24", "10.90.1.4/24", 0}, + {"dummy_cli4", "10.90.0.5/24", "10.90.1.5/24", 0}, + {"dummy_cli5", "10.90.0.6/24", "10.90.1.6/24", 1}, + {"dummy_cli6", "10.90.0.7/24", "10.90.1.7/24", 1}, + {"dummy_cli7", "10.90.0.8/24", "10.90.1.8/24", 0}, + {"dummy_cli8", "10.90.0.9/24", "10.90.1.9/24", 0}, +}; +#define N_ROUNDS (int)(sizeof(rounds)/sizeof(rounds[0])) + +/* ── helpers ── */ +static int wf(const char* p, const char* f, ...) { + va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1; + va_start(ap, f); vfprintf(fp, f, ap); va_end(ap); fclose(fp); return 0; +} +static char* gv(const char* p, const char* k) { + struct utun_config* c = parse_config(p); if (!c) return NULL; + char* r = strcmp(k, "pub") == 0 ? u_strdup(c->global.my_public_key_hex) + : u_strdup(c->global.my_private_key_hex); + free_config(c); return r; +} +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 int count_links_to_srv(void) { + int n = 0; struct ll_entry* e = ctx.client->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + if (ce->conn->peer_node_id == ctx.srv_node_id) { + struct ETCP_LINK* l = ce->conn->links; + while (l) { if (l->initialized && l->link_status) n++; l = l->next; } + } e = e->next; } + return n; +} +static int all_links_are_type(void) { + struct ll_entry* e = ctx.client->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + if (ce->conn->peer_node_id == ctx.srv_node_id) { + struct ETCP_LINK* l = ce->conn->links; + while (l) { if (l->initialized && l->link_status && l->is_tcp != ctx.keep_tcp) return 0; l = l->next; } + } e = e->next; } + return 1; +} + +/* ── filter sockets ── */ +static void keep_only_socket_type(int keep_tcp) { + struct ETCP_SOCKET *s = ctx.client->etcp_sockets, *rm[16]; + int n = 0; + while (s) { if (s->local_addr.ss_family == AF_INET && s->is_tcp != keep_tcp && n < 16) rm[n++] = s; s = s->next; } + for (int i = 0; i < n; i++) { + ncd_remove_socket_links(ctx.client, rm[i]); + etcp_socket_remove(rm[i]); + } +} + +/* ── server node (2 addrs: UDP + TCP) ── */ +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)); + if (!nq || !ni) { u_free(nq); u_free(ni); return NULL; } + ni->group_ref_count = 1; ni->node_id = ctx.srv_node_id; ni->group_id = TOPO_GROUP_UTUN; + memcpy(ni->public_key, ctx.srv_pubkey, SC_PUBKEY_SIZE); + nq->ll.size = sizeof(struct TOPO_GROUP_NODE) - sizeof(struct ll_entry); + nq->node_id = ctx.srv_node_id; + + 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; } + 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; } + 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; + + topo_node_registry_store(ctx.client->topo_groups, ni); + return nq; +} + +/* ── traffic ── */ +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); +} +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; } + ctx.pong_count++; ctx.total_recv++; + queue_entry_free(entry); +} +static void traffic_send_timer(void* arg) { + (void)arg; + if (ctx.result) return; + if (ctx.srv_conn && ctx.srv_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(ctx.srv_conn, e); ctx.send_count++; } else queue_entry_free(e); } + } + if (!ctx.result) timeout_handle = uasync_set_timeout(ctx.ua, TRAF_SEND_TB, NULL, 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\n", + ctx.round, ctx.send_count, ctx.total_recv, d, d * 10); fflush(stderr); + if (!ctx.result) timeout_handle = uasync_set_timeout(ctx.ua, TRAF_MON_TB, NULL, traffic_monitor_timer, "traf_mon"); +} + +/* ── etcp_connect callback ── */ +static void connect_cb(void* arg, struct ETCP_CONN* conn, int type) { + (void)arg; + if (type == ETCP_CONNECT_EARLY || type == ETCP_CONNECT_LATE) + { if (conn) { ctx.srv_conn = conn; ctx.step = 9; } } +} + +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"); +} + +/* ── 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); +} +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); +} +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); +} +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); +} +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); +} + +/* ═══════════════════════════════════════════════════════════ + * Phases + * ═══════════════════════════════════════════════════════════ */ +static void phase_check_ip(void* arg); +static void phase_change_ip(void* arg); +static void phase_check_del(void* arg); +static void phase_del_prev(void* arg); +static void phase_check_add(void* arg); +static void phase_add(void* arg); +static void phase_del_last(void* arg); +static void phase_done(void* arg); + +static void phase_done(void* arg) { + (void)arg; + if (ctx.result) return; + fprintf(stderr, "=== ALL PASSED ===\n"); fflush(stderr); + ctx.result = 2; +} + +static void phase_del_last(void* arg) { + (void)arg; + if (ctx.result) return; + ctx.step = 11; + ip_link_del(rounds[N_ROUNDS-1].ifname); + uasync_set_timeout(ctx.ua, STEP_TB, NULL, (timeout_cb)phase_del_last, "chk_del_last"); + return; /* перепланируем один раз для проверки */ +} + +/* ══ check wrappers ══ */ +static void do_check(void) { + int l = count_links_to_srv(); + if (l < 1) { fail("no links"); return; } + if (!all_links_are_type()) { fail("wrong link type"); return; } +} + +static void phase_check_del_last(void* arg) { + (void)arg; + ctx.step = 12; + int l = count_links_to_srv(); + if (l > 0) { fail("links survived last del"); return; } + /* повторная проверка через 300ms — после второго захода считаем ОК */ + static int cnt = 0; + if (++cnt < 2) { uasync_set_timeout(ctx.ua, STEP_TB, NULL, (timeout_cb)phase_check_del_last, "chk_del_last"); return; } + fprintf(stderr, " r=%d s=%d: no links after del_last (OK)\n", ctx.round, ctx.step); fflush(stderr); + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_done, "done"); +} + +static void phase_check_ip(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 6; + do_check(); if (ctx.result) return; + + 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"); + } 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"); + } else { + ctx.round++; + uasync_set_timeout(ctx.ua, STEP_TB, NULL, 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; + ip_addr_del(rounds[ctx.round].ifname, old_ip); + ip_addr_add(rounds[ctx.round].ifname, new_ip); + keep_only_socket_type(ctx.keep_tcp); + uasync_set_timeout(ctx.ua, STEP_TB, NULL, 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; + fprintf(stderr, " r=%d s=%d: del_prev OK links=%d\n", ctx.round, ctx.step, count_links_to_srv()); + fflush(stderr); + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip"); +} + +static void phase_del_prev(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 3; + ip_link_del(ctx.prev_iface); + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_del, "chk_del"); +} + +static void phase_check_add(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 2; + /* wait for etcp_connect on round 0 */ + if (ctx.round == 0 && ctx.step == 2 && !ctx.srv_conn) { + uasync_set_timeout(ctx.ua, STEP_TB / 3, NULL, phase_check_add, "chk_add"); + return; + } + if (ctx.round == 0 && !ctx.srv_conn) { fail("etcp_connect didn't fire"); return; } + do_check(); if (ctx.result) return; + fprintf(stderr, " r=%d s=%d: add OK links=%d type=%s\n", + ctx.round, ctx.step, count_links_to_srv(), ctx.keep_tcp ? "TCP" : "UDP"); + fflush(stderr); + + if (ctx.round == 0) start_traffic(); + + if (ctx.prev_iface[0]) { + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_del_prev, "del_prev"); + } else { + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip"); + } +} + +static void phase_add(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 1; + ctx.keep_tcp = rounds[ctx.round].keep_tcp; + 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); + keep_only_socket_type(ctx.keep_tcp); + + 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); + } + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_add, "chk_add"); +} + +/* ═══════════════════════════════════════════════════════════ + * 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); } + atexit(cleanup_ifaces); + cleanup_ifaces(); + test_mkdtemp(tdir); + snprintf(scf, sizeof(scf), "%s/s.conf", tdir); + snprintf(ccf, sizeof(ccf), "%s/c.conf", tdir); + + /* create db directories */ + { char dbd[256]; snprintf(dbd, sizeof(dbd), "%s/db_srv", tdir); utun_mkdir(dbd, 0755); + snprintf(dbd, sizeof(dbd), "%s/db_cli", tdir); utun_mkdir(dbd, 0755); } + + /* server: fixed [server] on dummy_srv */ + ip_link_add("dummy_srv"); + ip_addr_add("dummy_srv", "10.90.0.1/24"); + ip_link_up("dummy_srv"); + + /* client: first interface */ + ip_link_add(rounds[0].ifname); + ip_addr_add(rounds[0].ifname, rounds[0].ip1); + ip_link_up(rounds[0].ifname); + + /* write initial configs */ + 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" + "[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); + config_ensure_keys_and_node_id(scf); + config_ensure_keys_and_node_id(ccf); + + { struct utun_config* cs = parse_config(scf); ctx.srv_node_id = cs->global.my_node_id; free_config(cs); } + { struct utun_config* cc = parse_config(ccf); (void)cc->global.my_node_id; free_config(cc); } + + char *spub = gv(scf, "pub"), *spriv = gv(scf, "priv"); + char *cpub = gv(ccf, "pub"), *cpriv = gv(ccf, "priv"); + + wf(scf, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun_srv\n" + "tun_test_mode=1\nauto_sockets=no\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, tdir); + wf(ccf, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun_cli\n" + "tun_test_mode=1\nauto_sockets=yes\ndb_path=%s/db_cli\n[allowed_keys]\nallow_all=1\n", + cpriv, cpub, tdir); + + /* verify config */ + { struct utun_config* ck = parse_config(ccf); + fprintf(stderr, " [setup] client: auto_sockets=%d node_id=0x%016llx\n", + ck ? ck->global.auto_sockets : -1, + (unsigned long long)(ck ? ck->global.my_node_id : 0)); + free_config(ck); + ck = parse_config(scf); + fprintf(stderr, " [setup] server: auto_sockets=%d node_id=0x%016llx\n", + ck ? ck->global.auto_sockets : -1, + (unsigned long long)(ck ? ck->global.my_node_id : 0)); + fflush(stderr); + free_config(ck); } + + struct utun_config* cs2 = parse_config(scf); + if (cs2 && cs2->global.my_public_key_hex) + sc_hex_to_binary(cs2->global.my_public_key_hex, ctx.srv_pubkey, SC_PUBKEY_SIZE); + free_config(cs2); + + 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); } + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); + utun_instance_set_tun_init_enabled(0); + setup(); + + ctx.ua = uasync_create(); + ctx.server = utun_instance_create(ctx.ua, scf); + ctx.client = utun_instance_create(ctx.ua, ccf); + if (!ctx.server || !ctx.client) goto done; + + utun_instance_init(ctx.server); + utun_instance_init(ctx.client); + + etcp_bind(ctx.server, ETCP_RT_ID_TEST, srv_traffic_handler); + etcp_bind(ctx.client, ETCP_RT_ID_TEST, cli_traffic_handler); + + /* store init state */ + ctx.prev_iface[0] = '\0'; + 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"); + + timeout_handle = uasync_set_timeout(ctx.ua, TIMEOUT_TB, NULL, to_cb, "to"); + + { uint64_t start = get_time_tb(); + while (!ctx.result && (int)(get_time_tb() - start) < TIMEOUT_TB + 50000) + uasync_poll(ctx.ua, POLL_MS); } + + fprintf(stderr, "final result=%d\n", ctx.result); fflush(stderr); + +done: + if (timeout_handle) uasync_cancel_timeout(ctx.ua, timeout_handle); + 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; } + cleanup(); + return (ctx.result == 2) ? 0 : 1; +}