Browse Source

etcp_connect + BGP_READY callbacks + router reconnect test

- etcp_connect.c: async connect API with EARLY/LATE/BGP_READY callbacks
- etcp_connect.h: declaration (via etcp_api.h)
- connect_deliver: fix to preserve unfired callback nodes
- etcp_set_routing_exchange_state: fires bgp_ready_cbk on TABLE_COMPLETE
- route_bgp: TABLE_COMPLETE handler + route_bgp_send_table_complete
- etcp.c: connections_count tracking, cleanup on close
- test_etcp_router_reconnect: a-b-c topology with b restart (etcp_connect)
- test_etcp_connect: unit tests for async connect API
- etcp_router: inflight/retrans/close improvements
- Various proxy/NAT/router cleanups
chatgui
Evgeny 3 months ago
parent
commit
0b5a710fe6
  1. 28
      AGENTS.md
  2. 38
      doc/etcp_arch.md
  3. 1
      lib/debug_config.c
  4. 4
      lib/debug_config.h
  5. 2
      readme.md
  6. 1
      src/Makefile.am
  7. 16
      src/conn_mgr.c
  8. 12
      src/etcp.c
  9. 3
      src/etcp.h
  10. 7
      src/etcp_api.c
  11. 36
      src/etcp_api.h
  12. 248
      src/etcp_connect.c
  13. 3
      src/etcp_connections.c
  14. 376
      src/etcp_router.c
  15. 39
      src/etcp_router.h
  16. 14
      src/nat_transport.c
  17. 10
      src/proxy/icmp_proxy.c
  18. 12
      src/proxy/socks_proxy.c
  19. 10
      src/proxy/tcp_proxy_client.c
  20. 12
      src/proxy/tcp_proxy_server.c
  21. 8
      src/proxy/udp_proxy.c
  22. 19
      src/route_bgp.c
  23. 1
      src/route_bgp.h
  24. 9
      src/routing.c
  25. 5
      src/utun_instance.c
  26. 5
      src/utun_instance.h
  27. 10
      tests/Makefile.am
  28. 4
      tests/bbr_integration/test_bbr_integration.c
  29. 4
      tests/test_etcp_congestion.c
  30. 242
      tests/test_etcp_connect.c
  31. 4
      tests/test_etcp_reinit_inflight.c
  32. 287
      tests/test_etcp_router_reconnect.c
  33. 14
      tests/test_etcp_router_unit.c
  34. 2
      tests/test_icmp_proxy.c
  35. 10
      tests/test_nat_transport.c
  36. 2
      tests/test_udp_proxy.c

28
AGENTS.md

@ -333,6 +333,34 @@ SOCKET=14, CONTROL=15, DUMP=16, TRAFFIC=17, DEBUG=18, GENERAL=19, NAT=20
Всегда когда начинаешь диагностику ознакомься со скиллом (используй skills)
Диагностика и поиск ошибок - это в первую очередь продумывание отладочных механизмов которые покажут понятную картину и точное представление об ошибке. Рассуждать надо как диагностикой добиться точной картины. И как приавльнее добавить диагностику чтобы не спамила лишними сообщениями и была понятной, логичной и информативной.
Важное правило отладки: вся отладка сводится к тому чтобы сделать удобные логи которые наглядно показывают поведение и проблемы. без лишнего мусора, компактно и по существу. очень желательно с деталями которые сильно повышают качество отладки.
при отладке нельзя: гадать и долго пытаться разбирать код пытаяясь найти причину. гораздо надёжнее с помощью логов понято точное поведение.
Если логов много - записывай в файл и потом анализируй.
Если видишь спам-логи - подумай как их выборочно отключить чтобы не мешали.
Для отладки добавляй в DEBUG_CATEGORY_DEBUG диагностические сообщения где надо. И включи эту категорию в настройках. Убирай только после проверки (путём запуска) когда все ошибки устранены.
сообщение выводится по или - либо debug_set_level(DEBUG_LEVEL_TRACE) - выводится ВСЁ независимо от настроек по категориям. Тоесть берется max(global level, category lavel)
Твой бич - ты постоянно гадаешь и анализируешь код. Что приводит к снежному кобу ошибок и неверных гипотез. 10 раз повторяю - только логи логи логи и никакого гадания. Лоооги!!! правильные логи покажут всё с предельной точностью. Вся суть отладки - информативные логи
И смотри логи. не задавливай их grep-ом. лучше больше. единственное с чем борись - это бесполезные спам логи.
Но полезные смотри всегда, и всегда оставляй логи которые выводятся нечасто.
Частые логи - это трафик которые >100 раз повторяются. Но которые мало раз обязательно оставлять и выводить.
можешь в тесте включить логи в нужный интервал чтобы не спамить. debug_set_level(DEBUG_LEVEL_TRACE) и выключить debug_set_level(DEBUG_LEVEL_NONE)
Если видишь проблему и ее решение неочевидно то выстраивай диагностику вокруг неё пока не будет очевидно где и что происходит не так. Это базовое и обязательное требование к отладке.
Подробные логи ты должен выводить не обрезая всегда и анализировать. Если лог большой - выведи в файл и анализируй файл.
Рассуждение должно быть примерно таким:
- данные повреждаются при отправке.
- где и почему непонятно. надо продумать ключевые точки где получим максимум полезной информации и не было большого объёма вывода лога.
- где лучше? каждый пакет логировать - много, но будет предельно точная картина
- Еще хорошо бы время - так мы заодно сможем найти проблемы с производительностью, найти проблемные места застревания кода.
- Но это большой объём. Можно ли без него? можно но это будет малоэффективно и хороших альтернатив пока не видно.
- Значит логируем весь трафик, а заодно включим полную трассировку функций.
- Запустили. Записали весь процесс в файл. Теперь можно анализировать. Давай сперва посчитаю количество строк с дампом: ... 50000.
- давай сделаю простой скрипт который дамп прочитает и сохранит в файл. И сверю с оригинальным файлом
- Давай заодно возьму произвольный фрагмент лога и изучу еа предмет явных проблем. Особенно инициализацию и освобождение ресурсов.
Сколько времени занимает, нет ли заклиниваний, нет ли ошибок или странностей
### Прочие правила:
- sed для редактирования исходников - запрещено
- Проверяй на дублирование кода - не сделано ли это уже в другом месте

38
doc/etcp_arch.md

@ -22,12 +22,9 @@ utun поддерживает встроенный ipv4 роутинг, може
Есть статические подключения (которые прописаны в конфиге) и есть динамические - которые устанавливаются в процессе по необходимости.
Обмен маршрутами происходит только через статические подключения.
toto:
- задать ограничение bandwidth для каждого сокета (общее).
## etcp_connections:
- обслуживает encrypt/decrypt и установку защищенного подключения.
- обслуживает encrypt/decrypt и установку защищенного подключения per-link.
- для одного etcp соединения может испольоваться несколько подключений одновременно (load balancing / filover)
## etcp_loadbalancer:
@ -50,9 +47,36 @@ toto:
- обеспечивает ретрансмиссии при передаче и сборку в правильной последовательности при приёме
- взаимодействует с pkt_normalizer и с udp сокетами для отправки-получения пакетов
## etcp_metric:
- обеспечивает обновление метрик каналов etcp_connections
- вызывается при получении ack
## stcp_* :
- аналогично etcp только упрощенный вариант, используя tcp в качестве транспорта.
## etcp_router:
- устанавливает связь между узлами, маршрутизирует трафик через цепочку узлов
## route_lib, route6_lib :
- управляет таблицами маршрутизации (ipv4/ipv6 lookup)
## route_node:
- хранит список доступных узлов (полученные от соседей) и их адреса/pubkeys для подключения.
- содержит апи для поиска-чтения-обновления узла
## route_bgp:
- составляет карту маршрутизации между узлами и обновляет ее при подключении-отключении узлов
## route_ping:
- модуль для определения типа nat (strict или eim) используя вспомогательный узел
Библиотеки:
## uasync
- модуль для асинхронной работы.
реализует сокеты, таймауты и сигналы из других потоков
- быстро работает с большим числом таймаутов (timeout_heap)
## ll_queue
- модуль работы с очередями, использует uasync.
- реализует fifo/lifo очереди, быстрый поиск по индексу, различные callback - сигнализация о состояниях очереди,
- congestion control: если много желающих записать -> очередь round-robin
Тесты:
# Правила работы с очередями в тестах:

1
lib/debug_config.c

@ -128,6 +128,7 @@ static const struct {
{"keepalive", DEBUG_CATEGORY_KEEPALIVE},
{"etcp_route", DEBUG_CATEGORY_ETCPROUTE},
{"bbr", DEBUG_CATEGORY_BBR},
{"etcp_dump", DEBUG_CATEGORY_ETCP_DUMP},
{"all", DEBUG_CATEGORY_ALL},
{NULL, DEBUG_CATEGORY_NONE}
};

4
lib/debug_config.h

@ -61,7 +61,8 @@ typedef int debug_category_t;
#define DEBUG_CATEGORY_KEEPALIVE 21 // Link keepalive logic
#define DEBUG_CATEGORY_ETCPROUTE 22 // ETCP routing
#define DEBUG_CATEGORY_BBR 23 // BBR congestion control
#define DEBUG_CATEGORY_COUNT 24 // Total number of categories
#define DEBUG_CATEGORY_ETCP_DUMP 24 // ETCP packet dump
#define DEBUG_CATEGORY_COUNT 25 // Total number of categories
#define DEBUG_CATEGORY_ALL (-1) // special value for all categories
/* Debug configuration structure */
@ -87,6 +88,7 @@ extern debug_config_t g_debug_config;
void debug_config_init(void);
/* Set debug level */
// итоговый level = max (global level - здесь задается, category level)
void debug_set_level(debug_level_t level);
/* Set debug level for specific category (0 = disabled, otherwise uses that level) */

2
readme.md

@ -3,6 +3,8 @@
Идея: объединить узлы в единую локальную сеть.
Чтобы узлы сами находили оптимальные линки между собой, пробивали NAT где это можно, где нельзя - выбирали посредника с хорошей связью; чтобы подсто было добавлять новые узлы и адреса узлов сразу виделись во всей сети.
Похож на libp2p но трафик превращает в случайную последовательность кторую сложно классифицировать
Потенциальные варианты использования:
- локальная сеть. хочу объединить офис, сотрудников (включая их работчие подсети), телефоны, офисы, сервер доступа итд в одно адресное пространство.
- чат с файлообменником. узлы сети - это мемберы группы и одновременно сети.

1
src/Makefile.am

@ -35,6 +35,7 @@ utun_CORE_SOURCES = \
pkt_normalizer.c \
packet_dump.c \
etcp_api.c \
etcp_connect.c \
control_server.c \
msg_transport.c \
firewall.c \

16
src/conn_mgr.c

@ -78,7 +78,7 @@ struct CONN_MGR* conn_mgr_init(struct UTUN_INSTANCE* instance) {
mgr->entry_capacity = 8;
mgr->entries = u_calloc(mgr->entry_capacity, sizeof(struct CONN_MGR_ENTRY));
if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_init: entries alloc failed"); u_free(mgr); return NULL; }
etcp_router_bind(instance, ETCP_ID_CONN_MGR, conn_mgr_router_recv_handler);
etcp_router_bind(instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler);
mgr->bg_ping_timer = uasync_set_timeout(instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping");
mgr->candidate_ping_timer = uasync_set_timeout(instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping");
mgr->bg_ping_cycle_start_tb = get_time_tb();
@ -91,7 +91,7 @@ void conn_mgr_destroy(struct CONN_MGR* mgr) {
if (!mgr) return;
if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; }
if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; }
etcp_router_unbind(mgr->instance, ETCP_ID_CONN_MGR);
etcp_router_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR);
for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_destroy(&mgr->entries[i]);
struct cm_exchange_pending* ep = mgr->exchange_pending;
while (ep) {
@ -236,7 +236,7 @@ int conn_mgr_disconnect_node(struct CONN_MGR* mgr, uint64_t node_id) {
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id);
if (!entry) return CONN_MGR_ERR_NOT_FOUND;
struct CONN_MGR_DISCONNECT pkt;
pkt.cmd = ETCP_ID_CONN_MGR;
pkt.cmd = ETCP_RT_ID_CONN_MGR;
pkt.subcmd = CONN_MGR_SUBCMD_DISCONNECT;
pkt.node_id = node_id;
struct ll_entry* qe = queue_entry_new(sizeof(pkt));
@ -296,7 +296,7 @@ int conn_mgr_send(struct CONN_MGR* mgr, uint64_t node_id, struct ll_entry* entry
struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, node_id);
if (!e || e->state != CONN_MGR_STATE_CONNECTED) { queue_entry_free(entry); return -1; }
e->last_traffic_tb = get_time_tb();
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(mgr->instance, node_id, ETCP_ID_CONN_MGR);
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(mgr->instance, node_id, ETCP_RT_ID_CONN_MGR);
if (!rconn) { queue_entry_free(entry); return -1; }
int ret = etcp_router_conn_send(rconn, entry->dgram, entry->len);
queue_entry_free(entry);
@ -627,7 +627,7 @@ static void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) {
uint8_t* pkt = u_malloc(pkt_size);
if (!pkt) { cm_deliver_result(entry, CONN_MGR_ERR_INTERNAL); return; }
struct CONN_MGR_DIRECT_REQ* req = (struct CONN_MGR_DIRECT_REQ*)pkt;
req->cmd = ETCP_ID_CONN_MGR;
req->cmd = ETCP_RT_ID_CONN_MGR;
req->subcmd = CONN_MGR_SUBCMD_DIRECT_REQ;
req->request_id = req_id;
req->addr_count = direct_count;
@ -690,7 +690,7 @@ static int cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry) {
entry->main.phase = 3;
struct CONN_MGR_INTERM_EXCHANGE_REQ req;
memset(&req, 0, sizeof(req));
req.cmd = ETCP_ID_CONN_MGR;
req.cmd = ETCP_RT_ID_CONN_MGR;
req.subcmd = CONN_MGR_SUBCMD_INTERM_EXCHANGE_REQ;
req.request_id = req_id;
req.candidate_count = mgr->best_candidate_count > 4 ? 4 : mgr->best_candidate_count;
@ -756,7 +756,7 @@ static void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CONN_
struct CONN_MGR_INTERM_SELECTED pkt;
memset(&pkt, 0, sizeof(pkt));
pkt.cmd = ETCP_ID_CONN_MGR;
pkt.cmd = ETCP_RT_ID_CONN_MGR;
pkt.subcmd = CONN_MGR_SUBCMD_INTERM_SELECTED;
pkt.request_id = 0;
pkt.count = sel;
@ -930,7 +930,7 @@ static void cm_handle_interm_exchange_req(struct ETCP_CONN* conn,
}
struct CONN_MGR_INTERM_EXCHANGE_RESP resp;
memset(&resp, 0, sizeof(resp));
resp.cmd = ETCP_ID_CONN_MGR;
resp.cmd = ETCP_RT_ID_CONN_MGR;
resp.subcmd = CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP;
resp.request_id = req->request_id;
resp.my_count = mgr->best_candidate_count > 4 ? 4 : mgr->best_candidate_count;

12
src/etcp.c

@ -263,12 +263,12 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n
static void etcp_on_up(struct ETCP_CONN* etcp) {
// if (!etcp->initialized) return;
DEBUG_WARN(DEBUG_CATEGORY_BGP, "[%s] initialized=%d", etcp->log_name, etcp->initialized);
DEBUG_WARN(DEBUG_CATEGORY_BGP, "[%s] UP links_up=%d initialized=%d", etcp->log_name, etcp->links_up, etcp->initialized);
if (etcp->up_cbk) etcp->up_cbk(etcp, etcp->up_arg);
}
static void etcp_on_down(struct ETCP_CONN* etcp) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "[%s]", etcp->log_name);
DEBUG_WARN(DEBUG_CATEGORY_BGP, "[%s] DOWN links_up=%d", etcp->log_name, etcp->links_up);
if (etcp->down_cbk) etcp->down_cbk(etcp, etcp->down_arg);
}
@ -338,6 +338,14 @@ void etcp_connection_close(struct ETCP_CONN* etcp) {
u_free(etcp->name);
{ struct UTUN_INSTANCE* inst = etcp->instance;
if (inst && inst->connections) {
struct ETCP_CONN** pp = &inst->connections;
while (*pp && *pp != etcp) pp = &(*pp)->next;
if (*pp) { *pp = etcp->next; if (inst->connections_count) inst->connections_count--; }
}
}
// Clear next pointer to prevent dangling references
etcp->next = NULL;

3
src/etcp.h

@ -175,7 +175,7 @@ struct ETCP_CONN {
// uint32_t total_packets_sent; // Total packets sent counter - Not used
// Flags
uint8_t routing_exchange_active; // 0 - не активен, 1 - надо инициировать обмен маршрутами (клиент), 2 - обмен маршрутами активен
uint8_t routing_exchange_active; // 0-не активен, 1-надо инициировать (клиент), 2-обмен идёт, 3-завершён, 4-пропущен (нет BGP)
uint8_t got_initial_pkt; //
uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен)
uint32_t session_id; // случайный ID сессии (генерируется клиентом) для защиты от ложного reinit
@ -189,6 +189,7 @@ struct ETCP_CONN {
void* up_arg; // аргумент для ready_cbk
etcp_on_conn_ready down_cbk; // callback при готовности соединения
void* down_arg; // аргумент для ready_cbk
void (*bgp_ready_cbk)(struct ETCP_CONN* conn); // вызывается когда BGP готов (завершён или пропущен)
uint32_t cnt_ack_hit_inf; // счетчик удлений из inflight
uint32_t cnt_ack_hit_sndq; // счетчик удалений inflight пакетов из sndq

7
src/etcp_api.c

@ -25,6 +25,13 @@ void etcp_set_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg
if (inst) { inst->etcp_new_conn_cbk = fn; inst->etcp_new_conn_arg = arg; }
}
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) {
if (!conn) return;
conn->routing_exchange_active = new_state;
if (new_state >= 3 && conn->bgp_ready_cbk)
conn->bgp_ready_cbk(conn);
}
int etcp_bind(struct UTUN_INSTANCE* inst, uint8_t id, etcp_recv_fn callback) {
if (!inst || !callback || (unsigned)id >= ETCP_MAX_BINDINGS) return -1;
inst->api_bindings.callbacks[id] = callback;

36
src/etcp_api.h

@ -22,17 +22,20 @@
#define ETCP_MAX_BINDINGS 256
// ETCP packet IDs
// === Raw ETCP packet IDs (etcp_bind/etcp_send) ===
#define ETCP_ID_DATA 0x00 // Пакет для передачи адресату
#define ETCP_ID_ROUTE_ENTRY 0x01 // Элемент роутинг-таблицы
#define ETCP_ID_NAT 0x02 // NAT трафик между узлами
#define ETCP_ID_SVC_ROUTE 0x03 // Маршрутизируемые сервисные пакеты (etcp_router)
#define ETCP_ID_TCP_PROXY 0x04 // TCP proxy exit node (server)
#define ETCP_ID_UDP_PROXY 0x05 // UDP datagram прокси (client ↔ exit)
#define ETCP_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit)
#define ETCP_ID_TCP_PROXY_CLIENT 0x07 // TCP proxy client (клиент)
#define ETCP_ID_MSG_TRANSPORT 0x10 // msg_transport — локальный IPC транспорт сообщений
#define ETCP_ID_CONN_MGR 0x11 // Connection Manager — management connections
#define ETCP_ID_ROUTE_ENTRY 0x01 // Элемент роутинг-таблицы (BGP)
#define ETCP_ID_SVC_ROUTE 0x03 // Транспорт роутера (etcp_router)
// === Router-сервисы (etcp_router_bind/etcp_route_send) ===
#define ETCP_RT_ID_DATA 0x00 // routing.c — маршрутизация данных
#define ETCP_RT_ID_NAT 0x02 // NAT трафик между узлами
#define ETCP_RT_ID_SVC_ROUTE 0x03 // etcp_router — транспорт роутера
#define ETCP_RT_ID_TCP_PROXY 0x04 // TCP proxy (клиент ↔ exit)
#define ETCP_RT_ID_UDP_PROXY 0x05 // UDP datagram прокси (client ↔ exit)
#define ETCP_RT_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit)
#define ETCP_RT_ID_MSG_TRANSPORT 0x10 // msg_transport — локальный IPC транспорт сообщений
#define ETCP_RT_ID_CONN_MGR 0x11 // Connection Manager — management connections
// Forward declarations
struct ETCP_CONN;
@ -51,8 +54,19 @@ typedef void (*etcp_recv_fn)(struct ETCP_CONN* conn, struct ll_entry* entry);
// Forward declaration
struct ETCP_CONN;
struct UTUN_INSTANCE;
struct NODEINFO_Q;
typedef void (*etcp_cbk_fn)(struct ETCP_CONN* conn, void* arg);
// ---- Background connection initialization ----
#define ETCP_CONNECT_EARLY 1
#define ETCP_CONNECT_LATE 2
#define ETCP_CONNECT_BGP_READY 4 // BGP-синхронизация завершена
typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int type);
int etcp_connect(struct UTUN_INSTANCE* instance, struct NODEINFO_Q* node,
etcp_connect_callback_t cb, void* arg, uint8_t flags);
/**
* @brief Установить callback при создании нового ETCP соединения
*
@ -118,6 +132,8 @@ void etcp_conn_set_ready_cbk(struct ETCP_CONN* conn, etcp_cbk_fn callback_fn, vo
void etcp_conn_set_up_cbk (struct ETCP_CONN* conn, etcp_cbk_fn callback_fn, void* arg);
void etcp_conn_set_down_cbk (struct ETCP_CONN* conn, etcp_cbk_fn callback_fn, void* arg);
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state);
/**
* @brief Внутренняя функция etcp: Коллбэк для очередей output ll_queue normalizer
*

248
src/etcp_connect.c

@ -0,0 +1,248 @@
#include "etcp_api.h"
#include "etcp.h"
#include "utun_instance.h"
#include "route_node.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "../lib/socket_compat.h"
#include <string.h>
#define DEBUG_CATEGORY_ETCP_CONNECT DEBUG_CATEGORY_ETCP
struct etcp_connect_cb_node {
etcp_connect_callback_t cb;
void* arg;
uint8_t flags;
struct etcp_connect_cb_node* next;
};
struct ETCP_CONNECT {
struct ETCP_CONNECT* next;
struct UTUN_INSTANCE* instance;
uint64_t node_id;
struct ETCP_CONN* conn;
struct etcp_connect_cb_node* cb_list;
void* initial_timer;
void* settle_timer;
uint16_t min_rtt;
uint8_t early_delivered : 1;
uint8_t done : 1;
};
static struct ETCP_CONNECT* connect_find(struct UTUN_INSTANCE* inst, uint64_t node_id) {
struct ETCP_CONNECT* ctx = inst->pending_connects;
while (ctx) { if (ctx->node_id == node_id) return ctx; ctx = ctx->next; }
return NULL;
}
static void connect_deliver(struct ETCP_CONNECT* ctx, int type) {
struct etcp_connect_cb_node** pp = &ctx->cb_list;
while (*pp) {
struct etcp_connect_cb_node* node = *pp;
if (type == 0 || (node->flags & (uint8_t)type)) {
*pp = node->next;
node->cb(node->arg, (type == 0) ? NULL : ctx->conn, type);
u_free(node);
} else {
pp = &node->next;
}
}
}
static void connect_cancel(struct ETCP_CONNECT* ctx) {
if (!ctx) return;
struct UTUN_INSTANCE* inst = ctx->instance;
if (inst) {
struct ETCP_CONNECT** pp = &inst->pending_connects;
while (*pp && *pp != ctx) pp = &(*pp)->next;
if (*pp) *pp = ctx->next;
if (ctx->initial_timer) { uasync_cancel_timeout(inst->ua, ctx->initial_timer); ctx->initial_timer = NULL; }
if (ctx->settle_timer) { uasync_cancel_timeout(inst->ua, ctx->settle_timer); ctx->settle_timer = NULL; }
}
u_free(ctx);
}
static void connect_create_links_v4(struct ETCP_CONNECT* ctx, struct NODEINFO_Q* node) {
const struct NODEINFO_IPV4_ADDR* addrs; int ac = get_node_v4_addrs(node, &addrs);
for (int i = 0; i < ac; i++) {
if (addrs[i].port == 0) continue;
{ int zero = 1; for (int j = 0; j < 4; j++) if (addrs[i].addr[j] != 0) { zero = 0; break; } if (zero) continue; }
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin));
sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4);
sin.sin_port = htons(addrs[i].port);
struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
struct ETCP_SOCKET* s = ctx->instance->etcp_sockets;
while (s) {
if (s->local_addr.ss_family == AF_INET)
etcp_link_new(ctx->conn, s, &sa, 0);
s = s->next;
}
}
}
static void connect_create_links_v6(struct ETCP_CONNECT* ctx, struct NODEINFO_Q* node) {
const struct NODEINFO_IPV6_ADDR* addrs; int ac = get_node_v6_addrs(node, &addrs);
for (int i = 0; i < ac; i++) {
if (addrs[i].port == 0) continue;
int zero = 1;
for (int j = 0; j < 16; j++) if (addrs[i].addr[j] != 0) { zero = 0; break; }
if (zero) continue;
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6));
sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, addrs[i].addr, 16);
sin6.sin6_port = htons(addrs[i].port);
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
struct ETCP_SOCKET* s = ctx->instance->etcp_sockets;
while (s) {
if (s->local_addr.ss_family == AF_INET6)
etcp_link_new(ctx->conn, s, &sa, 0);
s = s->next;
}
}
}
static void connect_bgp_ready_cb(struct ETCP_CONN* conn) {
struct ETCP_CONNECT* ctx = connect_find(conn->instance, conn->peer_node_id);
if (!ctx || ctx->done) return;
connect_deliver(ctx, ETCP_CONNECT_BGP_READY);
}
static void connect_initial_timeout_cb(void* arg) {
struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg;
if (!ctx || ctx->done) return;
ctx->done = 1;
ctx->initial_timer = NULL;
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] timeout for node 0x%016llx", (unsigned long long)ctx->node_id);
struct ETCP_CONN* conn = ctx->conn; ctx->conn = NULL;
connect_deliver(ctx, 0);
if (conn) etcp_connection_close(conn);
connect_cancel(ctx);
}
static void connect_settle_timeout_cb(void* arg) {
struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg;
if (!ctx || ctx->done) return;
ctx->done = 1;
ctx->settle_timer = NULL;
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] settle for node 0x%016llx, min_rtt=%u",
(unsigned long long)ctx->node_id, ctx->min_rtt);
struct ETCP_LINK* link = ctx->conn->links;
while (link) {
struct ETCP_LINK* next = link->next;
if (!link->initialized) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] removing failed link=%p for node 0x%016llx",
(void*)link, (unsigned long long)ctx->node_id);
etcp_link_close(link);
}
link = next;
}
ctx->conn->ready_cbk = NULL;
ctx->conn->ready_arg = NULL;
connect_deliver(ctx, ETCP_CONNECT_LATE);
connect_cancel(ctx);
}
static void connect_ready_cb(struct ETCP_CONN* conn, void* arg) {
struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg;
if (!ctx || ctx->done) return;
struct ETCP_LINK* link = conn->links;
while (link) {
if (link->initialized && link->rtt_last && link->rtt_last < ctx->min_rtt)
ctx->min_rtt = link->rtt_last;
link = link->next;
}
if (ctx->early_delivered) return;
ctx->early_delivered = 1;
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] early ready for node 0x%016llx, min_rtt=%u",
(unsigned long long)ctx->node_id, ctx->min_rtt);
connect_deliver(ctx, ETCP_CONNECT_EARLY);
if (ctx->initial_timer) { uasync_cancel_timeout(ctx->instance->ua, ctx->initial_timer); ctx->initial_timer = NULL; }
uint32_t settle_tb;
if (ctx->min_rtt == 0xFFFF) settle_tb = 20000;
else {
settle_tb = (uint32_t)ctx->min_rtt * 8;
if (settle_tb < 5000) settle_tb = 5000;
if (settle_tb > 20000) settle_tb = 20000;
}
ctx->settle_timer = uasync_set_timeout(ctx->instance->ua, (int)settle_tb, ctx,
connect_settle_timeout_cb, "etcp_connect_settle");
}
int etcp_connect(struct UTUN_INSTANCE* inst, struct NODEINFO_Q* node,
etcp_connect_callback_t cb, void* arg, uint8_t flags) {
if (!inst || !node || !cb) return -1;
uint64_t node_id = node->node.node_id;
struct ETCP_CONN* existing = inst->connections;
while (existing) {
if (existing->peer_node_id == node_id) {
struct ETCP_LINK* l = existing->links;
while (l) { if (l->link_state == 3) break; l = l->next; }
if (l) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx already connected, delivering immediately",
(unsigned long long)node_id);
if (flags & ETCP_CONNECT_EARLY) cb(arg, existing, ETCP_CONNECT_EARLY);
if (flags & ETCP_CONNECT_LATE) cb(arg, existing, ETCP_CONNECT_LATE);
if ((flags & ETCP_CONNECT_BGP_READY) && existing->routing_exchange_active >= 3)
cb(arg, existing, ETCP_CONNECT_BGP_READY);
return 0;
}
}
existing = existing->next;
}
struct ETCP_CONNECT* ctx = connect_find(inst, node_id);
if (ctx) {
struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node));
if (!cn) return -1;
cn->cb = cb; cn->arg = arg; cn->flags = flags; cn->next = ctx->cb_list; ctx->cb_list = cn;
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx already connecting, added cb",
(unsigned long long)node_id);
return 0;
}
struct ETCP_CONN* conn = etcp_connection_create(inst, NULL);
if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] failed to create ETCP_CONN for node 0x%016llx",
(unsigned long long)node_id); cb(arg, NULL, 0); return -1; }
conn->next = inst->connections; inst->connections = conn; inst->connections_count++;
if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] sc_init_ctx failed for node 0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); cb(arg, NULL, 0); return -1;
}
if (sc_set_peer_public_key(&conn->crypto_ctx, node->node.public_key, 0) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] sc_set_peer_public_key failed for node 0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); cb(arg, NULL, 0); return -1;
}
ctx = u_calloc(1, sizeof(struct ETCP_CONNECT));
if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] alloc failed");
etcp_connection_close(conn); cb(arg, NULL, 0); return -1; }
ctx->instance = inst; ctx->node_id = node_id; ctx->conn = conn;
ctx->min_rtt = 0xFFFF;
{
struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node));
if (!cn) { u_free(ctx); etcp_connection_close(conn); cb(arg, NULL, 0); return -1; }
cn->cb = cb; cn->arg = arg; cn->flags = flags; ctx->cb_list = cn;
}
conn->ready_cbk = connect_ready_cb;
conn->ready_arg = ctx;
conn->bgp_ready_cbk = connect_bgp_ready_cb;
connect_create_links_v4(ctx, node);
connect_create_links_v6(ctx, node);
if (!conn->links) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] no links created for node 0x%016llx",
(unsigned long long)node_id);
connect_deliver(ctx, 0);
etcp_connection_close(conn); u_free(ctx); return -1;
}
ctx->initial_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, ctx, connect_initial_timeout_cb, "etcp_connect_init");
ctx->next = inst->pending_connects; inst->pending_connects = ctx;
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] started for node 0x%016llx", (unsigned long long)node_id);
return 0;
}

3
src/etcp_connections.c

@ -1,4 +1,5 @@
#include "etcp_connections.h"
#include "etcp_api.h"
#include "../lib/socket_compat.h"
#include "../lib/platform_compat.h"
#include "../lib/getmyip.h"
@ -2010,7 +2011,7 @@ int init_connections(struct UTUN_INSTANCE* instance) {
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "init_connections: no peer public key configured for client %s", client->name);
}
etcp_conn->routing_exchange_active=1;// инициируем обмен маршрутами
etcp_set_routing_exchange_state(etcp_conn, 1); // инициируем обмен маршрутами
// Create links for this client
struct CFG_CLIENT_LINK* client_link = client->links;

376
src/etcp_router.c

@ -2,6 +2,15 @@
// Упрощённый TCP: восстановление порядка (recv_q), дедупликация, без переповторов
// ACK: периодическая отправка rx_seq (не чаще 100ms), idle-таймер для последнего seq
// Inflight контроль через send_q: при переполнении — очередь + retry-таймер, без дропов
/*
/----- ack ack <-- id
|(pause) |
in ---> {Q} -> [src, etcp] --> ... --> [dst,etcp] -> {asm_q} -> {buf q} ---> out
*/
#include "etcp_router.h"
#include "etcp.h"
#include "utun_instance.h"
@ -26,14 +35,20 @@ static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t*
static void router_send_to(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t svc_id, uint8_t flag_bits);
static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn);
static void router_send_resume_cb(void* arg);
static void router_no_route_retry_cb(void* arg);
static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn);
static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn);
static void router_close_finalize(void* arg);
static void free_entry(struct ll_entry* e) { queue_dgram_free(e); queue_entry_free(e); }
static int router_verify_signature(struct UTUN_INSTANCE* inst, struct ll_entry* entry, struct SVC_ROUTE_HDR* hdr, size_t* pl_len);
static void router_handle_ack(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CONN* rconn, uint32_t seq, uint64_t src_node_id);
static void router_forward_transit(struct UTUN_INSTANCE* inst, struct ll_entry* entry, struct SVC_ROUTE_HDR* hdr);
static int router_check_peer_restart(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CONN** prconn, struct SVC_ROUTE_HDR* hdr);
static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, uint32_t seq, const uint8_t* pl, size_t pl_len);
static void router_retransmit_one(struct ETCP_ROUTER_CONN* rconn, struct ROUTER_INFLIGHT* inf);
static void router_retrans_schedule(struct ETCP_ROUTER_CONN* rconn);
static void router_retrans_timer_cb(void* arg);
static void router_incoming_q_cb(struct ll_queue* q, void* arg);
// ====================================================================
// Управление ROUTER_CONN
@ -64,7 +79,7 @@ static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t*
uint8_t* dgram = u_malloc(total_len);
if (!dgram) { if (!flag_bits) rconn->tx_seq--; rconn->c_pkts_send_err++; return -1; }
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)dgram;
hdr->cmd = ETCP_ID_SVC_ROUTE;
hdr->cmd = ETCP_RT_ID_SVC_ROUTE;
hdr->dst_node_id = rconn->remote_node_id;
hdr->src_node_id = inst->node_id;
hdr->seq = seq;
@ -135,6 +150,26 @@ static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t*
DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_send_one: etcp_send failed ret=%d svc_id=%u seq=%u", ret, rconn->svc_id, seq);
} else {
rconn->c_pkts_sent++;
if (!flag_bits && pl_len > 0) {
struct ROUTER_INFLIGHT* inf = u_malloc(sizeof(struct ROUTER_INFLIGHT));
if (inf) {
memset(&inf->ll, 0, sizeof(inf->ll));
inf->ll.size = 4;
inf->seq = seq;
*(uint32_t*)inf->ll.data = seq;
inf->last_sent_tb = get_time_tb();
inf->send_count = 1;
inf->payload = u_malloc(pl_len);
if (inf->payload) {
memcpy(inf->payload, payload, pl_len);
inf->payload_len = pl_len;
queue_data_put_with_index(rconn->inflight_q, &inf->ll);
} else {
u_free(inf);
}
if (!rconn->retrans_timer) router_retrans_schedule(rconn);
}
}
}
return ret;
}
@ -145,7 +180,7 @@ static void router_send_to(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t svc
if (!conn) return;
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)u_malloc(SVC_ROUTE_HDR_SIZE);
if (!hdr) return;
hdr->cmd = ETCP_ID_SVC_ROUTE;
hdr->cmd = ETCP_RT_ID_SVC_ROUTE;
hdr->dst_node_id = dst;
hdr->src_node_id = inst->node_id;
hdr->seq = 0;
@ -181,9 +216,18 @@ static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn) {
}
// Закрыть conn: отправить CLOSE удалённой стороне, уведомить локальный сервис, очистить
static void router_close_finalize(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
void (*cb)(void*) = rconn->close_callback;
void* cb_arg = rconn->close_callback_arg;
queue_entry_free(&rconn->ll);
if (cb) cb(cb_arg);
}
static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
if (!rconn) return;
if (!rconn || rconn->closed) return;
struct UTUN_INSTANCE* inst = rconn->inst;
rconn->closed = 1;
router_send_close_to_service(rconn);
// Отправляем CLOSE удалённой стороне
struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, rconn->remote_node_id);
@ -192,6 +236,7 @@ static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
if (rconn->idle_ack_timer) { uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
if (rconn->no_route_timer) { uasync_cancel_timeout(inst->ua, rconn->no_route_timer); rconn->no_route_timer = NULL; }
if (rconn->recv_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
@ -202,14 +247,31 @@ static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->send_q); rconn->send_q = NULL;
}
if (rconn->inflight_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->inflight_q)) != NULL) {
struct ROUTER_INFLIGHT* inf = (struct ROUTER_INFLIGHT*)f;
if (inf->payload) u_free(inf->payload);
u_free(inf);
}
queue_free(rconn->inflight_q); rconn->inflight_q = NULL;
}
if (rconn->retrans_timer) { uasync_cancel_timeout(inst->ua, rconn->retrans_timer); rconn->retrans_timer = NULL; }
if (rconn->incoming_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->incoming_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_resume_callback(rconn->incoming_q);
queue_free(rconn->incoming_q); rconn->incoming_q = NULL;
}
if (inst->router_conns)
queue_remove_data(inst->router_conns, &rconn->ll);
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_conn: closed remote=%016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
queue_entry_free(&rconn->ll);
uasync_call_soon(inst->ua, rconn, router_close_finalize);
}
static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->no_route || rconn->closed) return;
int sq = queue_entry_count(rconn->send_q);
if (sq > 0) DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "SEND_Q_DRAIN: svc_id=%u send_q=%d inflight=%d",
rconn->svc_id, sq, (int32_t)(rconn->tx_seq - rconn->tx_acked));
@ -234,10 +296,126 @@ static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn) {
static void router_send_resume_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed) return;
rconn->send_resume_timer = NULL;
router_drain_send_q(rconn);
}
static void router_no_route_retry_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed) return;
rconn->no_route_timer = NULL;
struct ETCP_CONN* c = route_bgp_find_conn_for_node(rconn->inst->bgp, rconn->remote_node_id);
if (c) {
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router: route appeared for %016llx svc_id=%u, draining send_q=%d",
(unsigned long long)rconn->remote_node_id, rconn->svc_id, queue_entry_count(rconn->send_q));
rconn->no_route = 0;
router_drain_send_q(rconn);
return;
}
rconn->no_route_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_NO_ROUTE_RETRY_TB, rconn, router_no_route_retry_cb, "router_no_route");
}
// ====================================================================
// Ретрансмиты
// ====================================================================
static void router_retransmit_one(struct ETCP_ROUTER_CONN* rconn, struct ROUTER_INFLIGHT* inf) {
struct UTUN_INSTANCE* inst = rconn->inst;
size_t total_len = SVC_ROUTE_HDR_SIZE + inf->payload_len;
uint8_t* dgram = u_malloc(total_len);
if (!dgram) { rconn->c_pkts_send_err++; return; }
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)dgram;
hdr->cmd = ETCP_RT_ID_SVC_ROUTE;
hdr->dst_node_id = rconn->remote_node_id;
hdr->src_node_id = inst->node_id;
hdr->seq = inf->seq;
hdr->svc_id = rconn->svc_id;
hdr->flags = (rconn->sess_id << ROUTER_SESS_ID_SHIFT);
struct ETCP_CONN* conn = NULL;
struct ROUTE_BGP* bgp = inst->bgp;
if (bgp) {
struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, rconn->remote_node_id);
if (nq && nq->conn_mgr_type == CONN_TYPE_INDIRECT && nq->conn_mgr_intermediariy_count > 0) {
for (uint8_t i = 0; i < nq->conn_mgr_intermediariy_count; i++) {
conn = route_bgp_find_conn_for_node(bgp, nq->conn_mgr_intermediaries[i]);
if (conn) break;
}
}
}
if (!conn) conn = route_bgp_find_conn_for_node(bgp, rconn->remote_node_id);
if (!conn) {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router_retransmit: no route to %016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
u_free(dgram); rconn->c_pkts_send_err++; return;
}
memcpy(dgram + SVC_ROUTE_HDR_SIZE, inf->payload, inf->payload_len);
struct ll_entry* entry = queue_entry_new(0);
if (!entry) { u_free(dgram); rconn->c_pkts_send_err++; return; }
entry->dgram = dgram;
entry->len = total_len;
rconn->last_dgram_ts = get_current_timestamp();
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "RETRANS: svc_id=%u seq=%u attempt=%d → %016llx",
rconn->svc_id, inf->seq, inf->send_count + 1, (unsigned long long)rconn->remote_node_id);
int ret = etcp_send(conn, entry);
if (ret != 0) {
queue_dgram_free(entry); queue_entry_free(entry);
rconn->c_pkts_send_err++;
} else {
inf->last_sent_tb = get_time_tb();
inf->send_count++;
rconn->c_retrans_done++;
rconn->c_pkts_sent++;
}
}
static void router_retrans_schedule(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->retrans_timer) return;
uint64_t now = get_time_tb();
if (rconn->last_ack_changed_tb == 0) rconn->last_ack_changed_tb = now;
int64_t remaining = (int64_t)ROUTER_RETRANS_TIMEOUT_TB - (int64_t)(now - rconn->last_ack_changed_tb);
if (remaining <= 0) remaining = 1;
rconn->retrans_timer = uasync_set_timeout(rconn->inst->ua,
(uint32_t)remaining, rconn, router_retrans_timer_cb, "router_retrans");
}
static void router_retrans_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed) return;
rconn->retrans_timer = NULL;
uint64_t now = get_time_tb();
uint64_t elapsed = now - rconn->last_ack_changed_tb;
if (elapsed >= ROUTER_RETRANS_TIMEOUT_TB) {
rconn->no_ack_count++;
if (rconn->no_ack_count >= ROUTER_NO_ACK_MAX_RETRANS) {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router: no ACK progress for %d retrans cycles, closing svc_id=%u remote=%016llx",
rconn->no_ack_count, rconn->svc_id, (unsigned long long)rconn->remote_node_id);
router_close_and_notify(rconn);
return;
}
uint32_t s = rconn->tx_acked;
while ((int32_t)(s - rconn->tx_acked) < (int32_t)(rconn->tx_seq - rconn->tx_acked)) {
struct ROUTER_INFLIGHT* inf = (struct ROUTER_INFLIGHT*)queue_find_data_by_index(rconn->inflight_q, &s);
if (!inf) { s++; continue; }
router_retransmit_one(rconn, inf);
s++;
}
rconn->last_ack_changed_tb = now;
}
if (queue_entry_count(rconn->inflight_q) > 0) {
int64_t remaining = (int64_t)ROUTER_RETRANS_TIMEOUT_TB - (int64_t)(now - rconn->last_ack_changed_tb);
if (remaining <= 0) remaining = 1;
rconn->retrans_timer = uasync_set_timeout(rconn->inst->ua,
(uint32_t)remaining, rconn, router_retrans_timer_cb, "router_retrans");
}
}
static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, int force) {
if (!force && queue_entry_count(rconn->send_q) >= ROUTER_MAX_SEND_Q_PACKETS) {
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_send_q: FULL svc_id=%u count=%d — backpressure",
@ -280,7 +458,7 @@ static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) {
}
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)u_malloc(SVC_ROUTE_HDR_SIZE);
if (!hdr) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_ack: u_malloc failed"); return; }
hdr->cmd = ETCP_ID_SVC_ROUTE;
hdr->cmd = ETCP_RT_ID_SVC_ROUTE;
hdr->dst_node_id = rconn->remote_node_id;
hdr->src_node_id = inst->node_id;
hdr->seq = rconn->rx_seq;
@ -324,6 +502,7 @@ static void router_schedule_ack(struct ETCP_ROUTER_CONN* rconn) {
static void router_ack_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed) return;
rconn->ack_timer = NULL;
if (rconn->rx_seq != rconn->last_sent_ack_seq)
router_ack_do_send(rconn);
@ -339,6 +518,7 @@ static void router_ack_timer_cb(void* arg) {
static void router_idle_ack_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed) return;
rconn->idle_ack_timer = NULL;
if (rconn->rx_seq != rconn->last_sent_ack_seq)
router_ack_do_send(rconn);
@ -385,6 +565,39 @@ static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN
}
}
// ====================================================================
// Очередь приёма входящего трафика (между сетью и recv_q)
// ====================================================================
static void router_incoming_q_cb(struct ll_queue* q, void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed) { queue_resume_callback(q); return; }
struct ll_entry* e = queue_data_get(q);
if (!e) { queue_resume_callback(q); return; }
uint32_t seq = *(uint32_t*)e->data;
const uint8_t* pl = e->dgram;
size_t pl_len = e->len;
if (queue_find_data_by_index(rconn->recv_q, &seq)) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "router_incoming_q: dup in recv_q seq=%u, dropping", seq);
queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(q); return;
}
struct ll_entry* rq = queue_entry_new(4);
if (!rq) { queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(q); return; }
*(uint32_t*)rq->data = seq;
rq->dgram = e->dgram; rq->len = e->len; e->dgram = NULL;
queue_entry_free(e);
queue_data_put_with_index(rconn->recv_q, rq);
struct ETCP_CONN* conn = route_bgp_find_conn_for_node(rconn->inst->bgp, rconn->remote_node_id);
if (seq == rconn->rx_seq) router_try_assembly(rconn, conn);
router_schedule_ack(rconn);
queue_resume_callback(q);
}
// ====================================================================
// Извлечённые функции для etcp_router_recv_cb
// ====================================================================
@ -444,10 +657,27 @@ static void router_handle_ack(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CON
(void)src_node_id;
if (!rconn) return;
if ((int32_t)(seq - rconn->tx_acked) >= 0) {
uint32_t old_acked = rconn->tx_acked;
rconn->tx_acked = seq;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "SEND_Q_ACKED: svc_id=%u ack=%u inflight=%d send_q=%d from %016llx",
uint32_t clean_seq = old_acked;
while ((int32_t)(clean_seq - rconn->tx_acked) < 0) {
struct ROUTER_INFLIGHT* inf = (struct ROUTER_INFLIGHT*)queue_find_data_by_index(rconn->inflight_q, &clean_seq);
if (inf) {
queue_remove_data(rconn->inflight_q, &inf->ll);
if (inf->payload) u_free(inf->payload);
u_free(inf);
}
clean_seq++;
}
rconn->last_ack_changed_tb = get_time_tb();
rconn->no_ack_count = 0;
if (queue_entry_count(rconn->inflight_q) == 0 && rconn->retrans_timer) {
uasync_cancel_timeout(rconn->inst->ua, rconn->retrans_timer);
rconn->retrans_timer = NULL;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "SEND_Q_ACKED: svc_id=%u ack=%u inflight=%d send_q=%d inflight_q=%d from %016llx",
rconn->svc_id, seq, (int32_t)(rconn->tx_seq - rconn->tx_acked),
queue_entry_count(rconn->send_q), (unsigned long long)src_node_id);
queue_entry_count(rconn->send_q), queue_entry_count(rconn->inflight_q), (unsigned long long)src_node_id);
rconn->c_ack_recv++;
} else {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router: stale ACK seq=%u tx_acked=%u from %016llx, ignoring",
@ -502,11 +732,12 @@ static int router_check_peer_restart(struct UTUN_INSTANCE* inst, struct ETCP_ROU
static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn,
uint32_t seq, const uint8_t* pl, size_t pl_len) {
(void)conn;
rconn->last_dgram_ts = get_current_timestamp();
if (!rconn->peer_sync_done) { rconn->peer_sync_done = 1; rconn->rx_seq = seq; }
{
{// SEQ out of bounds
int32_t d = (int32_t)(seq - rconn->rx_seq);
if (d > ROUTER_MAX_INFLIGHT || d < -ROUTER_MAX_INFLIGHT) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router: seq=%u out of bounds, rx_seq=%u (d=%d), dropping",
@ -515,7 +746,8 @@ static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETC
}
}
if (((int32_t)(rconn->rx_seq - seq) > 0) || queue_find_data_by_index(rconn->recv_q, &seq)) {
if (((int32_t)(rconn->rx_seq - seq) > 0) || queue_find_data_by_index(rconn->recv_q, &seq)) {// DUP or delta seq<=0
router_schedule_ack(rconn);// сценарий потеряного ack
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "router: dup seq=%u rx_seq=%u, dropping", seq, rconn->rx_seq);
rconn->c_dup_dropped++; return;
}
@ -527,17 +759,14 @@ static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETC
if (!qe->dgram) { queue_dgram_free(qe); queue_entry_free(qe); return; }
qe->len = pl_len; memcpy(qe->dgram, pl, pl_len);
DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "router: queued seq=%u rx_seq=%u recv_q=%d from %016llx svc_id=%u",
seq, rconn->rx_seq, queue_entry_count(rconn->recv_q),
DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "router: incoming_q seq=%u rx_seq=%u from %016llx svc_id=%u",
seq, rconn->rx_seq,
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
queue_data_put_with_index(rconn->recv_q, qe);
if (seq == rconn->rx_seq) router_try_assembly(rconn, conn);
router_schedule_ack(rconn);
queue_data_put(rconn->incoming_q, qe);
}
// ====================================================================
// etcp_router_recv_cb — обработчик ETCP_ID_SVC_ROUTE
// etcp_router_recv_cb — обработчик ETCP_RT_ID_SVC_ROUTE
// ====================================================================
static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
@ -608,21 +837,7 @@ void etcp_router_destroy(struct UTUN_INSTANCE* inst) {
struct ll_entry* entry;
while ((entry = queue_data_get(inst->router_conns)) != NULL) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
router_send_close_to_service(rconn);
if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
if (rconn->idle_ack_timer) { uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
if (rconn->recv_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->recv_q);
}
if (rconn->send_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->send_q);
}
queue_entry_free(entry);
router_close_and_notify(rconn);
}
queue_free(inst->router_conns);
inst->router_conns = NULL;
@ -673,9 +888,17 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_
struct ETCP_CONN* etcp_conn = route_bgp_find_conn_for_node(inst->bgp, dst_node_id);
if (!etcp_conn) {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "etcp_route_send: no BGP route to %016llx svc_id=%u, dropping",
(unsigned long long)dst_node_id, svc_id);
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, dst_node_id, svc_id);
if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
int eq_ret = router_enqueue_send(rconn, entry->dgram + 1, payload_len, force);
queue_dgram_free(entry); queue_entry_free(entry);
if (eq_ret == 0) {
rconn->no_route = 1;
if (!rconn->no_route_timer)
rconn->no_route_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_NO_ROUTE_RETRY_TB, rconn, router_no_route_retry_cb, "router_no_route");
return 0;
}
return -1;
}
@ -741,6 +964,8 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t remote_node_i
if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
if (rconn->idle_ack_timer) { uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
if (rconn->no_route_timer) { uasync_cancel_timeout(inst->ua, rconn->no_route_timer); rconn->no_route_timer = NULL; }
if (rconn->retrans_timer) { uasync_cancel_timeout(inst->ua, rconn->retrans_timer); rconn->retrans_timer = NULL; }
struct ll_entry* f;
if (rconn->send_q) {
@ -751,9 +976,27 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t remote_node_i
while ((f = queue_data_get(rconn->recv_q)) != NULL) { u_check(f, "restart:recv_q", "router"); queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->recv_q);
}
if (rconn->inflight_q) {
while ((f = queue_data_get(rconn->inflight_q)) != NULL) {
struct ROUTER_INFLIGHT* inf = (struct ROUTER_INFLIGHT*)f;
if (inf->payload) u_free(inf->payload);
u_free(inf);
}
queue_free(rconn->inflight_q);
}
if (rconn->incoming_q) {
while ((f = queue_data_get(rconn->incoming_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_resume_callback(rconn->incoming_q);
queue_free(rconn->incoming_q);
}
rconn->send_q = queue_new(inst->ua, 0, 0, 0, "router_send_q");
rconn->recv_q = queue_new(inst->ua, ROUTER_RECVQ_HASH_SIZE, 0, 4, "router_recv_q");
rconn->inflight_q = queue_new(inst->ua, ROUTER_INFLIGHT_HASH_SIZE, 0, 4, "router_inflight_q");
rconn->incoming_q = queue_new(inst->ua, 0, 0, 0, "router_incoming_q");
queue_set_threshold(rconn->send_q, ROUTER_MAX_SEND_Q_PACKETS - 1, 0);
queue_set_callback(rconn->incoming_q, router_incoming_q_cb, rconn);
queue_set_threshold(rconn->incoming_q, 0, 0);
rconn->incoming_data_ready = 0;
rconn->tx_seq = 0;
rconn->rx_seq = 0;
@ -768,10 +1011,15 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t remote_node_i
rconn->c_pkts_rcvd = 0;
rconn->c_ack_sent = 0;
rconn->c_ack_recv = 0;
rconn->c_retrans_done = 0;
rconn->c_dup_dropped = 0;
rconn->c_oob_dropped = 0;
rconn->c_stale_ack = 0;
rconn->c_sign_fail = 0;
rconn->last_ack_changed_tb = 0;
rconn->no_route = 0;
rconn->no_ack_count = 0;
rconn->closed = 0;
}
void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t node_id, uint8_t svc_id,
@ -829,10 +1077,26 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst,
rconn->c_oob_dropped = 0;
rconn->c_stale_ack = 0;
rconn->c_sign_fail = 0;
rconn->c_retrans_done = 0;
rconn->retrans_timer = NULL;
rconn->last_ack_changed_tb = 0;
rconn->no_route = 0;
rconn->no_route_timer = NULL;
rconn->no_ack_count = 0;
rconn->closed = 0;
rconn->close_callback = NULL;
rconn->close_callback_arg = NULL;
rconn->recv_q = queue_new(inst->ua, ROUTER_RECVQ_HASH_SIZE, 0, 4, "router_recv_q");
if (!rconn->recv_q) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_conn_get: queue_new(recv_q) failed"); queue_entry_free(&rconn->ll); return NULL; }
rconn->send_q = queue_new(inst->ua, 0, 0, 0, "router_send_q");
if (!rconn->send_q) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_conn_get: queue_new(send_q) failed"); queue_free(rconn->recv_q); queue_entry_free(&rconn->ll); return NULL; }
rconn->inflight_q = queue_new(inst->ua, ROUTER_INFLIGHT_HASH_SIZE, 0, 4, "router_inflight_q");
if (!rconn->inflight_q) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_conn_get: queue_new(inflight_q) failed"); queue_free(rconn->send_q); queue_free(rconn->recv_q); queue_entry_free(&rconn->ll); return NULL; }
rconn->incoming_q = queue_new(inst->ua, 0, 0, 0, "router_incoming_q");
if (!rconn->incoming_q) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_conn_get: queue_new(incoming_q) failed"); queue_free(rconn->inflight_q); queue_free(rconn->send_q); queue_free(rconn->recv_q); queue_entry_free(&rconn->ll); return NULL; }
queue_set_callback(rconn->incoming_q, router_incoming_q_cb, rconn);
queue_set_threshold(rconn->incoming_q, 0, 0);
rconn->incoming_data_ready = 0;
queue_set_threshold(rconn->send_q, ROUTER_MAX_SEND_Q_PACKETS - 1, 0);
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_conn: new conn remote=%016llx svc_id=%u",
(unsigned long long)remote_node_id, svc_id);
@ -871,31 +1135,31 @@ void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) {
router_close_and_notify(rconn);
}
void etcp_router_conn_close_async(struct ETCP_ROUTER_CONN* rconn,
void (*on_close_done)(void* arg),
void* close_arg) {
if (!rconn) return;
rconn->close_callback = on_close_done;
rconn->close_callback_arg = close_arg;
router_close_and_notify(rconn);
}
void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id) {
if (!inst || !inst->router_conns) return;
struct ll_entry* entry;
while ((entry = queue_data_get(inst->router_conns)) != NULL) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
if (rconn->remote_node_id == remote_node_id) {
router_send_close_to_service(rconn);
if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
if (rconn->idle_ack_timer) { uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
if (rconn->recv_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->recv_q);
}
if (rconn->send_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_free(rconn->send_q);
}
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_conn: closed by node remote=%016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
queue_entry_free(&rconn->ll);
} else {
queue_data_put(inst->router_conns, entry);
// Iterate via hash chains to avoid inf-loop from put-back in FIFO iteration
struct ll_entry* to_close[256];
int count = 0;
struct ll_queue* q = inst->router_conns;
for (uint32_t slot = 0; slot < q->hash_size && count < 256; slot++) {
struct ll_entry* entry = q->hash_table[slot];
while (entry) {
struct ll_entry* next = entry->hash_next;
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
if (rconn->remote_node_id == remote_node_id)
to_close[count++] = entry;
entry = next;
}
}
for (int i = 0; i < count; i++)
router_close_and_notify((struct ETCP_ROUTER_CONN*)to_close[i]);
}

39
src/etcp_router.h

@ -14,7 +14,7 @@
#pragma pack(push, 1)
struct SVC_ROUTE_HDR {
uint8_t cmd; // ETCP_ID_SVC_ROUTE (0x03)
uint8_t cmd; // ETCP_ID_SVC_ROUTE (0x03) / ETCP_RT_ID_SVC_ROUTE
uint64_t dst_node_id;
uint64_t src_node_id;
uint32_t seq; // data: tx_seq; ACK: rx_seq (ожидаемый seq)
@ -54,16 +54,32 @@ struct ETCP_ROUTER_CONN {
struct ll_queue* send_q; // очередь ожидающих отправки (inflight полон)
void* send_resume_timer; // таймер возобновления отправки
uint8_t send_blocked; // 1 = inflight полон, ждём ack/таймер
uint8_t no_route; // 1 = нет BGP-маршрута, передача приостановлена
void* no_route_timer; // таймер 20ms проверки появления маршрута
uint8_t no_ack_count; // счётчик последовательных ретрансмиссий без ACK
uint8_t closed; // 1 = в процессе закрытия, таймеры игнорируют
void* close_callback; // пользовательский коллбэк завершения закрытия
void* close_callback_arg; // аргумент для close_callback
uint8_t sess_id; // наш session id (0-3)
uint8_t peer_sess_id; // последний sess_id от peer'а
uint8_t start_sent; // 0 = нужно отправить START в первом data
uint8_t peer_sync_done; // 0 = rx_seq ещё не синхронизирован с удалённым seq (авто-создание/рестарт)
// Ретрансмиты: inflight очередь
struct ll_queue* inflight_q; // хеш по seq (4 байта), непрерывный блок без SACK
void* retrans_timer; // таймер проверки застоя ACK
uint64_t last_ack_changed_tb; // когда tx_acked последний раз менялся (0.1ms)
// Очередь приёма входящего трафика (между сетью и recv_q)
struct ll_queue* incoming_q; // FIFO очередь входящих пакетов
uint8_t incoming_data_ready; // 1 = данные в incoming_q и callback был вызван
uint32_t c_pkts_sent; // успешные отправки данных
uint32_t c_pkts_send_err; // ошибки отправки
uint32_t c_pkts_rcvd; // получено и доставлено данных
uint32_t c_ack_sent; // отправлено ACK
uint32_t c_ack_recv; // получено ACK (последовательных)
uint32_t c_retrans_done; // число выполненых ретрансмитов
uint32_t c_dup_dropped; // дропнуто дубликатов seq
uint32_t c_oob_dropped; // дропнуто out-of-bounds seq
uint32_t c_stale_ack; // устаревших ACK
@ -78,6 +94,22 @@ struct ETCP_ROUTER_CONN {
#define ROUTER_SEND_RESUME_TB 500 // retry интервал send_q: 50ms
#define ROUTER_MAX_SEND_Q_PACKETS 64 // порог backpressure на send_q
// Ретрансмиты
#define ROUTER_RETRANS_TIMEOUT_TB 3000 // 300ms в timebase (0.1ms)
#define ROUTER_INFLIGHT_HASH_SIZE 1024
#define ROUTER_NO_ROUTE_RETRY_TB 200 // 20ms проверка появления маршрута
#define ROUTER_NO_ACK_MAX_RETRANS 17 // 17 × 300ms ≈ 5s таймаут без ACK
// Inflight запись — копия отправленного пакета для возможного ретрансмита
struct ROUTER_INFLIGHT {
struct ll_entry ll; // индекс по seq (4 байта, offset 0)
uint32_t seq;
uint64_t last_sent_tb; // время последней отправки (0.1ms)
uint8_t send_count; // число переотправок
uint8_t* payload; // копия данных
size_t payload_len;
};
// Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS)
struct ETCP_ROUTER_BINDINGS {
etcp_recv_fn callbacks[SVC_ROUTE_MAX_BINDINGS];
@ -112,6 +144,11 @@ int etcp_router_conn_send_signed(struct ETCP_ROUTER_CONN* rconn,
// Закрыть seq-подключение
void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn);
// Асинхронное закрытие: коллбэк вызывается после освобождения памяти rconn
void etcp_router_conn_close_async(struct ETCP_ROUTER_CONN* rconn,
void (*on_close_done)(void* arg),
void* close_arg);
// Закрыть все router_conn для указанного remote_node_id (peer умер)
void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id);

14
src/nat_transport.c

@ -20,7 +20,7 @@ static ip_str_t ip_host_to_str(uint32_t ip_host) {
// ==================== Callbacks ====================
// CLIENT: NAT TUN output → encapsulate in ETCP_ID_NAT → send to provider via etcp_router
// CLIENT: NAT TUN output → encapsulate in ETCP_RT_ID_NAT → send to provider via etcp_router
static void nat_transport_client_tun_out_cb(struct ll_queue* q, void* arg) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg;
if (!inst) { queue_resume_callback(q); return; }
@ -35,7 +35,7 @@ static void nat_transport_client_tun_out_cb(struct ll_queue* q, void* arg) {
uint8_t* new_dgram = u_malloc(total_len);
if (!new_dgram) { queue_dgram_free(pkt); queue_entry_free(pkt); queue_resume_callback(q); return; }
new_dgram[0] = ETCP_ID_NAT;
new_dgram[0] = ETCP_RT_ID_NAT;
memcpy(new_dgram + 1, &tr->self_node_id, 8);
memcpy(new_dgram + 9, pkt->dgram + 1, ip_len);
@ -78,7 +78,7 @@ static void nat_transport_provider_tun_out_cb(struct ll_queue* q, void* arg) {
uint8_t* new_dgram = u_malloc(total_len);
if (!new_dgram) { queue_dgram_free(pkt); queue_entry_free(pkt); queue_resume_callback(q); return; }
new_dgram[0] = ETCP_ID_NAT;
new_dgram[0] = ETCP_RT_ID_NAT;
memcpy(new_dgram + 1, &inst->nat_tr.self_node_id, 8);
memcpy(new_dgram + 9, pkt->dgram + 1, ip_len);
@ -98,7 +98,7 @@ static void nat_transport_provider_tun_out_cb(struct ll_queue* q, void* arg) {
}
}
// ETCP_ID_NAT receive via etcp_router: CLIENT gets response, PROVIDER gets request
// ETCP_RT_ID_NAT receive via etcp_router: CLIENT gets response, PROVIDER gets request
static void nat_transport_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!conn || !entry || !entry->dgram || entry->len < NAT_SVC_HDR_SIZE) {
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
@ -173,8 +173,8 @@ int nat_transport_init(struct UTUN_INSTANCE* inst) {
return -1;
}
if (etcp_router_bind(inst, ETCP_ID_NAT, nat_transport_etcp_recv_cb) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NAT, "Failed to bind ETCP_ID_NAT via etcp_router");
if (etcp_router_bind(inst, ETCP_RT_ID_NAT, nat_transport_etcp_recv_cb) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NAT, "Failed to bind ETCP_RT_ID_NAT via etcp_router");
tun_close(tr->nat_tun); tr->nat_tun = NULL;
eim_nat_destroy_ctx(&inst->nat);
return -1;
@ -200,7 +200,7 @@ void nat_transport_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->nat_tr.initialized) return;
struct nat_transport_ctx* tr = &inst->nat_tr;
etcp_router_unbind(inst, ETCP_ID_NAT);
etcp_router_unbind(inst, ETCP_RT_ID_NAT);
if (tr->nat_tun) { tun_close(tr->nat_tun); tr->nat_tun = NULL; }

10
src/proxy/icmp_proxy.c

@ -111,7 +111,7 @@ static void raw_read_cb(socket_t sock, void* arg) {
if (!e) return;
e->dgram = u_malloc(ICMP_PROXY_HDR_SIZE + payload_len);
if (!e->dgram) { queue_entry_free(e); return; }
e->dgram[0] = ETCP_ID_ICMP_PROXY;
e->dgram[0] = ETCP_RT_ID_ICMP_PROXY;
e->dgram[1] = ICMP_PROXY_SUBCMD_REPLY;
memcpy(e->dgram + 2, &r->client_node_id, 8);
memcpy(e->dgram + 10, &r->dst_ip, 4);
@ -149,7 +149,7 @@ static void exit_handle_request(struct ETCP_CONN* conn, struct ll_entry* entry)
if (e) {
e->dgram = u_malloc(ICMP_PROXY_HDR_SIZE + payload_len);
if (e->dgram) {
e->dgram[0] = ETCP_ID_ICMP_PROXY;
e->dgram[0] = ETCP_RT_ID_ICMP_PROXY;
e->dgram[1] = ICMP_PROXY_SUBCMD_REPLY;
memcpy(e->dgram + 2, &client_node_id, 8);
memcpy(e->dgram + 10, &dst_ip, 4);
@ -214,7 +214,7 @@ int icmp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id,
if (!e) return -1;
e->dgram = u_malloc(ICMP_PROXY_HDR_SIZE + payload_len);
if (!e->dgram) { queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_ICMP_PROXY;
e->dgram[0] = ETCP_RT_ID_ICMP_PROXY;
e->dgram[1] = ICMP_PROXY_SUBCMD_REQUEST;
memcpy(e->dgram + 2, &inst->node_id, 8);
memcpy(e->dgram + 10, &dst_ip, 4);
@ -322,7 +322,7 @@ int icmp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) {
}
}
etcp_router_bind(inst, ETCP_ID_ICMP_PROXY, icmp_proxy_recv_cb);
etcp_router_bind(inst, ETCP_RT_ID_ICMP_PROXY, icmp_proxy_recv_cb);
ctx->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "icmp_proxy initialized (exit=%d raw_sock=%d)", ctx->is_exit, (int)ctx->raw_sock);
return 0;
@ -330,7 +330,7 @@ int icmp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) {
void icmp_proxy_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !g_icmp_ctx) return;
etcp_router_unbind(inst, ETCP_ID_ICMP_PROXY);
etcp_router_unbind(inst, ETCP_RT_ID_ICMP_PROXY);
if (g_icmp_ctx->expire_timer) { uasync_cancel_timeout(g_icmp_ctx->ua, g_icmp_ctx->expire_timer); g_icmp_ctx->expire_timer = NULL; }
if (g_icmp_ctx->raw_sock != SOCKET_INVALID) {
if (g_icmp_ctx->raw_read_id) { uasync_remove_socket_t(g_icmp_ctx->ua, g_icmp_ctx->raw_sock); g_icmp_ctx->raw_read_id = NULL; }

12
src/proxy/socks_proxy.c

@ -51,7 +51,7 @@ static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, ui
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_TCP_PROXY;
e->dgram[0] = ETCP_RT_ID_TCP_PROXY;
e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
@ -331,7 +331,7 @@ static void process_http_request(struct socks_proxy_conn* c) {
u_free(pkt);
} else {
c->tx_buf = pkt; c->tx_len = total;
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
}
c->buf_len = 0;
@ -395,7 +395,7 @@ static void on_read_cb(struct ll_queue* q, void* arg) {
if (c->tx_buf) { memcpy(c->tx_buf, e->dgram, e->len); c->tx_len = e->len; }
else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: tx_buf malloc=%u failed sid=%08x — drop", e->len, c->stream_id); }
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
}
}
@ -408,7 +408,7 @@ static void tx_waiter_cb(struct ll_queue* q, void* arg) {
u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0;
queue_resume_callback(c->tc->read_queue);
} else {
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
}
}
@ -522,7 +522,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count,
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: REM_CLOSED sid=%08x", stream_id);
c->rem_closed = 1;
etcp_router_cancel_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter);
etcp_router_cancel_send_ready(c->inst, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter);
tcp_conn_push_close(c->tc);
return 1;
}
@ -530,7 +530,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count,
if (subcmd == TCP_PROXY_SUBCMD_ERROR) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: ERROR from exit sid=%08x", stream_id);
c->rem_closed = 1;
etcp_router_cancel_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter);
etcp_router_cancel_send_ready(c->inst, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter);
tcp_conn_push_close(c->tc);
return 1;
}

10
src/proxy/tcp_proxy_client.c

@ -70,7 +70,7 @@ static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, u
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_client_send_msg: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_client_send_msg: malloc(%zu) failed subcmd=%02x sid=%08x", TCP_PROXY_HDR_SIZE + len, subcmd, sid); queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_TCP_PROXY;
e->dgram[0] = ETCP_RT_ID_TCP_PROXY;
e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
@ -517,7 +517,7 @@ void tcp_proxy_client_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* en
struct tcp_proxy_client* proxy = inst ? inst->tcp_proxy_client : NULL;
if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) {
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_ID_TCP_PROXY_CLIENT) {
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_RT_ID_TCP_PROXY) {
uint64_t peer_id;
memcpy(&peer_id, entry->dgram + 1, 8);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE_ALL from %016llx — clearing client conns for peer",
@ -624,10 +624,10 @@ struct tcp_proxy_client* tcp_proxy_client_create(struct UTUN_INSTANCE* inst, str
}
if (inst) {
if (etcp_router_bind(inst, ETCP_ID_TCP_PROXY_CLIENT, tcp_proxy_client_router_recv_cb) != 0) {
if (etcp_router_bind(inst, ETCP_RT_ID_TCP_PROXY, tcp_proxy_client_router_recv_cb) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy client: etcp_router_bind failed");
} else {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client: etcp_router bind registered for ID=0x%02x", ETCP_ID_TCP_PROXY_CLIENT);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client: etcp_router bind registered for ID=0x%02x", ETCP_RT_ID_TCP_PROXY);
if (need_tun) { udp_proxy_init(inst, ua); icmp_proxy_init(inst, ua); }
}
}
@ -643,7 +643,7 @@ void tcp_proxy_client_destroy(struct tcp_proxy_client* p) {
int total = p->conn_count + p->socks_conn_count + p->http_conn_count;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy client destroying: lwip=%d socks=%d http=%d",
p->conn_count, p->socks_conn_count, p->http_conn_count);
if (p->inst) etcp_router_unbind(p->inst, ETCP_ID_TCP_PROXY_CLIENT);
if (p->inst) etcp_router_unbind(p->inst, ETCP_RT_ID_TCP_PROXY);
udp_proxy_destroy(p->inst);
icmp_proxy_destroy(p->inst);

12
src/proxy/tcp_proxy_server.c

@ -53,7 +53,7 @@ static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: malloc(%zu) failed", TCP_PROXY_HDR_SIZE + len); queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_TCP_PROXY_CLIENT;
e->dgram[0] = ETCP_RT_ID_TCP_PROXY;
e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
@ -222,7 +222,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
return;
}
} else {
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT,
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY,
&rc->pause_waiter, pause_resume_cb, rc);
return;
}
@ -259,7 +259,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry",
(int)rc->tc->sock, rc->stream_id, rc->tx_len);
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT,
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY,
&rc->pause_waiter, pause_resume_cb, rc);
}
}
@ -305,7 +305,7 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
struct tcp_proxy_server_conn** prev = &rc->ctx->conns;
while (*prev) { if (*prev == rc) { *prev = rc->next; rc->ctx->conn_count--; break; } prev = &(*prev)->next; }
}
if (rc->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, &rc->pause_waiter);
if (rc->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY, &rc->pause_waiter);
if (rc->tx_buf) { u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; }
if (rc->close_timer) { uasync_cancel_timeout(rc->ua, rc->close_timer); rc->close_timer = NULL; }
if (rc->diag_timer) { uasync_cancel_timeout(rc->ua, rc->diag_timer); rc->diag_timer = NULL; }
@ -453,7 +453,7 @@ void tcp_proxy_server_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL;
if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) {
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_ID_TCP_PROXY) {
if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_RT_ID_TCP_PROXY) {
uint64_t peer_id;
memcpy(&peer_id, entry->dgram + 1, 8);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE_ALL from %016llx — clearing server conns for peer",
@ -509,7 +509,7 @@ int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) {
ctx->inst = inst;
if (!ctx->enabled) return 0;
g_tcp_proxy_server_ctx = ctx;
etcp_router_bind(inst, ETCP_ID_TCP_PROXY, tcp_proxy_server_recv_cb);
etcp_router_bind(inst, ETCP_RT_ID_TCP_PROXY, tcp_proxy_server_recv_cb);
if (!inst->config->global.tcp_proxy_client_enabled) {
udp_proxy_init(inst, inst->ua);
icmp_proxy_init(inst, inst->ua);

8
src/proxy/udp_proxy.c

@ -52,7 +52,7 @@ static void flow_read_cb(socket_t sock, void* arg) {
if (!e) return;
e->dgram = u_malloc(UDP_PROXY_HDR_SIZE + n);
if (!e->dgram) { queue_entry_free(e); return; }
e->dgram[0] = ETCP_ID_UDP_PROXY;
e->dgram[0] = ETCP_RT_ID_UDP_PROXY;
e->dgram[1] = UDP_PROXY_SUBCMD_DATA;
memcpy(e->dgram + 2, &f->client_node_id, 8);
// Ответ src = оригинальный dst_ip:dest_port
@ -160,7 +160,7 @@ int udp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id,
if (!e) return -1;
e->dgram = u_malloc(UDP_PROXY_HDR_SIZE + payload_len);
if (!e->dgram) { queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_ID_UDP_PROXY;
e->dgram[0] = ETCP_RT_ID_UDP_PROXY;
e->dgram[1] = UDP_PROXY_SUBCMD_DATA;
memcpy(e->dgram + 2, &inst->node_id, 8);
memcpy(e->dgram + 10, &src_ip, 4);
@ -245,7 +245,7 @@ int udp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) {
ctx->is_exit = inst->tcp_proxy_server.enabled;
g_udp_ctx = ctx;
etcp_router_bind(inst, ETCP_ID_UDP_PROXY, udp_proxy_recv_cb);
etcp_router_bind(inst, ETCP_RT_ID_UDP_PROXY, udp_proxy_recv_cb);
ctx->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "udp_proxy initialized (exit=%d)", ctx->is_exit);
return 0;
@ -253,7 +253,7 @@ int udp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) {
void udp_proxy_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !g_udp_ctx) return;
etcp_router_unbind(inst, ETCP_ID_UDP_PROXY);
etcp_router_unbind(inst, ETCP_RT_ID_UDP_PROXY);
if (g_udp_ctx->expire_timer) { uasync_cancel_timeout(g_udp_ctx->ua, g_udp_ctx->expire_timer); g_udp_ctx->expire_timer = NULL; }
struct udp_flow* f = g_udp_ctx->flows;
while (f) { struct udp_flow* n = f->next;

19
src/route_bgp.c

@ -47,10 +47,24 @@ static void route_bgp_send_table_request(struct ROUTE_BGP* bgp, struct ETCP_CONN
}
}
static void route_bgp_send_table_complete(struct ROUTE_BGP* bgp, struct ETCP_CONN* conn) {
if (!bgp || !conn) return;
struct BGP_ROUTE_REQUEST* req = u_calloc(1, sizeof(struct BGP_ROUTE_REQUEST));
if (!req) return;
req->cmd = ETCP_ID_ROUTE_ENTRY;
req->subcmd = ROUTE_SUBCMD_TABLE_COMPLETE;
struct ll_entry* e = queue_entry_new(0);
if (!e) { u_free(req); return; }
e->dgram = (uint8_t*)req;
e->len = sizeof(struct BGP_ROUTE_REQUEST);
if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); }
}
static void route_bgp_add_to_senders(struct ROUTE_BGP* bgp, struct ETCP_CONN* conn);
static bool route_bgp_should_send_to(const struct NODEINFO_Q* nq, uint64_t target_id);
static void route_bgp_send_full_table(struct ROUTE_BGP* bgp, struct ETCP_CONN* conn);
static void route_bgp_handle_request_table(struct ROUTE_BGP* bgp, struct ETCP_CONN* conn);
static void route_bgp_send_table_complete(struct ROUTE_BGP* bgp, struct ETCP_CONN* conn);
static void nodeinfo_dump_log(const uint8_t* data, size_t len) {
if (!data || len < sizeof(struct BGP_NODEINFO_PACKET)) return;
@ -103,6 +117,7 @@ static const char* bgp_subcmd_name(uint8_t subcmd) {
case ROUTE_SUBCMD_PING_RESP: return "PING_RESP";
case ROUTE_SUBCMD_NAT_INFO: return "NAT_INFO";
case ROUTE_SUBCMD_NAT_CHECK_REQ: return "NAT_CHECK_REQ";
case ROUTE_SUBCMD_TABLE_COMPLETE: return "TABLE_COMPLETE";
default: return "?";
}
}
@ -190,6 +205,9 @@ static void route_bgp_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
route_bgp_handle_nat_info(bgp, from_conn, data, entry->len);
} else if (subcmd == ROUTE_SUBCMD_NAT_CHECK_REQ) {
route_bgp_handle_nat_check_req(bgp, from_conn, data, entry->len);
} else if (subcmd == ROUTE_SUBCMD_TABLE_COMPLETE) {
etcp_set_routing_exchange_state(from_conn, 3);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP initial sync complete with %s", from_conn->log_name);
}
queue_dgram_free(entry);
@ -1000,6 +1018,7 @@ static void route_bgp_handle_request_table(struct ROUTE_BGP* bgp, struct ETCP_CO
route_bgp_send_nodeinfo(bgp->local_node, conn);
route_bgp_send_full_table(bgp, conn);
route_bgp_add_to_senders(bgp, conn);
route_bgp_send_table_complete(bgp, conn);
}
static void route_bgp_handle_nat_info(struct ROUTE_BGP* bgp, struct ETCP_CONN* from_conn, const uint8_t* data, size_t len) {

1
src/route_bgp.h

@ -19,6 +19,7 @@
#define ROUTE_SUBCMD_WITHDRAW 0x06 // узел стал недоступен
#define ROUTE_SUBCMD_NAT_INFO 0x09 // информация о типе NAT клиента
#define ROUTE_SUBCMD_NAT_CHECK_REQ 0x0A // запрос от клиента на проверку NAT для сокета
#define ROUTE_SUBCMD_TABLE_COMPLETE 0x0B // завершение начальной синхронизации таблицы
#define MAX_HOPS 16
#define BGP_NODES_HASH_SIZE 256

9
src/routing.c

@ -22,9 +22,6 @@
#define IP_HDR_DST_ADDR_OFFSET 16
#define IP_HDR_MIN_SIZE 20
// ETCP packet ID for routing data
#define ETCP_ID_DATA 0x00
// Extract destination IP from IPv4 packet
// Returns 0 if not IPv4 or packet too small
static uint32_t extract_dst_ip(uint8_t* data, size_t len) {
@ -254,8 +251,8 @@ int routing_bind(struct UTUN_INSTANCE* instance) {
DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "routing_bind: instance is NULL");
return -1;
}
if (etcp_router_bind(instance, ETCP_ID_DATA, routing_pkt_from_etcp_cb) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "routing_bind: failed to bind ETCP_ID_DATA via etcp_router for node %016llx",
if (etcp_router_bind(instance, ETCP_RT_ID_DATA, routing_pkt_from_etcp_cb) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "routing_bind: failed to bind ETCP_RT_ID_DATA via etcp_router for node %016llx",
(unsigned long long)instance->node_id);
return -1;
}
@ -272,7 +269,7 @@ void routing_destroy(struct UTUN_INSTANCE* instance) {
(unsigned long long)instance->node_id);
// Unbind DATA handler from etcp_router
etcp_router_unbind(instance, ETCP_ID_DATA);
etcp_router_unbind(instance, ETCP_RT_ID_DATA);
// Clean up route table if exists
if (instance->rt) {

5
src/utun_instance.c

@ -265,7 +265,8 @@ struct UTUN_INSTANCE* utun_instance_create_from_config(struct UASYNC* ua, struct
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "Failed to allocate UTUN_INSTANCE");
return NULL;
}
instance->etcp_connect_timeout_tb = 20000;
// Initialize using common function (config ownership transferred to instance)
if (instance_init_common(instance, ua, config) != 0) {
// Cleanup on error - caller still owns config since we failed
@ -517,7 +518,7 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
if (instance->config->global.msg_transport_sock.ss_family != 0) {
int mt_ret = msg_transport_init(&instance->msg_t, instance, instance->ua,
&instance->config->global.msg_transport_sock,
8, ETCP_ID_MSG_TRANSPORT);
8, ETCP_RT_ID_MSG_TRANSPORT);
if (mt_ret != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "Failed to initialize msg_transport, continuing without IPC transport");
instance->msg_t = NULL;

5
src/utun_instance.h

@ -34,6 +34,7 @@ struct control_server;
struct msg_transport;
struct PING_CONTEXT;
struct CONN_MGR;
struct ETCP_CONNECT;
// uTun instance configuration
struct UTUN_INSTANCE {
@ -117,6 +118,10 @@ struct UTUN_INSTANCE {
// TCP proxy server (exit node)
struct tcp_proxy_server tcp_proxy_server;
// Pending background connections (etcp_connect API)
struct ETCP_CONNECT* pending_connects;
uint32_t etcp_connect_timeout_tb; // Initial timeout in 0.1ms units, default 20000 (2s)
};
// Functions

10
tests/Makefile.am

@ -32,6 +32,7 @@ check_PROGRAMS = \
test_lwip_tcp \
test_u_async_performance \
test_etcp_router \
test_etcp_router_reconnect \
test_etcp_bbr \
test_etcp_ping \
test_route_ping \
@ -47,6 +48,7 @@ check_PROGRAMS = \
test_bgp_route_exchange \
test_bgp_triangle \
test_conn_mgr \
test_etcp_connect \
test_bbr_integration \
test_intensive_memory_pool \
test_tcp_io \
@ -74,6 +76,7 @@ ETCP_CORE_OBJS = \
$(top_builddir)/src/utun-etcp_loadbalancer.o \
$(top_builddir)/src/utun-pkt_normalizer.o \
$(top_builddir)/src/utun-etcp_api.o \
$(top_builddir)/src/utun-etcp_connect.o \
$(top_builddir)/src/utun-etcp_debug.o \
$(top_builddir)/src/utun-etcp_dump.o
@ -198,6 +201,9 @@ test_etcp_router_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $
test_etcp_router_unit_SOURCES = test_etcp_router_unit.c
test_etcp_router_unit_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_etcp_router_reconnect_SOURCES = test_etcp_router_reconnect.c
test_etcp_router_reconnect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_tcp_proxy_server_SOURCES = test_tcp_proxy_server.c
test_tcp_proxy_server_CFLAGS = -I$(top_srcdir)/lib
test_tcp_proxy_server_LDADD = $(COMMON_LIBS)
@ -330,6 +336,10 @@ test_conn_mgr_SOURCES = test_conn_mgr.c
test_conn_mgr_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_conn_mgr_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_etcp_connect_SOURCES = test_etcp_connect.c
test_etcp_connect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_etcp_connect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_bbr_integration_SOURCES = bbr_integration/test_bbr_integration.c
test_bbr_integration_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_bbr_integration_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)

4
tests/bbr_integration/test_bbr_integration.c

@ -200,7 +200,7 @@ static void send_burst(struct test_ctx* ctx) {
while (sent < BURST_MAX) {
struct ll_entry* e = ll_alloc_lldgram(1 + PAYLOAD_SIZE);
if (!e) break;
e->dgram[0] = ETCP_ID_DATA;
e->dgram[0] = ETCP_RT_ID_DATA;
e->len = 1 + PAYLOAD_SIZE;
if (etcp_send(conn, e) == 0) { ctx->bytes_sent += PAYLOAD_SIZE; sent++; }
else { queue_entry_free(e); break; }
@ -359,7 +359,7 @@ int main(void) {
printf("Init sender...\n");
if (utun_instance_init(ctx.sender) < 0) { printf("ERROR: sender init failed\n"); return 1; }
etcp_bind(ctx.receiver, ETCP_ID_DATA, on_recv);
etcp_bind(ctx.receiver, ETCP_RT_ID_DATA, on_recv);
etcp_set_new_conn_cbk(ctx.receiver, NULL, &ctx);
printf("Creating dummynet on port %d ...\n", DN_PORT);

4
tests/test_etcp_congestion.c

@ -192,7 +192,7 @@ static void send_burst(struct test_ctx* ctx) {
while (sent < 64) {
struct ll_entry* e = ll_alloc_lldgram(sizeof(uint8_t) + PAYLOAD_SIZE);
if (!e) break;
e->dgram[0] = ETCP_ID_DATA;
e->dgram[0] = ETCP_RT_ID_DATA;
e->len = 1 + PAYLOAD_SIZE;
if (etcp_send(conn, e) == 0) { ctx->bytes_sent += e->len - 1; sent++; }
else { queue_entry_free(e); break; }
@ -328,7 +328,7 @@ int main(void) {
printf("Init sender...\n");
if (utun_instance_init(ctx.sender) < 0) { printf("Sender init failed\n"); return 1; }
etcp_bind(ctx.receiver, ETCP_ID_DATA, on_recv);
etcp_bind(ctx.receiver, ETCP_RT_ID_DATA, on_recv);
etcp_set_new_conn_cbk(ctx.receiver, NULL, &ctx);
printf("Creating dummynet...\n");

242
tests/test_etcp_connect.c

@ -0,0 +1,242 @@
#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>
#endif
#include "../src/etcp.h"
#include "../src/etcp_connections.h"
#include "../src/etcp_api.h"
#include "../src/config_parser.h"
#include "../src/config_updater.h"
#include "../src/utun_instance.h"
#include "../src/route_bgp.h"
#include "../src/route_node.h"
#include "../src/secure_channel.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 NID_A 0xAAAA000000000001ULL
#define NID_B 0xBBBB000000000002ULL
static struct UTUN_INSTANCE* g_a = NULL, *g_b = NULL;
static struct UASYNC* ua = NULL;
static volatile int result = 0;
static void* ttimer = NULL;
static char tdir[] = "/tmp/utun_ec_XXXXXX";
static char ca[256], cb[256];
static int pa = 0, pb = 0;
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 int lks(struct UTUN_INSTANCE* i) {
int n = 0; struct ETCP_CONN* c = i->connections;
while (c) { struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } c = c->next; }
return n;
}
static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 1; }
static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); result = 1; }
static void test1(void* arg); static void test2(void* arg);
static void test3(void* arg); static void test4(void* arg);
static void test5(void* arg); static void test6(void* arg);
static void test7(void* arg); static void test8(void* arg);
static void test9(void* arg); static void test10(void* arg);
/* ----- callback state per test ----- */
static volatile int cb_type = 0, cb_conn_ok = 0, cb_count = 0;
static void connect_ccb(void* a, struct ETCP_CONN* conn, int type) { cb_type = type; cb_conn_ok = (conn != NULL); cb_count++; }
/* ----- helper: create minimal NODEINFO_Q with 1 v4 addr ----- */
static struct NODEINFO_Q* mknode(uint64_t nid, const uint8_t pubkey[32],
uint8_t a, uint8_t b, uint8_t c, uint8_t d, uint16_t port) {
size_t extra = sizeof(struct NODEINFO_IPV4_SOCKET_META) + sizeof(struct NODEINFO_IPV4_ADDR);
size_t data_sz = sizeof(struct NODEINFO_Q) - sizeof(struct ll_entry) + extra;
struct NODEINFO_Q* nq = (struct NODEINFO_Q*)queue_entry_new(data_sz);
if (!nq) return NULL;
struct NODEINFO* ni = &nq->node;
ni->node_id = nid; ni->ver = 0;
memcpy(ni->public_key, pubkey, SC_PUBKEY_SIZE);
ni->local_v4_sockets = 1;
ni->local_v4_addrs = 1;
uint8_t* dyn = (uint8_t*)(ni + 1);
struct NODEINFO_IPV4_SOCKET_META* meta = (struct NODEINFO_IPV4_SOCKET_META*)dyn;
meta->id = 0; meta->config_type = CFG_SERVER_TYPE_PUBLIC; meta->nat_type = NAT_TYPE_UNKNOWN;
dyn += sizeof(struct NODEINFO_IPV4_SOCKET_META);
struct NODEINFO_IPV4_ADDR* addr = (struct NODEINFO_IPV4_ADDR*)dyn;
addr->addr[0] = a; addr->addr[1] = b; addr->addr[2] = c; addr->addr[3] = d;
addr->port = port; addr->type = ADDR_TYPE_INTERFACE; addr->socket_id = 0;
return nq;
}
/* ======================== Test 1: already connected → immediate EARLY+LATE ======================== */
static void test1(void* arg) {
(void)arg; if (result) return;
if (lks(g_a) < 1 || lks(g_b) < 1) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1"); return; }
struct NODEINFO_Q* nb = nodeinfo_find_by_id(g_a->bgp, NID_B);
if (!nb) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1b"); return; }
fprintf(stderr, "Test 1: already connected → EARLY|LATE\n"); fflush(stderr);
cb_type = 0; cb_conn_ok = 0; cb_count = 0;
int r = etcp_connect(g_a, nb, connect_ccb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE);
if (r != 0) { fail("etcp_connect returned error"); return; }
uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test2, "t2");
}
static void test2(void* arg) {
(void)arg; if (result) return;
if (!cb_count) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test2, "t2"); return; }
if (cb_count != 2) { fail("test1: expected 2 callbacks"); return; }
if (!cb_conn_ok) { fail("test1: conn=NULL in callback"); return; }
fprintf(stderr, " OK: %d callbacks, conn=OK\n", cb_count); fflush(stderr);
/* advance to test3 */
uasync_call_soon(ua, NULL, test3);
}
/* ======================== Test 3: unreachable node → timeout error ======================== */
static void test3(void* arg) {
(void)arg; if (result) return;
fprintf(stderr, "Test 3: unreachable node → error\n"); fflush(stderr);
struct SC_MYKEYS fk; sc_generate_keypair(&fk);
uint64_t fnid = 0xF000000000000001ULL;
struct NODEINFO_Q* fn = mknode(fnid, fk.public_key, 127,0,0,1, 49999);
if (!fn) { fail("mknode failed"); return; }
queue_data_put_with_index(g_a->bgp->nodes, &fn->ll);
g_a->etcp_connect_timeout_tb = 2000;
cb_type = -1; cb_conn_ok = -1; cb_count = 0;
int r = etcp_connect(g_a, fn, connect_ccb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE);
if (r != 0) { fail("etcp_connect(unreachable) returned error"); return; }
uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test4, "t4");
}
static void test4(void* arg) {
(void)arg; if (result) return;
if (!cb_count) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test4, "t4"); return; }
if (cb_count != 1) { fail("test3: expected 1 callback"); return; }
if (cb_conn_ok != 0) { fail("test3: expected conn=NULL"); return; }
if (cb_type != 0) { fail("test3: expected type=0 (error)"); return; }
fprintf(stderr, " OK: error callback delivered (NULL, type=0)\n"); fflush(stderr);
uasync_call_soon(ua, NULL, test5);
}
/* ======================== Test 5: double connect to unreachable → both callbacks ======================== */
static int t5_cb_count = 0;
static void t5_ccb(void* a, struct ETCP_CONN* conn, int type) {
(void)a; (void)conn; if (type == 0) t5_cb_count++;
}
static void test5(void* arg) {
(void)arg; if (result) return;
fprintf(stderr, "Test 5: double connect to unreachable → both callbacks\n"); fflush(stderr);
struct SC_MYKEYS fk; sc_generate_keypair(&fk);
uint64_t fnid = 0xF000000000000002ULL;
struct NODEINFO_Q* fn = mknode(fnid, fk.public_key, 127,0,0,1, 49998);
if (!fn) { fail("mknode failed"); return; }
queue_data_put_with_index(g_a->bgp->nodes, &fn->ll);
g_a->etcp_connect_timeout_tb = 2000;
t5_cb_count = 0;
int r1 = etcp_connect(g_a, fn, t5_ccb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE);
int r2 = etcp_connect(g_a, fn, t5_ccb, NULL, ETCP_CONNECT_EARLY);
if (r1 != 0 || r2 != 0) { fail("etcp_connect double call returned error"); return; }
uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test6, "t6");
}
static void test6(void* arg) {
(void)arg; if (result) return;
if (t5_cb_count < 2) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test6, "t6"); return; }
if (t5_cb_count != 2) { fail("test5: expected 2 error callbacks"); return; }
fprintf(stderr, " OK: both callbacks fired\n"); fflush(stderr);
uasync_call_soon(ua, NULL, test7);
}
/* ======================== Test 7: EARLY-only flag ======================== */
static void test7(void* arg) {
(void)arg; if (result) return;
struct NODEINFO_Q* nb = nodeinfo_find_by_id(g_a->bgp, NID_B);
if (!nb || lks(g_a) < 1) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test7, "t7"); return; }
fprintf(stderr, "Test 7: EARLY-only flag\n"); fflush(stderr);
cb_type = 0; cb_conn_ok = 0; cb_count = 0;
int r = etcp_connect(g_a, nb, connect_ccb, NULL, ETCP_CONNECT_EARLY);
if (r != 0) { fail("etcp_connect(EARLY) returned error"); return; }
uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test8, "t8");
}
static void test8(void* arg) {
(void)arg; if (result) return;
if (!cb_count) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test8, "t8"); return; }
if (cb_count != 1) { fail("test7: expected exactly 1 callback"); return; }
if (cb_type != ETCP_CONNECT_EARLY) { fail("test7: expected EARLY type"); return; }
fprintf(stderr, " OK: only EARLY fired\n"); fflush(stderr);
uasync_call_soon(ua, NULL, test9);
}
/* ======================== Test 9: LATE-only flag ======================== */
static void test9(void* arg) {
(void)arg; if (result) return;
struct NODEINFO_Q* nb = nodeinfo_find_by_id(g_a->bgp, NID_B);
if (!nb || lks(g_a) < 1) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test9, "t9"); return; }
fprintf(stderr, "Test 9: LATE-only flag\n"); fflush(stderr);
cb_type = 0; cb_conn_ok = 0; cb_count = 0;
int r = etcp_connect(g_a, nb, connect_ccb, NULL, ETCP_CONNECT_LATE);
if (r != 0) { fail("etcp_connect(LATE) returned error"); return; }
uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test10, "t10");
}
static void test10(void* arg) {
(void)arg; if (result) return;
if (!cb_count) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test10, "t10"); return; }
if (cb_count != 1) { fail("test9: expected exactly 1 callback"); return; }
if (cb_type != ETCP_CONNECT_LATE) { fail("test9: expected LATE type"); return; }
fprintf(stderr, " OK: only LATE fired\n"); fflush(stderr);
fprintf(stderr, "=== ALL PASSED ===\n"); fflush(stderr);
result = 2;
}
/* ======================== setup / cleanup / main ======================== */
static void setup(void) {
test_mkdtemp(tdir);
int base = 48000 + (getpid() % 10000); pa = base; pb = base + 1;
snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir);
wf(ca, "[global]\nmy_node_id=0x%llx\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_A, pa);
wf(cb, "[global]\nmy_node_id=0x%llx\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_B, pb);
config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb);
char *p0 = gv(ca,"pub"), *r0 = gv(ca,"priv"), *p1 = gv(cb,"pub"), *r1 = gv(cb,"priv");
wf(ca, "[global]\nmy_node_id=0x%llx\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", NID_A, r0, p0, pa, p1, pb);
wf(cb, "[global]\nmy_node_id=0x%llx\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_B, r1, p1, pb);
u_free(p0); u_free(r0); u_free(p1); u_free(r1);
}
static void cleanup(void) { test_unlink(ca); test_unlink(cb); test_rmdir(tdir); }
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN);
utun_instance_set_tun_init_enabled(0); setup();
ua = uasync_create();
g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb);
if (!g_a || !g_b) goto done;
utun_instance_init(g_a); utun_instance_init(g_b);
uasync_call_soon(ua, NULL, test1);
ttimer = uasync_set_timeout(ua, TIMEOUT_TB, NULL, to_cb, "to");
{ uint64_t start = get_time_tb();
while (!result && (int)(get_time_tb() - start) < TIMEOUT_TB + 50000) { uasync_poll(ua, POLL_MS); } }
fprintf(stderr, "final result=%d\n", result); fflush(stderr);
done:
if (ttimer) uasync_cancel_timeout(ua, ttimer);
if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); }
if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); }
if (ua) { uasync_destroy(ua, 0); ua = NULL; }
cleanup();
return (result == 2) ? 0 : 1;
}

4
tests/test_etcp_reinit_inflight.c

@ -209,7 +209,7 @@ static void monitor(void* arg) {
ctx->receiver = create_instance(ctx->ua, 0x2222222222222222ULL, s_priv, s_pub);
add_server(ctx->receiver, "srv1", SRV_PORT);
utun_instance_init(ctx->receiver);
etcp_bind(ctx->receiver, ETCP_ID_DATA, on_recv);
etcp_bind(ctx->receiver, ETCP_RT_ID_DATA, on_recv);
}
ctx->phase = 4;
break;
@ -296,7 +296,7 @@ int main(void) {
if (utun_instance_init(ctx.receiver) < 0) { printf("receiver init failed\n"); return 1; }
if (utun_instance_init(ctx.sender) < 0) { printf("sender init failed\n"); return 1; }
etcp_bind(ctx.receiver, ETCP_ID_DATA, on_recv);
etcp_bind(ctx.receiver, ETCP_RT_ID_DATA, on_recv);
printf("Creating dummynet...\n");
ctx.dn = dummynet_create(ctx.ua, "127.0.0.1", DN_PORT);

287
tests/test_etcp_router_reconnect.c

@ -0,0 +1,287 @@
// test_etcp_router_reconnect.c — ETCP router test: a-b-c, b restarts
// Топология: a [server] ←—[клиент b]—→ c [server], один UASYNC
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include "../lib/platform_compat.h"
#include "test_utils.h"
#ifdef _WIN32
#include <windows.h>
#include <direct.h>
#include <process.h>
#define getpid _getpid
#else
#include <unistd.h>
#endif
#include "../src/etcp.h"
#include "../src/etcp_connections.h"
#include "../src/etcp_router.h"
#include "../src/etcp_api.h"
#include "../src/route_bgp.h"
#include "../src/route_node.h"
#include "../src/config_parser.h"
#include "../src/utun_instance.h"
#include "../src/routing.h"
#include "../src/tun_if.h"
#include "../src/secure_channel.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#define TEST_SVC_ID 0x10
#define TEST_TIMEOUT_MS 60000
#define PHASE_PACKETS 500
#define MAX_PAYLOAD 1500
#define MIN_PAYLOAD 20
#define MAX_BURST 20
#define MIN_BURST 2
#define POLL_TB 50
static struct UTUN_INSTANCE* g_a = NULL;
static struct UTUN_INSTANCE* g_b = NULL;
static struct UTUN_INSTANCE* g_c = NULL;
static struct UASYNC* ua = NULL;
static char temp_dir[] = "/tmp/utun_test_XXXXXX";
static char a_conf[512], b_conf[512], c_conf[512];
static int a_port, b_port_a, b_port_c, c_port;
static const uint64_t node_a = 0xAAAAAAAA00000001ULL;
static const uint64_t node_b = 0xBBBBBBBB00000001ULL;
static const uint64_t node_c = 0xCCCCCCCC00000001ULL;
static struct SC_MYKEYS keys_a, keys_b, keys_c;
static char privhex_a[65], privhex_b[65], privhex_c[65], pubhex_a[65], pubhex_b[65], pubhex_c[65];
enum { ST_INIT, ST_PHASE1_SEND, ST_PHASE1_WAIT, ST_B_KILL, ST_B_RESTART, ST_PHASE3_SEND, ST_PHASE3_WAIT, ST_DONE };
static int g_state = ST_INIT;
static int g_test_ok = 0;
static int g_fail_code = 0;
static uint32_t g_total_sent_phase1 = 0, g_total_sent_phase3 = 0;
static uint32_t g_phase_sent = 0, g_seq = 0;
static int g_tick_counter = 0;
static uint32_t g_rcvd_total = 0, g_expected_seq = 0;
static int g_bgp_ready_mask = 0;
static int g_loop_cnt = 0;
static void hex_encode(const uint8_t* bin, int len, char* out) {
static const char hex[] = "0123456789abcdef";
for (int i = 0; i < len; i++) { out[i*2]=hex[bin[i]>>4]; out[i*2+1]=hex[bin[i]&0xF]; }
out[len*2]=0;
}
static void gen_payload(uint32_t seq, uint8_t* buf, int len) {
for (int i=0;i<len;i++) buf[i]=(uint8_t)((i^seq^0xA5)&0xFF);
}
static struct NODEINFO_Q* mknode(uint64_t nid, const uint8_t pubkey[32],
uint8_t a, uint8_t b, uint8_t c, uint8_t d, uint16_t port) {
size_t extra = sizeof(struct NODEINFO_IPV4_SOCKET_META) + sizeof(struct NODEINFO_IPV4_ADDR);
size_t data_sz = sizeof(struct NODEINFO_Q) - sizeof(struct ll_entry) + extra;
struct NODEINFO_Q* nq = (struct NODEINFO_Q*)queue_entry_new(data_sz);
if (!nq) return NULL;
struct NODEINFO* ni = &nq->node;
ni->node_id = nid; ni->ver = 0;
memcpy(ni->public_key, pubkey, SC_PUBKEY_SIZE);
ni->local_v4_sockets = 1;
ni->local_v4_addrs = 1;
uint8_t* dyn = (uint8_t*)(ni + 1);
struct NODEINFO_IPV4_SOCKET_META* meta = (struct NODEINFO_IPV4_SOCKET_META*)dyn;
meta->id = 0; meta->config_type = CFG_SERVER_TYPE_PUBLIC; meta->nat_type = NAT_TYPE_UNKNOWN;
dyn += sizeof(struct NODEINFO_IPV4_SOCKET_META);
struct NODEINFO_IPV4_ADDR* addr = (struct NODEINFO_IPV4_ADDR*)dyn;
addr->addr[0] = a; addr->addr[1] = b; addr->addr[2] = c; addr->addr[3] = d;
addr->port = port; addr->type = ADDR_TYPE_INTERFACE; addr->socket_id = 0;
return nq;
}
static int create_temp_configs(void) {
if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp fail\n"); return -1; }
int base = 42000+(getpid()%15000);
a_port=base; b_port_a=base+1; b_port_c=base+2; c_port=base+3;
if (sc_generate_keypair(&keys_a)!=SC_OK||sc_generate_keypair(&keys_b)!=SC_OK||sc_generate_keypair(&keys_c)!=SC_OK) { fprintf(stderr,"keygen fail\n"); return -1; }
hex_encode(keys_a.private_key,32,privhex_a); hex_encode(keys_a.public_key,32,pubhex_a);
hex_encode(keys_b.private_key,32,privhex_b); hex_encode(keys_b.public_key,32,pubhex_b);
hex_encode(keys_c.private_key,32,privhex_c); hex_encode(keys_c.public_key,32,pubhex_c);
snprintf(a_conf,sizeof(a_conf),"%s/a.conf",temp_dir);
FILE* f=fopen(a_conf,"w"); if(!f){perror("a_conf");return -1;}
fprintf(f,"[global]\nmy_node_id=0x%016llX\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n\n"
"[server:hub_a]\naddr=127.0.0.1:%d\ntype=public\n\n[allowed_keys]\nallow_all=1\n",(unsigned long long)node_a,privhex_a,pubhex_a,a_port);
fclose(f);
snprintf(b_conf,sizeof(b_conf),"%s/b.conf",temp_dir);
f=fopen(b_conf,"w"); if(!f){perror("b_conf");return -1;}
fprintf(f,"[global]\nmy_node_id=0x%016llX\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n\n"
"[server:s_a]\naddr=127.0.0.1:%d\ntype=public\n\n[server:s_c]\naddr=127.0.0.1:%d\ntype=public\n",
(unsigned long long)node_b,privhex_b,pubhex_b,b_port_a,b_port_c);
fclose(f);
snprintf(c_conf,sizeof(c_conf),"%s/c.conf",temp_dir);
f=fopen(c_conf,"w"); if(!f){perror("c_conf");return -1;}
fprintf(f,"[global]\nmy_node_id=0x%016llX\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.3/24\ntun_ifname=tun97\n\n"
"[server:hub_c]\naddr=127.0.0.1:%d\ntype=public\n\n[allowed_keys]\nallow_all=1\n",(unsigned long long)node_c,privhex_c,pubhex_c,c_port);
fclose(f);
return 0;
}
static void cleanup(void) { test_unlink(a_conf); test_unlink(b_conf); test_unlink(c_conf); test_rmdir(temp_dir); }
static void recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || !entry->dgram) { fprintf(stderr, "RECV: null entry\n"); g_test_ok=-1; g_fail_code=10; return; }
if (!conn && entry->len == 9) { queue_dgram_free(entry); queue_entry_free(entry); return; }
if (entry->len < 7) { fprintf(stderr, "RECV: short entry len=%u\n", (unsigned)entry->len); g_test_ok=-1; g_fail_code=10; queue_dgram_free(entry); queue_entry_free(entry); return; }
uint32_t seq=0; memcpy(&seq,entry->dgram+1,4);
uint16_t sz=0; memcpy(&sz, entry->dgram+5,2);
if (entry->len != (size_t)(7+sz)) { fprintf(stderr, "FAIL: size mismatch seq=%u\n", seq); g_test_ok=-1; g_fail_code=11; queue_dgram_free(entry); queue_entry_free(entry); return; }
if (seq != g_expected_seq) { fprintf(stderr, "FAIL: seq broken seq=%u expected=%u\n", seq, g_expected_seq); g_test_ok=-1; g_fail_code=12; queue_dgram_free(entry); queue_entry_free(entry); return; }
uint8_t exp[MAX_PAYLOAD]; gen_payload(seq, exp, sz);
if (memcmp(entry->dgram+7, exp, sz) != 0) { fprintf(stderr, "FAIL: corrupt seq=%u\n", seq); g_test_ok=-1; g_fail_code=13; queue_dgram_free(entry); queue_entry_free(entry); return; }
g_expected_seq++; g_rcvd_total++;
queue_dgram_free(entry); queue_entry_free(entry);
}
static int send_one_pkt(void) {
int len = MIN_PAYLOAD + (rand()%(MAX_PAYLOAD-MIN_PAYLOAD+1));
uint8_t* buf = u_malloc(7+len); if(!buf) return -1;
buf[0]=TEST_SVC_ID; memcpy(buf+1,&g_seq,4);
uint16_t s16=(uint16_t)len; memcpy(buf+5,&s16,2);
gen_payload(g_seq, buf+7, len);
struct ll_entry* e=queue_entry_new(0); if(!e){u_free(buf);return -1;}
e->dgram=buf; e->len=7+len;
int ret = etcp_route_send(g_a, node_c, e, 0);
if (ret == 0) g_seq++;
return ret;
}
static void connect_bgp_ready_cb(void* arg, struct ETCP_CONN* conn, int type) {
(void)conn;
if (type == ETCP_CONNECT_BGP_READY) {
int id = *(int*)arg; g_bgp_ready_mask |= (1 << id);
}
}
static const char* state_name(int s) {
static const char* n[]={"INIT","P1SEND","P1WAIT","B_KILL","B_RESTART","P3SEND","P3WAIT","DONE"}; return s<8?n[s]:"?";
}
static void send_burst(void) {
int b = MIN_BURST+(rand()%(MAX_BURST-MIN_BURST+1));
if (b>(int)(PHASE_PACKETS-g_phase_sent)) b=PHASE_PACKETS-g_phase_sent;
for (int i=0;i<b;i++) { if (send_one_pkt()==0) g_phase_sent++; else break; }
}
static void state_step(void) {
switch (g_state) {
case ST_INIT:
if (g_bgp_ready_mask == 3) {
struct ETCP_CONN* rc = route_bgp_find_conn_for_node(g_a->bgp, node_c);
if (rc) g_state=ST_PHASE1_SEND;
}
break;
case ST_PHASE1_SEND:
if (g_phase_sent < PHASE_PACKETS && g_tick_counter <= 0) {
send_burst(); g_tick_counter = 1; g_total_sent_phase1 = g_phase_sent;
}
break;
case ST_PHASE1_WAIT:
if (g_rcvd_total >= g_total_sent_phase1) g_state=ST_B_KILL;
break;
case ST_B_KILL:
g_b->running=0; utun_instance_destroy(g_b); g_b=NULL;
{ struct ETCP_CONN *c,*n;
for(c=g_a->connections;c;c=n){n=c->next;if(c->peer_node_id==node_b)etcp_connection_close(c);}
for(c=g_c->connections;c;c=n){n=c->next;if(c->peer_node_id==node_b)etcp_connection_close(c);} }
g_phase_sent=0; g_bgp_ready_mask=0; g_state=ST_B_RESTART;
break;
case ST_B_RESTART:
if (!g_b) {
g_b=utun_instance_create(ua,b_conf);
if (!g_b||utun_instance_init(g_b)<0) { fprintf(stderr,"FAIL: b restart\n"); g_test_ok=-1; g_fail_code=20; return; }
g_b->etcp_connect_timeout_tb = 300000;
struct NODEINFO_Q* nq_a = mknode(node_a, keys_a.public_key, 127,0,0,1, (uint16_t)a_port);
struct NODEINFO_Q* nq_c = mknode(node_c, keys_c.public_key, 127,0,0,1, (uint16_t)c_port);
if (!nq_a || !nq_c) { fprintf(stderr,"FAIL: mknode\n"); g_test_ok=-1; return; }
static int cc_a=0,cc_c=1;
etcp_connect(g_b,nq_a,connect_bgp_ready_cb,&cc_a,ETCP_CONNECT_BGP_READY);
etcp_connect(g_b,nq_c,connect_bgp_ready_cb,&cc_c,ETCP_CONNECT_BGP_READY);
}
if (g_bgp_ready_mask == 3) {
struct ETCP_CONN* rc = route_bgp_find_conn_for_node(g_a->bgp, node_c);
if (rc) { g_phase_sent=0; g_state=ST_PHASE3_SEND; }
}
break;
case ST_PHASE3_SEND:
if (g_phase_sent < PHASE_PACKETS && g_tick_counter <= 0) {
send_burst(); g_tick_counter = 1; g_total_sent_phase3 = g_phase_sent;
}
break;
case ST_PHASE3_WAIT:
if (g_rcvd_total >= g_total_sent_phase1 + g_total_sent_phase3) g_state=ST_DONE;
break;
case ST_DONE:
if (g_rcvd_total == g_total_sent_phase1+g_total_sent_phase3 && g_rcvd_total>0) g_test_ok=1;
else { fprintf(stderr, "FAIL: sent=%u rcvd=%u\n", g_total_sent_phase1+g_total_sent_phase3, g_rcvd_total); g_test_ok=-1; g_fail_code=30; }
return;
}
if (g_tick_counter>0) g_tick_counter--;
}
static void timeout_cb(void* arg) {
(void)arg;
fprintf(stderr, "TIMEOUT: state=%s sent_p1=%u sent_p3=%u rcvd=%u\n", state_name(g_state), g_total_sent_phase1, g_total_sent_phase3, g_rcvd_total);
g_test_ok=-1; g_fail_code=99;
}
int main(void) {
srand((unsigned)time(NULL));
if (create_temp_configs()!=0) return 1;
printf("=== ETCP Router Reconnect Test ===\nPorts: a=%d b_a=%d b_c=%d c=%d\n", a_port, b_port_a, b_port_c, c_port);
debug_config_init();
debug_set_level(DEBUG_LEVEL_ERROR);
utun_instance_set_tun_init_enabled(0);
ua=uasync_create(); if(!ua){cleanup();return 1;}
g_a=utun_instance_create(ua,a_conf); if(!g_a||utun_instance_init(g_a)<0){fprintf(stderr,"FAIL: a\n");goto fail;}
g_c=utun_instance_create(ua,c_conf); if(!g_c||utun_instance_init(g_c)<0){fprintf(stderr,"FAIL: c\n");goto fail;}
g_b=utun_instance_create(ua,b_conf); if(!g_b||utun_instance_init(g_b)<0){fprintf(stderr,"FAIL: b\n");goto fail;}
etcp_router_bind(g_c, TEST_SVC_ID, recv_handler);
struct NODEINFO_Q* nq_a = mknode(node_a, keys_a.public_key, 127,0,0,1, (uint16_t)a_port);
struct NODEINFO_Q* nq_c = mknode(node_c, keys_c.public_key, 127,0,0,1, (uint16_t)c_port);
if (!nq_a || !nq_c) { fprintf(stderr,"FAIL: mknode\n"); goto fail; }
g_b->etcp_connect_timeout_tb = 300000;
{ static int cc_a=0,cc_c=1;
etcp_connect(g_b,nq_a,connect_bgp_ready_cb,&cc_a,ETCP_CONNECT_BGP_READY);
etcp_connect(g_b,nq_c,connect_bgp_ready_cb,&cc_c,ETCP_CONNECT_BGP_READY); }
void* to_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS*10, NULL, timeout_cb, "to");
int max_loops = TEST_TIMEOUT_MS * 10 / POLL_TB;
while (!g_test_ok && g_loop_cnt < max_loops) {
uasync_poll(ua, POLL_TB);
state_step();
if (g_state==ST_PHASE1_SEND && g_phase_sent>=PHASE_PACKETS) g_state=ST_PHASE1_WAIT;
else if (g_state==ST_PHASE3_SEND && g_phase_sent>=PHASE_PACKETS) g_state=ST_PHASE3_WAIT;
g_loop_cnt++;
}
uasync_cancel_timeout(ua,to_id);
if(g_a){g_a->running=0;utun_instance_destroy(g_a);}
if(g_b){g_b->running=0;utun_instance_destroy(g_b);}
if(g_c){g_c->running=0;utun_instance_destroy(g_c);}
if(ua) uasync_destroy(ua,0);
cleanup();
if(g_test_ok==1){printf("=== TEST PASSED ===\n");return 0;}
printf("=== TEST FAILED: code=%d ===\n",g_fail_code);
return 1;
fail:
if(g_a){g_a->running=0;utun_instance_destroy(g_a);}
if(g_b){g_b->running=0;utun_instance_destroy(g_b);}
if(g_c){g_c->running=0;utun_instance_destroy(g_c);}
if(ua) uasync_destroy(ua,0);
cleanup();
return 1;
}

14
tests/test_etcp_router_unit.c

@ -64,7 +64,7 @@ static void inject(etcp_recv_fn recv_cb, struct UTUN_INSTANCE* inst,
const uint8_t* pl, size_t pl_len, int is_ack, uint8_t flags) {
struct SVC_ROUTE_HDR hdr;
memset(&hdr, 0, sizeof(hdr));
hdr.cmd = ETCP_ID_SVC_ROUTE;
hdr.cmd = ETCP_RT_ID_SVC_ROUTE;
hdr.dst_node_id = inst->node_id;
hdr.src_node_id = src;
hdr.seq = seq;
@ -98,7 +98,7 @@ static void inject(etcp_recv_fn recv_cb, struct UTUN_INSTANCE* inst,
fake_conn.instance = &inst; \
TESTASSERT(etcp_router_init(&inst) == 0); \
TESTASSERT(etcp_router_bind(&inst, TEST_SVC_ID, test_handler) == 0); \
recv_cb = inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE]; \
recv_cb = inst.api_bindings.callbacks[ETCP_RT_ID_SVC_ROUTE]; \
TESTASSERT(recv_cb != NULL); \
} while(0)
@ -179,7 +179,7 @@ static void inject_signed(etcp_recv_fn recv_cb, struct UTUN_INSTANCE* inst,
struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)dgram;
memset(hdr, 0, SVC_ROUTE_HDR_SIZE);
hdr->cmd = ETCP_ID_SVC_ROUTE;
hdr->cmd = ETCP_RT_ID_SVC_ROUTE;
hdr->dst_node_id = inst->node_id;
hdr->src_node_id = src;
hdr->seq = seq;
@ -221,11 +221,11 @@ static int test_init_destroy(void) {
if (etcp_router_init(&inst) != 0) FAIL("init");
if (inst.router_conns == NULL) FAIL("router_conns NULL after init");
if (inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE] == NULL) FAIL("recv_cb not bound");
if (inst.api_bindings.callbacks[ETCP_RT_ID_SVC_ROUTE] == NULL) FAIL("recv_cb not bound");
etcp_router_destroy(&inst);
if (inst.router_conns != NULL) FAIL("router_conns not NULL after destroy");
if (inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE] != NULL) FAIL("recv_cb still bound");
if (inst.api_bindings.callbacks[ETCP_RT_ID_SVC_ROUTE] != NULL) FAIL("recv_cb still bound");
uasync_destroy(ua, 0);
PASS();
@ -752,7 +752,7 @@ static int test_server_reinit(void) {
memset(&inst.router_bindings, 0, sizeof(inst.router_bindings));
if (etcp_router_init(&inst) != 0) FAIL("reinit failed");
etcp_router_bind(&inst, TEST_SVC_ID, test_handler);
recv_cb = inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE];
recv_cb = inst.api_bindings.callbacks[ETCP_RT_ID_SVC_ROUTE];
if (!recv_cb) FAIL("recv_cb not bound");
struct ETCP_ROUTER_CONN* c = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID);
@ -853,7 +853,7 @@ static int test_sign_too_short(void) {
rx_reset(0xF5);
struct SVC_ROUTE_HDR hdr;
memset(&hdr, 0, sizeof(hdr));
hdr.cmd = ETCP_ID_SVC_ROUTE;
hdr.cmd = ETCP_RT_ID_SVC_ROUTE;
hdr.dst_node_id = inst.node_id;
hdr.src_node_id = SIGNER_NODE;
hdr.seq = 0;

2
tests/test_icmp_proxy.c

@ -123,7 +123,7 @@ static void monitor(void* arg) {
if (cli_ok && exit_ok) {
g_test_phase = 1;
// Override client handler for test verification
etcp_router_bind(cli, ETCP_ID_ICMP_PROXY, cli_recv_cb);
etcp_router_bind(cli, ETCP_RT_ID_ICMP_PROXY, cli_recv_cb);
for (int i = 0; i < ICMP_PAYLOAD_SIZE; i++) send_buf[i] = (uint8_t)(rand() & 0xFF);
icmp_proxy_send_to_exit(cli, exit_node_id,
inet_addr("127.0.0.1"), inet_addr("127.0.0.100"), test_echo_id, test_echo_seq, send_buf, ICMP_PAYLOAD_SIZE);

10
tests/test_nat_transport.c

@ -285,12 +285,12 @@ static int test_provider_egress(void) {
}
DEBUG_INFO(DEBUG_CATEGORY_NAT, "Client conn peer_node_id=0x%016llx", (unsigned long long)client_conn->peer_node_id);
// Build ETCP_ID_NAT packet via etcp_route_send (new format: svc_id + src_node_id + ip_data)
// Build ETCP_RT_ID_NAT packet via etcp_route_send (new format: svc_id + src_node_id + ip_data)
size_t ip_len;
uint8_t* raw_ip = build_udp_pkt(0x0A0000FE, 40000, 0x08080808, 53, &ip_len);
size_t total = NAT_SVC_HDR_SIZE + ip_len;
uint8_t* dgram = u_malloc(total);
dgram[0] = ETCP_ID_NAT;
dgram[0] = ETCP_RT_ID_NAT;
memcpy(dgram + 1, &inst_client->nat_tr.self_node_id, 8);
memcpy(dgram + 9, raw_ip, ip_len);
free(raw_ip);
@ -306,7 +306,7 @@ static int test_provider_egress(void) {
queue_entry_free(entry);
return 0;
}
DEBUG_INFO(DEBUG_CATEGORY_NAT, "ETCP_ID_NAT sent from client to provider via etcp_router");
DEBUG_INFO(DEBUG_CATEGORY_NAT, "ETCP_RT_ID_NAT sent from client to provider via etcp_router");
// Poll to let provider process
int cycles = 0;
@ -455,12 +455,12 @@ static int test_full_roundtrip(void) {
memset(&inst_provider->nat.table[p], 0, sizeof(struct eim_nat_entry));
inst_provider->nat.next_port = inst_provider->nat.port_start;
// === Step 1: Send ETCP_ID_NAT from client to provider via etcp_router (egress) ===
// === Step 1: Send ETCP_RT_ID_NAT from client to provider via etcp_router (egress) ===
size_t ip_len;
uint8_t* raw_ip = build_udp_pkt(0x0A0000CD, 44444, 0x08080808, 80, &ip_len);
size_t total = NAT_SVC_HDR_SIZE + ip_len;
uint8_t* dgram = u_malloc(total);
dgram[0] = ETCP_ID_NAT;
dgram[0] = ETCP_RT_ID_NAT;
memcpy(dgram + 1, &inst_client->nat_tr.self_node_id, 8);
memcpy(dgram + 9, raw_ip, ip_len);
free(raw_ip);

2
tests/test_udp_proxy.c

@ -127,7 +127,7 @@ static void monitor(void* arg) {
if (cli_ok && exit_ok) {
g_test_phase = 1;
// Bind client handler for UDP replies (overrides what tcp_proxy_client_create set, for test verification)
etcp_router_bind(cli, ETCP_ID_UDP_PROXY, cli_recv_cb);
etcp_router_bind(cli, ETCP_RT_ID_UDP_PROXY, cli_recv_cb);
// Send UDP_REQUEST
for (int i = 0; i < PAYLOAD_SIZE; i++) send_buf[i] = (uint8_t)(rand() & 0xFF);
udp_proxy_send_to_exit(cli, exit_node_id,

Loading…
Cancel
Save