Browse Source

auto_socket: TCP links support, chat_sync retry on socket change, link dedup by is_tcp

topo_upd
evgeny 2 months ago
parent
commit
d4a170b5c7
  1. 4
      ARCHITECTURE.md
  2. 14
      src/chat/chat_sync.c
  3. 4
      src/chat/chat_sync.h
  4. 3
      src/routing_layer/topo_group_connect.c
  5. 324
      src/transport_layer/auto_socket.c
  6. 65
      src/transport_layer/auto_socket.h
  7. 24
      src/transport_layer/etcp_connections.c
  8. 3
      src/transport_layer/etcp_connections.h
  9. 139
      src/transport_layer/node_conn_direct.c
  10. 10
      tests/Makefile.am
  11. 474
      tests/test_auto_socket_dynamic.c

4
ARCHITECTURE.md

@ -9,6 +9,10 @@
1. Модуль подключения p2p (ETCP/STCP)
Позволяет подключиться двум узлам друг к другу. Весь трафик шифруется. Для подключения нужен публичный ключ удаленного узла.
Позволяет использовать несколько линков до узла параллельно и динамически их переключать (агрегация/балансировка нагрузки/failover)
Модуль может работать в автоматическим режиме: мониторит изменения на интерфейсах и при появлении нового сетевого интерфейса автоматически добавлять подключения через новые интерфейсы.
При этом подключения работат параллельно по всем доступным интерфейсам распределяя нагрузку (какой интерфейс быстрее - тот бОльшую нагрузку берет на себя).
Разуммется, бесшовная работа - если есть хотяюы один живой линк то соединение работает. В процессе могут добавляться-удаляться линки.
2. Логические группы узлов (topo_group)
Один сервис может работать с несколькими группами.

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

4
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

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

324
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 <stdio.h>
/* Проверяет через /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;

65
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 <stdint.h>
@ -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);

24
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);

3
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

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

10
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 )); \

474
tests/test_auto_socket_dynamic.c

@ -0,0 +1,474 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdarg.h>
#include "../lib/platform_compat.h"
#include "test_utils.h"
#ifndef _WIN32
#include <unistd.h>
#include <net/if.h>
#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;
}
Loading…
Cancel
Save