Browse Source

etcp: фикс UAF teardown (двойное закрытие линка) + вытеснение down-линка при коллизии адреса

- tcp_link_close_cb (серверная ветка): убрать преждевременный огонь link_status/on_link_down,
  etcp_link_close сам снимает линк из списка и стреляет коллбэки в безопасном порядке
  (устраняет двойное закрытие линка — UAF из callring 45502)
- etcp_connection_close: при callbacks_running откладывать закрытие через uasync_call_soon
  вместо *(int*)0=0
- callbacks_running: uint8_t-флаг -> int-счётчик вложенности (++/--)
- etcp_fire_conn_status / etcp_fire_link_status_cbk: поднимать callbacks_running вокруг огня
- insert_link_queue: при коллизии адреса с чужим conn — ERROR при одинаковом pubkey,
  отказ при up-линке, вытеснение (etcp_link_close) при down-линке

Попутно: доки/комментарии в etcp_router, ACK-интервал 100мс->10мс (ROUTER_ACK_INTERVAL_TB).
proxy
evgeny 6 days ago
parent
commit
cc2d2b5f78
  1. 3
      AGENTS.md
  2. 96
      src/routing_layer/etcp_router.c
  3. 16
      src/routing_layer/etcp_router.h
  4. 19
      src/transport_layer/etcp.c
  5. 2
      src/transport_layer/etcp.h
  6. 2
      src/transport_layer/etcp_api.c
  7. 92
      src/transport_layer/etcp_connections.c

3
AGENTS.md

@ -535,6 +535,9 @@ MEMBER_SYNC=27
- Обязательно наличие подробной диагностики в функциях кода которые не сильно спамят (не часто вызываются).
- Для каждого отладочного вывода подумай какая из доступной информация будет полезна чтобы можно было наиболее завершенно оценить состояние алгоритма и состояний влияющих на алгоритм.
## Оформление кода
- в .c: перед каждой функцией - краткое но понятное описание что делает функция (кроме совсем простых)
- в .h: в начале - описание модуля - для чего предназначен, как пользоваться, нюансы. Далее структуры с комментариями, далее публичные функции с описаниями (что делает, как пользоваться, какие аргументы, возврат, нюансы работы)
## chatgui (GUI Chat Client)

96
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;

16
src/routing_layer/etcp_router.h

@ -19,9 +19,9 @@ extern "C" {
#include <stdint.h>
#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-очередей

19
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 ===

2
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;

2
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) {

92
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 {

Loading…
Cancel
Save