diff --git a/AGENTS.md b/AGENTS.md index 20492f44..7d697e7f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -535,6 +535,9 @@ MEMBER_SYNC=27 - Обязательно наличие подробной диагностики в функциях кода которые не сильно спамят (не часто вызываются). - Для каждого отладочного вывода подумай какая из доступной информация будет полезна чтобы можно было наиболее завершенно оценить состояние алгоритма и состояний влияющих на алгоритм. +## Оформление кода +- в .c: перед каждой функцией - краткое но понятное описание что делает функция (кроме совсем простых) +- в .h: в начале - описание модуля - для чего предназначен, как пользоваться, нюансы. Далее структуры с комментариями, далее публичные функции с описаниями (что делает, как пользоваться, какие аргументы, возврат, нюансы работы) ## chatgui (GUI Chat Client) diff --git a/src/routing_layer/etcp_router.c b/src/routing_layer/etcp_router.c index 2648c28a..d8248c1a 100644 --- a/src/routing_layer/etcp_router.c +++ b/src/routing_layer/etcp_router.c @@ -1,8 +1,12 @@ // etcp_router.c — Сервисный слой маршрутизации поверх ETCP // Упрощённый TCP: восстановление порядка (recv_q), дедупликация, ретрансмиты (inflight_q) -// ACK: периодическая отправка rx_seq (не чаще 100ms), idle-таймер для последнего seq +// ACK: периодическая отправка rx_seq (не чаще 10ms, ROUTER_ACK_INTERVAL_TB), idle-таймер для последнего seq // Inflight контроль через send_q: при переполнении — очередь + retry-таймер, без дропов // Подпись/шифрование — в route_crypto.c (encode в начале отправки, decode перед коллбэком) +// +// Поток сервисной кодограммы: +// in -> {send_q} -> [src: etcp_send] -> ... -> [dst: recv_q] -> {incoming_q} -> out +// ack возвращается по обратному пути; (pause) — backpressure при заполнении очередей. /* @@ -62,8 +66,10 @@ static void router_update_inflight_limit(struct ETCP_ROUTER_CONN* rconn); static void router_update_minrtt(struct ETCP_ROUTER_CONN* rconn); // ==================================================================== -// Транзитные очереди — per (src, dst) pair, backpressure через waiter на send_input_q +// Транзитные очереди — per (group_id, src_node_id, dst_node_id) pair, backpressure через waiter на send_input_q +// ==================================================================== +// Найти транзитную очередь для пары (group_id, src, dst) на conn (NULL если не создана). static struct TRANSIT_QUEUE* transit_queue_find(struct ETCP_CONN* conn, uint64_t group_id, uint64_t src_node_id, uint64_t dst_node_id) { @@ -75,6 +81,7 @@ static struct TRANSIT_QUEUE* transit_queue_find(struct ETCP_CONN* conn, return (struct TRANSIT_QUEUE*)queue_find_data_by_index(conn->transit_queues, key); } +// Найти транзитную очередь, при отсутствии — создать (реестр транзитных очередей живёт на conn). static struct TRANSIT_QUEUE* transit_queue_get_or_create(struct ETCP_CONN* conn, uint64_t group_id, uint64_t src_node_id, uint64_t dst_node_id) { @@ -108,6 +115,7 @@ static struct TRANSIT_QUEUE* transit_queue_get_or_create(struct ETCP_CONN* conn, return tq; } +// Уничтожить транзитную очередь: снять waiter, дропнуть оставшиеся пакеты, удалить из реестра conn. static void transit_queue_destroy(struct ETCP_CONN* conn, struct TRANSIT_QUEUE* tq) { if (conn->send_input_q) queue_waiter_cancel(conn->send_input_q, &tq->waiter); int dropped = 0; @@ -121,6 +129,8 @@ static void transit_queue_destroy(struct ETCP_CONN* conn, struct TRANSIT_QUEUE* queue_entry_free(&tq->ll); } +// Waiter-коллбэк транзитной очереди: send_input_q освободился → шлём один пакет, +// при непустой очереди снова встаём в waiter, иначе уничтожаем очередь. static void transit_queue_drain_cb(struct ll_queue* q, void* arg) { struct TRANSIT_QUEUE* tq = (struct TRANSIT_QUEUE*)arg; struct ETCP_CONN* conn = tq->conn; @@ -138,6 +148,7 @@ static void transit_queue_drain_cb(struct ll_queue* q, void* arg) { transit_queue_destroy(conn, tq); } +// Удалить все транзитные очереди соединения (вызывается при закрытии ETCP_CONN). void etcp_router_transit_queues_destroy(struct ETCP_CONN* conn) { if (!conn || !conn->transit_queues) return; struct ll_entry* entry; @@ -178,6 +189,7 @@ void etcp_router_conn_destroyed(struct ETCP_CONN* conn) { // Управление ROUTER_CONN // ==================================================================== +// Найти состояние seq-подключения по (group_id, remote_node_id, svc_id). NULL если не создано. static struct ETCP_ROUTER_CONN* router_conn_find(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id, uint8_t svc_id) { @@ -210,10 +222,13 @@ static struct ETCP_CONN* router_route_conn(struct UTUN_INSTANCE* inst, uint64_t return instance_find_conn(inst, remote_node_id); } +// Тонкая обёртка над router_route_conn для rconn (next_hop для отправки). static struct ETCP_CONN* router_send_conn(struct ETCP_ROUTER_CONN* rconn) { return router_route_conn(rconn->inst, rconn->group_id, rconn->remote_node_id); } +// Есть ли физический маршрут до узла (прямой/indirect/глобальный): 1 — есть, 0 — нет. +// Локальная доставка (dst == self) всегда возвращает 1. int etcp_router_has_route(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id) { if (!inst) return 0; if (dst_node_id == inst->node_id) return 1; @@ -341,6 +356,8 @@ static void router_restart_send(struct ETCP_ROUTER_CONN* rconn) { rconn->send_restart_pending = 1; } +// Синтетический loopback-conn (static) для доставки сервисного коллбэка при локальных событиях +// (CLOSE, restart) без реального сетевого соединения. static struct ETCP_CONN* loopback_conn(struct UTUN_INSTANCE* inst) { static struct ETCP_CONN c; memset(&c, 0, sizeof(c)); @@ -348,6 +365,7 @@ static struct ETCP_CONN* loopback_conn(struct UTUN_INSTANCE* inst) { return &c; } +// Доставить локальному сервису событие закрытия conn (пустая кодограмма ROUTER_SVC_HDR_SIZE). static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn) { struct UTUN_INSTANCE* inst = rconn->inst; etcp_recv_fn cb = inst->router_bindings.callbacks[rconn->svc_id]; @@ -367,7 +385,7 @@ static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn) { cb(loopback_conn(rconn->inst), e); } -// Закрыть conn: отправить CLOSE удалённой стороне, уведомить локальный сервис, очистить +// Финальный шаг закрытия: освобождает ll entry (вызывается через uasync_call_soon, уже вне очередей). static void router_close_finalize(void* arg) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; queue_entry_free(&rconn->ll); @@ -445,6 +463,8 @@ static void router_conn_reset(struct ETCP_ROUTER_CONN* rconn) { rconn->last_ack_changed_tb = 0; } +// Закрыть conn: уведомить локальный сервис, отправить CLOSE удалённой стороне, очистить +// очереди/таймеры, удалить из реестра router_conns и отложить освобождение ll через call_soon. static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) { if (!rconn || rconn->closed) return; struct UTUN_INSTANCE* inst = rconn->inst; @@ -581,6 +601,7 @@ static void router_send_watchdog_disarm(struct ETCP_ROUTER_CONN* rconn) { if (rconn->watchdog_timer) { uasync_cancel_timeout(rconn->inst->ua, rconn->watchdog_timer); rconn->watchdog_timer = NULL; } } +// Таймер повторной проверки маршрута (20ms): при появлении — снять no_route и возобновить отправку. static void router_no_route_retry_cb(void* arg) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; if (rconn->closed) return; @@ -600,6 +621,7 @@ static void router_no_route_retry_cb(void* arg) { // Ретрансмиты // ==================================================================== +// Отправить повторно один inflight-пакет (финальный wire-пакет уже закодирован — шлём копию как есть). static void router_retransmit_one(struct ETCP_ROUTER_CONN* rconn, struct ROUTER_INFLIGHT* inf) { struct ETCP_CONN* conn = router_send_conn(rconn); if (!conn) { @@ -630,6 +652,7 @@ static void router_retransmit_one(struct ETCP_ROUTER_CONN* rconn, struct ROUTER_ } } +// Взвести таймер ретрансмита на остаток до ROUTER_RETRANS_TIMEOUT_TB от last_ack_changed_tb (идемпотентно). static void router_retrans_schedule(struct ETCP_ROUTER_CONN* rconn) { if (rconn->retrans_timer) return; uint64_t now = get_time_tb(); @@ -640,6 +663,8 @@ static void router_retrans_schedule(struct ETCP_ROUTER_CONN* rconn) { (uint32_t)remaining, rconn, router_retrans_timer_cb, "router_retrans"); } +// Таймер ретрансмита: при застое ACK (ROUTER_RETRANS_TIMEOUT_TB) переотправляет все inflight-пакеты. +// После ROUTER_NO_ACK_MAX_RETRANS циклов без прогресса — закрывает conn. static void router_retrans_timer_cb(void* arg) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; if (rconn->closed) return; @@ -699,10 +724,13 @@ static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* pa // minRTT helpers // ==================================================================== +// Эффективный лимит inflight для текущего режима (probe или обычный). static uint32_t router_effective_max_inflight(struct ETCP_ROUTER_CONN* rconn) { return rconn->inflight_limit; } +// Пересчитать inflight_limit (=MINRTT_PROBE_MAX_INFLIGHT в probe-режиме, иначе max_inflight) и +// скорректировать send_blocked / возобновить отправку при росте лимита. static void router_update_inflight_limit(struct ETCP_ROUTER_CONN* rconn) { uint32_t new_limit = rconn->minrtt_probe ? MINRTT_PROBE_MAX_INFLIGHT : rconn->max_inflight; if (new_limit == rconn->inflight_limit) return; @@ -719,6 +747,7 @@ static void router_update_inflight_limit(struct ETCP_ROUTER_CONN* rconn) { } } +// Установить рабочий max_inflight (не ниже MINRTT_PROBE_MAX_INFLIGHT), пересчитать лимит и состояние канала. void etcp_router_set_max_inflight(struct ETCP_ROUTER_CONN* rconn, uint16_t new_max) { if (!rconn) return; if (new_max < MINRTT_PROBE_MAX_INFLIGHT) new_max = MINRTT_PROBE_MAX_INFLIGHT; @@ -727,6 +756,8 @@ void etcp_router_set_max_inflight(struct ETCP_ROUTER_CONN* rconn, uint16_t new_m router_track_inflight_state(rconn); } +// Отслеживать состояние загрузки канала: фиксирует моменты перехода loaded/unloaded +// (inflight >= max_inflight/2) для последующей оценки minRTT. static void router_track_inflight_state(struct ETCP_ROUTER_CONN* rconn) { uint32_t inflight = (uint32_t)queue_entry_count(rconn->inflight_q); uint32_t threshold = rconn->max_inflight / 2; @@ -740,12 +771,14 @@ static void router_track_inflight_state(struct ETCP_ROUTER_CONN* rconn) { } } +// Обновить minRTT: в probe-режиме (старые замеры > MINRTT_PROBE_TIMEOUT_TB) ограничить inflight, +// при разгруженном канале достаточно долго — добавить свежий замер rtt в окно и пересчитать среднее. static void router_update_minrtt(struct ETCP_ROUTER_CONN* rconn) { uint64_t now = get_time_tb(); uint16_t w_idx = rconn->minrtt_window_idx; uint8_t w_cnt = rconn->minrtt_window_count; - // Probe check: oldest entry freshness + // Проверка probe: свежесть самого старого замера в окне uint8_t prev_probe = rconn->minrtt_probe; if (w_cnt > 0) { uint8_t oldest = (w_idx - w_cnt + MINRTT_WINDOW_SIZE) % MINRTT_WINDOW_SIZE; @@ -758,7 +791,7 @@ static void router_update_minrtt(struct ETCP_ROUTER_CONN* rconn) { } if (prev_probe != rconn->minrtt_probe) router_update_inflight_limit(rconn); - // Update minRTT if channel is unloaded long enough + // Обновить minRTT, если канал разгружен достаточно долго if (!rconn->inflight_was_loaded) { uint64_t unloaded_dur = now - rconn->last_inflight_unloaded_tb; uint64_t min_dur = (uint64_t)rconn->rtt * 2; @@ -784,10 +817,13 @@ static void router_update_minrtt(struct ETCP_ROUTER_CONN* rconn) { // ACK timers // ==================================================================== +// Тонкая обёртка для единообразного вызова отправки ACK из таймеров. static void router_ack_do_send(struct ETCP_ROUTER_CONN* rconn) { router_send_ack(rconn); } +// Отправить ACK с текущим rx_seq (header-only пакет, seq=rx_seq, payload пуст). +// timestamp — эхо последнего data-пакета пира с поправкой на локальное время (для RTT). static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) { struct UTUN_INSTANCE* inst = rconn->inst; struct ETCP_CONN* conn = rconn->recv_conn; @@ -834,6 +870,8 @@ static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) { rconn->last_ack_sent_tb = get_time_tb(); } +// Запланировать отправку ACK: снять idle-таймер, и либо отправить сразу (интервал истёк), +// либо взвести ack_timer на остаток ROUTER_ACK_INTERVAL_TB. static void router_schedule_ack(struct ETCP_ROUTER_CONN* rconn) { if (rconn->idle_ack_timer) { uasync_cancel_timeout(rconn->inst->ua, rconn->idle_ack_timer); @@ -851,6 +889,7 @@ static void router_schedule_ack(struct ETCP_ROUTER_CONN* rconn) { ROUTER_ACK_INTERVAL_TB - (uint32_t)elapsed, rconn, router_ack_timer_cb, "router_ack"); } +// Периодический ACK-таймер: отправить ACK при неподтверждённом rx_seq, иначе взвести idle-таймер. static void router_ack_timer_cb(void* arg) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; if (rconn->closed) return; @@ -867,6 +906,7 @@ static void router_ack_timer_cb(void* arg) { } } +// Idle-таймер: дослать последний ACK после паузы в приёме (чтобы пир не висел на ретрансмитах). static void router_idle_ack_timer_cb(void* arg) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; if (rconn->closed) return; @@ -948,6 +988,7 @@ static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* con cb(conn, svc_entry); } +// Собрать доставку из recv_q: последовательно доставлять пакеты по возрастанию rx_seq, пока есть следующий. static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn) { uint32_t next = rconn->rx_seq; while (1) { @@ -966,6 +1007,8 @@ static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN // Очередь приёма входящего трафика (между сетью и recv_q) // ==================================================================== +// Коллбэк incoming_q: переместить пакет в recv_q (с дедупликацией по seq), при совпадении с +// rx_seq — собрать доставку, затем запланировать ACK. 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; } @@ -998,6 +1041,8 @@ static void router_incoming_q_cb(struct ll_queue* q, void* arg) { // Функции приёма etcp_router_recv_cb // ==================================================================== +// Обработать ACK: обновить RTT/jitter, снять подтверждённые inflight-пакеты, при необходимости +// возобновить отправку (снять send_blocked). static void router_handle_ack(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CONN* rconn, uint32_t seq, uint64_t src_node_id, uint16_t ack_ts) { (void)src_node_id; @@ -1044,6 +1089,8 @@ static void router_handle_ack(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CON if (rconn->send_blocked) { rconn->send_blocked = 0; router_send_kick(rconn); } } +// Ретранслировать транзитный пакет к dst_node_id: отправить напрямую, при занятости send_input_q — +// поставить в транзитную очередь с backpressure. static void router_forward_transit(struct UTUN_INSTANCE* inst, struct ll_entry* entry, struct SVC_ROUTE_HDR* hdr) { struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, hdr->group_id); @@ -1093,6 +1140,9 @@ static void router_forward_transit(struct UTUN_INSTANCE* inst, struct ll_entry* queue_waiter_wait(next->send_input_q, &tq->waiter, transit_queue_drain_cb, tq); } +// Детект рестарта пира по START-пакету и эпохе reset_id (аналог etcp_conn_apply_peer_reset_id): +// master держит свою эпоху и шлёт RST, slave принимает эпоху мастера и ресетит локальное состояние. +// Возвращает -1 при ошибке пересоздания conn. static int router_check_peer_restart(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CONN** prconn, struct SVC_ROUTE_HDR* hdr, size_t wire_len, struct ETCP_CONN* conn) { (void)wire_len; @@ -1133,6 +1183,8 @@ static int router_check_peer_restart(struct UTUN_INSTANCE* inst, struct ETCP_ROU return 0; } +// Принять data-пакет: зафиксировать recv_conn, синхронизировать rx_seq при первом пакете сеанса, +// отсеять out-of-bounds/дубликаты и положить пакет в incoming_q. static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, uint32_t seq, const uint8_t* wire, size_t wire_len) { rconn->recv_conn = conn; @@ -1145,7 +1197,7 @@ static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETC rconn->peer_reset_id = ((const struct SVC_ROUTE_HDR*)wire)->reset_id; } - {// SEQ out of bounds + {// SEQ за пределами окна int32_t d = (int32_t)(seq - rconn->rx_seq); if (d > (int32_t)rconn->max_inflight || d < -(int32_t)rconn->max_inflight) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router: seq=%u out of bounds, rx_seq=%u (d=%d), dropping", @@ -1154,7 +1206,7 @@ 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)) {// DUP or delta seq<=0 + if (((int32_t)(rconn->rx_seq - seq) > 0) || queue_find_data_by_index(rconn->recv_q, &seq)) {// дубликат или seq <= rx_seq // Сценарий потерянного ACK: пир ретранслит, значит не получил наш ACK — дослать его. // router_schedule_ack здесь не годится: после первого ACK rx_seq==last_sent_ack_seq и он выходит рано. uint64_t now = get_time_tb(); @@ -1183,6 +1235,8 @@ static void router_handle_data_packet(struct ETCP_ROUTER_CONN* rconn, struct ETC // etcp_router_recv_cb — обработчик ETCP_RT_ID_SVC_ROUTE // ==================================================================== +// Входная точка приёма SVC_ROUTE-кодограмм: различает транзит, RST/CLOSE/ACK (пустой payload) +// и data-пакеты, направляя их в router_handle_data_packet. static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { if (!entry) return; struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL; @@ -1248,9 +1302,11 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) } // ==================================================================== -// Public API +// Публичное API // ==================================================================== +// Инициализация роутера: создать реестр router_conns, очистить bindings и зарегистрировать +// обработчик ETCP_ID_SVC_ROUTE. int etcp_router_init(struct UTUN_INSTANCE* inst) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "NULL instance"); return -1; } memset(&inst->router_bindings, 0, sizeof(inst->router_bindings)); @@ -1270,6 +1326,7 @@ int etcp_router_init(struct UTUN_INSTANCE* inst) { return 0; } +// Деинициализация: снять обработчик SVC_ROUTE, закрыть все conn и освободить реестр. void etcp_router_destroy(struct UTUN_INSTANCE* inst) { if (!inst) return; etcp_unbind(inst, ETCP_ID_SVC_ROUTE); @@ -1286,6 +1343,7 @@ void etcp_router_destroy(struct UTUN_INSTANCE* inst) { DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "etcp_router destroyed"); } +// Зарегистрировать обработчик сервиса. callback=NULL — снять обработчик (эквивалент unbind без ошибки). int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "NULL instance"); return -1; } if (!callback) { @@ -1300,6 +1358,7 @@ int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn ca return 0; } +// Снять обработчик сервиса. Возвращает -1 если обработчик не был зарегистрирован. int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "NULL instance"); return -1; } if (!inst->router_bindings.callbacks[svc_id]) { @@ -1311,6 +1370,8 @@ int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) { return 0; } +// Отправить сервисный пакет. entry->dgram[0] = svc_id, остальное — payload. Берёт на себя +// освобождение entry в любом исходе. dst==self — loopback-доставка. int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id, struct ll_entry* entry, int force, int mode) { if (!inst || !entry) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "NULL inst=%p entry=%p", (void*)inst, (void*)entry); @@ -1337,6 +1398,8 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_ return ret; } +// Рестарт локального состояния для peer+svc в группе (перезапуск удалённой стороны): уведомить +// сервис пустой кодограммой и сбросить conn. void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id, uint8_t svc_id) { struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, group_id, remote_node_id, svc_id); if (!rconn) return; @@ -1364,6 +1427,7 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t group_id, uin router_conn_reset(rconn); } +// Backpressure: зарегистрировать waiter на send_q для (group_id, node_id, svc_id). void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id, struct queue_waiter_handle* h, queue_threshold_callback_fn callback, void* arg) { @@ -1373,6 +1437,7 @@ void etcp_router_on_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, ui queue_waiter_wait(rconn->send_q, h, callback, arg); } +// Backpressure: отменить waiter на send_q для (group_id, node_id, svc_id). void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id, struct queue_waiter_handle* h) { if (!inst || !h) return; @@ -1381,6 +1446,7 @@ void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id queue_waiter_cancel(rconn->send_q, h); } +// Backpressure: есть ли место в send_q (без регистрации waiter). 1 — есть, 0 — нет. int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id) { if (!inst) return 0; struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, group_id, node_id, svc_id); @@ -1390,9 +1456,11 @@ int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, u } // ==================================================================== -// Seq-connection API +// API seq-подключений // ==================================================================== +// Найти/создать состояние seq-подключения по (group_id, remote_node_id, svc_id): выделяет conn, +// инициализирует случайную эпоху reset_id и очереди. struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id, uint8_t svc_id) { @@ -1427,6 +1495,8 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, return rconn; } +// Отправить данные (payload без svc_id) с авто-seq и контролем inflight. mode — биты +// ROUTE_CRYPTO_SIGN / ROUTE_CRYPTO_ENCRYPT (0 = обычный пакет). int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* data, size_t len, int mode) { if (!rconn || !data || len == 0) { @@ -1437,12 +1507,13 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn, return router_enqueue_send(rconn, data, len, 0, mode); } +// Закрыть seq-подключение (нормальное закрытие с уведомлением сервиса и пира). void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) { router_close_and_notify(rconn); } -// Reset retransmit state for all router connections to a node -// (no callbacks, no close notifications — lightweight, for conn reinit) +// Сбросить ретрансмиты для всех router_conn к узлу (без коллбэков и без CLOSE — облегчённо, +// для etcp_conn_reinit, чтобы избежать гонки с очисткой ETCP-очередей). void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t remote_node_id) { if (!inst || !inst->router_conns) return; struct ll_queue* q = inst->router_conns; @@ -1477,9 +1548,10 @@ void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t rem } } +// Закрыть все router_conn для (group_id, remote_node_id) (peer умер). void etcp_router_conn_close_all_for_node(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t remote_node_id) { if (!inst || !inst->router_conns) return; - // Iterate via hash chains to avoid inf-loop from put-back in FIFO iteration + // Итерируем по hash-цепочкам, чтобы избежать бесконечного цикла из-за повторного помещения в FIFO-итерации struct ll_entry* to_close[256]; int count = 0; struct ll_queue* q = inst->router_conns; diff --git a/src/routing_layer/etcp_router.h b/src/routing_layer/etcp_router.h index 16b4e23f..2b54e844 100644 --- a/src/routing_layer/etcp_router.h +++ b/src/routing_layer/etcp_router.h @@ -19,9 +19,9 @@ extern "C" { #include #include "etcp_api.h" -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== // 1. Протокол SVC_ROUTE (wire-формат и формат доставки в сервис) -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== #pragma pack(push, 1) struct SVC_ROUTE_HDR { @@ -58,9 +58,9 @@ struct SVC_ROUTE_HDR { #define ROUTER_FLAG_SIGNED 0x08 // пакет содержит Ed25519 подпись (64 байта после payload) #define ROUTER_FLAG_CLOSE 0x02 // нормальное закрытие conn -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== // 2. Внутренние структуры и константы (etcp_router.c, control_server.c) -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== #define TRANSIT_QUEUE_HASH 256 @@ -185,9 +185,9 @@ struct ETCP_ROUTER_BINDINGS { etcp_recv_fn callbacks[SVC_ROUTE_MAX_BINDINGS]; }; -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== // 3. Публичное API (для сервисов) -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== // Инициализация: etcp_bind(ETCP_ID_SVC_ROUTE) + очистка bindings + создание router_conns int etcp_router_init(struct UTUN_INSTANCE* inst); @@ -232,9 +232,9 @@ void etcp_router_cancel_send_ready(struct UTUN_INSTANCE* inst, uint64_t group_id // Backpressure: проверить есть ли место в send_q (без регистрации waiter) int etcp_router_send_q_has_room(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t node_id, uint8_t svc_id); -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== // 4. Служебные функции (для каркаса: etcp.c, topo_group.c) -// ═══════════════════════════════════════════════════════════════ +// ==================================================================== // Сбросить ретрансмиты для всех router_conn к узлу (без колбэков, без CLOSE) // Используется при etcp_conn_reinit чтобы избежать гонки с очисткой ETCP-очередей diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c index 5e785232..7550edc7 100644 --- a/src/transport_layer/etcp.c +++ b/src/transport_layer/etcp.c @@ -295,8 +295,10 @@ void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status) { static const char* names[] = { "NEW", "UP", "DOWN", "DELETE" }; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[%s] Connection status changed to %s", etcp->log_name, (status >= 0 && status < (int)(sizeof(names)/sizeof(names[0]))) ? names[status] : "?"); + etcp->callbacks_running++; struct etcp_status_cbk_entry* cbe = etcp->instance->conn_status_cbks; while (cbe) { struct etcp_status_cbk_entry* n = cbe->next; cbe->fn(etcp, status, cbe->arg); cbe = n; } + etcp->callbacks_running--; } void etcp_cbk_fire(struct ETCP_CONN* conn, int event) { @@ -305,10 +307,10 @@ void etcp_cbk_fire(struct ETCP_CONN* conn, int event) { static const char* names[] = { "INIT", "REINIT", "UP", "DOWN", "NODE_CHANGED" }; int idx = 0, e = event; while (e >>= 1) idx++; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] callback event: %s", conn->log_name, (idx >= 0 && idx < (int)(sizeof(names)/sizeof(names[0]))) ? names[idx] : "?"); - conn->callbacks_running = 1; + conn->callbacks_running++; struct etcp_cbk_entry* cbe = conn->cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; if (cbe->event_mask & event) cbe->fn(conn, event, cbe->arg); cbe = n; } - conn->callbacks_running = 0; + conn->callbacks_running--; } static void etcp_on_up(struct ETCP_CONN* etcp) { @@ -388,6 +390,13 @@ static void etcp_connection_free_deferred(void* arg) { etcp_connection_free_resources((struct ETCP_CONN*)arg); } +/* Отложенное закрытие conn: вызывается из etcp_connection_close, когда закрытие + * запрошено внутри цепочки коллбэков (callbacks_running) — синхронное закрытие + * в этом случае освободило бы conn/линки, которые коллбэк-цепочка ещё использует. */ +static void etcp_connection_close_deferred(void* arg) { + etcp_connection_close((struct ETCP_CONN*)arg); +} + // Close connection: phase 1 detach + deferred phase 2 cleanup void etcp_connection_close(struct ETCP_CONN* etcp) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); @@ -396,9 +405,9 @@ void etcp_connection_close(struct ETCP_CONN* etcp) { callring_put(CR_CCLOSE_IN, (uintptr_t)etcp, 0); if (etcp->callbacks_running) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] FATAL: etcp_connection_close called from inside callback chain — SEGFAULTING to show backtrace", - etcp->log_name); - *(volatile int*)0 = 0; + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] close requested inside callback chain — deferring", etcp->log_name); + uasync_call_soon(etcp->instance->ua, etcp, etcp_connection_close_deferred); + return; } // === PHASE 1: detach from external world === diff --git a/src/transport_layer/etcp.h b/src/transport_layer/etcp.h index 638329bc..929f9fc4 100644 --- a/src/transport_layer/etcp.h +++ b/src/transport_layer/etcp.h @@ -202,7 +202,7 @@ struct ETCP_CONN { uint8_t tx_state; // 0 - n/a, 1 - data_wait (queues empty), 2 - link_wait (link busy) uint8_t links_up; // 0 - канал не готов для передачи, 1 - канал готов для передачи (хотя бы один линк не down) uint8_t reset_done; // 0 - рукопожатие не завершено (реинит разрешён), 1 - соединение стабильно (реинит заблокирован) - uint8_t callbacks_running; // 1 - внутри итерации колбэк-цепочек, etcp_connection_close запрещён + int callbacks_running; // счётчик вложенности колбэк-цепочек (>0 — внутри цепочки, etcp_connection_close откладывается) // Unified callback chain with event mask (init/reinit/up/down/node_changed) struct etcp_cbk_entry* cbks; diff --git a/src/transport_layer/etcp_api.c b/src/transport_layer/etcp_api.c index 75edea78..29c1bf41 100644 --- a/src/transport_layer/etcp_api.c +++ b/src/transport_layer/etcp_api.c @@ -87,11 +87,13 @@ void etcp_remove_link_status_cbk(struct UTUN_INSTANCE* inst, etcp_link_status_cb void etcp_fire_link_status_cbk(struct ETCP_LINK* link, int old_state, int old_status) { if (!link || !link->etcp || !link->etcp->instance) return; + link->etcp->callbacks_running++; struct etcp_link_status_cbk_entry* cbe = link->etcp->instance->link_status_cbks; while (cbe) { cbe->fn(link->etcp, link, old_state, old_status, cbe->arg); cbe = cbe->next; } + link->etcp->callbacks_running--; } static void etcp_status_cbk_add_chain(struct etcp_status_cbk_entry** head, etcp_conn_status_fn fn, void* arg) { diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 8aa6a019..22258bc5 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -375,7 +375,9 @@ static void sockaddr_to_key(struct sockaddr_storage* addr, uint8_t key[LINK_ADDR // find_link_index, realloc_links, insert_link, remove_link — УДАЛЕНЫ. Заменены на ll_queue с хеш-индексом. -// Вставка линка в links_queue сокета по hash-ключу адреса; при дубле адреса другого коннекта — коллизия. +// Вставка линка в links_queue сокета по hash-ключу адреса. При занятом адресе: +// другой conn с тем же pubkey — ERROR (дубль-conn); другой conn с up-линком — отказ; +// другой conn с down-линком — вытеснение (etcp_link_close); тот же conn — замена stale-записи. static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); if (!e_sock || !link || !e_sock->links_queue) return -1; @@ -388,27 +390,51 @@ static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) if (dup_lqe->link && dup_lqe->link->etcp != link->etcp) { struct ETCP_CONN* new_c = link->etcp; struct ETCP_CONN* old_c = dup_lqe->link->etcp; + + /* одинаковые ключи на двух разных conn — баг «дубль-conn»: ERROR + отказ */ + if (!memcmp(new_c->crypto_ctx.peer_public_key, old_c->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE)) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, + "insert_link_queue: SAME-PUBKEY DUPLICATE conn at %s — new=%p [%s] state=%d peer=0x%016llx / old=%p [%s] state=%d init=%d up=%d", + sockaddr_storage_to_str(&link->remote_addr).str, + (void*)new_c, new_c->log_name, new_c->state, (unsigned long long)new_c->peer_node_id, + (void*)old_c, old_c->log_name, old_c->state, old_c->initialized, old_c->links_up); + return -1; + } + uint64_t new_qkey = 0, old_qkey = 0; if (new_c->conn_queue_entry) new_qkey = ((struct conn_queue_entry*)new_c->conn_queue_entry->data)->peer_node_id; if (old_c->conn_queue_entry) old_qkey = ((struct conn_queue_entry*)old_c->conn_queue_entry->data)->peer_node_id; - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, - "!!!!!!!!!!!! LINK ADDR COLLISION !!!!!!!!!!!! addr=%s " - "new: conn=%p [%s] peer=0x%016llx state=%d init=%d up=%d srv=%d tcp=%d qkey=0x%016llx " - "old: conn=%p [%s] peer=0x%016llx state=%d init=%d up=%d srv=%d tcp=%d qkey=0x%016llx", - sockaddr_storage_to_str(&link->remote_addr).str, - (void*)new_c, new_c->log_name, (unsigned long long)new_c->peer_node_id, - new_c->state, new_c->initialized, new_c->links_up, link->is_server, link->is_tcp, - (unsigned long long)new_qkey, - (void*)old_c, old_c->log_name, (unsigned long long)old_c->peer_node_id, - old_c->state, old_c->initialized, old_c->links_up, - dup_lqe->link->is_server, dup_lqe->link->is_tcp, - (unsigned long long)old_qkey); - return -1; + if (dup_lqe->link->link_status != 0) { + /* активный линк чужого conn — отказ (активный линк сносить нельзя) */ + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, + "!!!!!!!!!!!! LINK ADDR COLLISION !!!!!!!!!!!! addr=%s " + "new: conn=%p [%s] peer=0x%016llx state=%d init=%d up=%d srv=%d tcp=%d qkey=0x%016llx " + "old: conn=%p [%s] peer=0x%016llx state=%d init=%d up=%d srv=%d tcp=%d qkey=0x%016llx", + sockaddr_storage_to_str(&link->remote_addr).str, + (void*)new_c, new_c->log_name, (unsigned long long)new_c->peer_node_id, + new_c->state, new_c->initialized, new_c->links_up, link->is_server, link->is_tcp, + (unsigned long long)new_qkey, + (void*)old_c, old_c->log_name, (unsigned long long)old_c->peer_node_id, + old_c->state, old_c->initialized, old_c->links_up, + dup_lqe->link->is_server, dup_lqe->link->is_tcp, + (unsigned long long)old_qkey); + return -1; + } + + /* разные ключи, линк чужого conn в down — вытесняем полным закрытием: + etcp_link_close снимает запись из links_queue, стреляет link_status/on_link_down + и освобождает link; линк восстановится при ретрае пира. */ + DEBUG_WARN(DEBUG_CATEGORY_ETCP, + "insert_link_queue: evict stale DOWN link old=%p [%s] peer=0x%016llx for new peer=0x%016llx at %s", + (void*)dup_lqe->link, old_c->log_name, (unsigned long long)old_c->peer_node_id, + (unsigned long long)new_c->peer_node_id, sockaddr_storage_to_str(&link->remote_addr).str); + etcp_link_close(dup_lqe->link); + } else { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "insert_link_queue: replacing stale DUP addr in [%s] old_link=%p", e_sock->name, dup_lqe->link); + if (dup_lqe->link) dup_lqe->link->link_queue_entry = NULL; + queue_remove_data(e_sock->links_queue, dup_qe); + queue_entry_free(dup_qe); } - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "insert_link_queue: replacing stale DUP addr in [%s] old_link=%p", e_sock->name, dup_lqe->link); - if (dup_lqe->link) dup_lqe->link->link_queue_entry = NULL; - queue_remove_data(e_sock->links_queue, dup_qe); - queue_entry_free(dup_qe); } struct ll_entry* qe = queue_entry_new(sizeof(struct link_queue_entry)); @@ -1276,19 +1302,17 @@ static void tcp_link_close_cb(struct stcp_link *sl, int err, void *arg) { stcp_link_close(sl); link->tcp_link = NULL; - etcp_fire_link_status_cbk(link, old_state, 0); - if (is_server) { + /* Серверный линк: etcp_link_close сам снимает link из списка conn, затем + стреляет link_status и on_link_down, затем освобождает. Здесь коллбэки + НЕ стреляем — преждевременный огонь создавал окно реентерабельности, в + котором реентерабельный etcp_connection_close успевал освободить link, + а мы после этого закрывали его повторно (use-after-free). */ DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP server link %d down err=%d, closing", etcp->log_name, local_link_id, err); - etcp_on_link_down(etcp, link); - if (etcp->state != 2) { - callring_put(CR_TCPCB_FREE, (uintptr_t)link, 0); - etcp_link_close(link); - } else { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP server link %d already closed by conn teardown, skip", - etcp->log_name, local_link_id); - } + callring_put(CR_TCPCB_FREE, (uintptr_t)link, 0); + etcp_link_close(link); } else { + etcp_fire_link_status_cbk(link, old_state, 0); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d down err=%d, scheduling reconnect", etcp->log_name, local_link_id, err); etcp_on_link_down(etcp, link); if (etcp->state != 2) etcp_tcp_link_start_reconnect(link); @@ -2287,9 +2311,9 @@ static void etcp_process_packet(struct ETCP_SOCKET* e_sock, uint8_t* data, sockaddr_storage_to_str(&addr).str, (unsigned long long)peer_id, conn->log_name, *(uint64_t*)conn->crypto_ctx.peer_public_key, *(uint64_t*)sc.peer_public_key); DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "peer key mismatch for node %016llx, firing node_changed callbacks", (unsigned long long)peer_id); - conn->callbacks_running = 1; + conn->callbacks_running++; etcp_cbk_fire(conn, ETCP_CBK_EVENT_NODE_CHANGED); - conn->callbacks_running = 0; + conn->callbacks_running--; goto ec_fr; }// коллизия - peer id совпал а ключи разные. } @@ -2337,9 +2361,9 @@ static void etcp_process_packet(struct ETCP_SOCKET* e_sock, uint8_t* data, if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] link address match but pubkey mismatch, firing node_changed", conn->log_name); DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "pubkey mismatch on reused link for node %016llx", (unsigned long long)peer_id); - conn->callbacks_running = 1; + conn->callbacks_running++; etcp_cbk_fire(conn, ETCP_CBK_EVENT_NODE_CHANGED); - conn->callbacks_running = 0; + conn->callbacks_running--; errorcode = 67; goto ec_fr; } @@ -2351,9 +2375,9 @@ static void etcp_process_packet(struct ETCP_SOCKET* e_sock, uint8_t* data, conn->log_name, sockaddr_storage_to_str(&addr).str, old_conn->log_name); DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "link address conflict for node %016llx, firing node_changed callbacks on %s", (unsigned long long)peer_id, old_conn->log_name); - old_conn->callbacks_running = 1; + old_conn->callbacks_running++; etcp_cbk_fire(old_conn, ETCP_CBK_EVENT_NODE_CHANGED); - old_conn->callbacks_running = 0; + old_conn->callbacks_running--; errorcode = 67; goto ec_fr; } else {