diff --git a/AGENTS.md b/AGENTS.md index fb15bb56..fbabab36 100644 --- a/AGENTS.md +++ b/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 для редактирования исходников - запрещено - Проверяй на дублирование кода - не сделано ли это уже в другом месте diff --git a/doc/etcp_arch.md b/doc/etcp_arch.md index 729c27f6..8017f7ff 100644 --- a/doc/etcp_arch.md +++ b/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 + Тесты: # Правила работы с очередями в тестах: diff --git a/lib/debug_config.c b/lib/debug_config.c index f8b5e24e..87ca2887 100644 --- a/lib/debug_config.c +++ b/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} }; diff --git a/lib/debug_config.h b/lib/debug_config.h index 8808de92..73b247ee 100644 --- a/lib/debug_config.h +++ b/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) */ diff --git a/readme.md b/readme.md index 40abc118..ab45d263 100644 --- a/readme.md +++ b/readme.md @@ -3,6 +3,8 @@ Идея: объединить узлы в единую локальную сеть. Чтобы узлы сами находили оптимальные линки между собой, пробивали NAT где это можно, где нельзя - выбирали посредника с хорошей связью; чтобы подсто было добавлять новые узлы и адреса узлов сразу виделись во всей сети. +Похож на libp2p но трафик превращает в случайную последовательность кторую сложно классифицировать + Потенциальные варианты использования: - локальная сеть. хочу объединить офис, сотрудников (включая их работчие подсети), телефоны, офисы, сервер доступа итд в одно адресное пространство. - чат с файлообменником. узлы сети - это мемберы группы и одновременно сети. diff --git a/src/Makefile.am b/src/Makefile.am index 706463dd..9d20ce6a 100644 --- a/src/Makefile.am +++ b/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 \ diff --git a/src/conn_mgr.c b/src/conn_mgr.c index d7541324..eec8d4a8 100644 --- a/src/conn_mgr.c +++ b/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; diff --git a/src/etcp.c b/src/etcp.c index 256bf5d7..1de4cb01 100644 --- a/src/etcp.c +++ b/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; diff --git a/src/etcp.h b/src/etcp.h index 969f7a51..a8389bef 100644 --- a/src/etcp.h +++ b/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 diff --git a/src/etcp_api.c b/src/etcp_api.c index 9dca44a7..62b3922d 100644 --- a/src/etcp_api.c +++ b/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; diff --git a/src/etcp_api.h b/src/etcp_api.h index ee69c753..b2320bab 100644 --- a/src/etcp_api.h +++ b/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 * diff --git a/src/etcp_connect.c b/src/etcp_connect.c new file mode 100644 index 00000000..543f16a5 --- /dev/null +++ b/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 + +#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; +} diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 0ea84758..7434c8b7 100644 --- a/src/etcp_connections.c +++ b/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; diff --git a/src/etcp_router.c b/src/etcp_router.c index c0b6fb38..58fb44c0 100644 --- a/src/etcp_router.c +++ b/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]); } diff --git a/src/etcp_router.h b/src/etcp_router.h index 7da4ad26..7c3da3f8 100644 --- a/src/etcp_router.h +++ b/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); diff --git a/src/nat_transport.c b/src/nat_transport.c index 76681d44..0a9a5522 100644 --- a/src/nat_transport.c +++ b/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; } diff --git a/src/proxy/icmp_proxy.c b/src/proxy/icmp_proxy.c index c730d5cb..87dfabdf 100644 --- a/src/proxy/icmp_proxy.c +++ b/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; } diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index 8935f759..5ef3d1ee 100644 --- a/src/proxy/socks_proxy.c +++ b/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; } diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index f3d9abaf..c383849a 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/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); diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index efce0116..2e5556d8 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/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); diff --git a/src/proxy/udp_proxy.c b/src/proxy/udp_proxy.c index bb20dfeb..4885dcb0 100644 --- a/src/proxy/udp_proxy.c +++ b/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; diff --git a/src/route_bgp.c b/src/route_bgp.c index 168fcf42..2397db0e 100644 --- a/src/route_bgp.c +++ b/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) { diff --git a/src/route_bgp.h b/src/route_bgp.h index 4685d45f..4c6ac8e8 100644 --- a/src/route_bgp.h +++ b/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 diff --git a/src/routing.c b/src/routing.c index cd135715..0cfe2e80 100644 --- a/src/routing.c +++ b/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) { diff --git a/src/utun_instance.c b/src/utun_instance.c index ea6afe84..81ef3c3f 100644 --- a/src/utun_instance.c +++ b/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; diff --git a/src/utun_instance.h b/src/utun_instance.h index 6d663009..baa9c502 100644 --- a/src/utun_instance.h +++ b/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 diff --git a/tests/Makefile.am b/tests/Makefile.am index 807855f5..e51b5235 100644 --- a/tests/Makefile.am +++ b/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) diff --git a/tests/bbr_integration/test_bbr_integration.c b/tests/bbr_integration/test_bbr_integration.c index 110a2e13..a7515642 100644 --- a/tests/bbr_integration/test_bbr_integration.c +++ b/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); diff --git a/tests/test_etcp_congestion.c b/tests/test_etcp_congestion.c index 11eea1f5..ccc9eff3 100644 --- a/tests/test_etcp_congestion.c +++ b/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"); diff --git a/tests/test_etcp_connect.c b/tests/test_etcp_connect.c new file mode 100644 index 00000000..0bee3e51 --- /dev/null +++ b/tests/test_etcp_connect.c @@ -0,0 +1,242 @@ +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifndef _WIN32 +#include +#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; +} diff --git a/tests/test_etcp_reinit_inflight.c b/tests/test_etcp_reinit_inflight.c index 5b3a1c31..20270b69 100644 --- a/tests/test_etcp_reinit_inflight.c +++ b/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); diff --git a/tests/test_etcp_router_reconnect.c b/tests/test_etcp_router_reconnect.c new file mode 100644 index 00000000..4db3c2d6 --- /dev/null +++ b/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 +#include +#include +#include + +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifdef _WIN32 +#include +#include +#include +#define getpid _getpid +#else +#include +#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;inode; + 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;ibgp, 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; +} diff --git a/tests/test_etcp_router_unit.c b/tests/test_etcp_router_unit.c index 1fac2c79..3c9769a5 100644 --- a/tests/test_etcp_router_unit.c +++ b/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; diff --git a/tests/test_icmp_proxy.c b/tests/test_icmp_proxy.c index 917bc3a8..5077a207 100644 --- a/tests/test_icmp_proxy.c +++ b/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); diff --git a/tests/test_nat_transport.c b/tests/test_nat_transport.c index 62d582a3..44b43029 100644 --- a/tests/test_nat_transport.c +++ b/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); diff --git a/tests/test_udp_proxy.c b/tests/test_udp_proxy.c index bc057467..15e19a0c 100644 --- a/tests/test_udp_proxy.c +++ b/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,