Browse Source

etcp_router: round-robin send + контроль заполненности send_q через waiter на send_input_q

v2
evgeny 4 weeks ago
parent
commit
5b8c725c01
  1. 83
      doc/etcp_router_arch.md
  2. 271
      src/routing_layer/etcp_router.c
  3. 7
      src/routing_layer/etcp_router.h
  4. 6
      src/routing_layer/etcp_router_doc.md
  5. 10
      tests/test_etcp_router_unit.c

83
doc/etcp_router_arch.md

@ -20,7 +20,8 @@ ETCP_ROUTER_CONN (один на пару remote_node_id + svc_id)
│ │
│ ack_timer: void* периодический 100ms │
│ idle_ack_timer: void* idle таймаут 500ms │
│ send_resume_timer: void* retry send_q 50ms │
│ send_waiter: queue_waiter_handle waiter на send_input_q │
│ watchdog_timer: void* проверка инварианта 500ms │
│ send_blocked: u8 inflight полон │
└─────────────────────────────────────────────────────┘
@ -50,40 +51,41 @@ SVC_ROUTE_HDR (22 байта)
│
▼
┌──────────────────────────────┐
│ inflight < ROUTER_MAX(256)? │
│ tx_seq - tx_acked < 256 ? │
└──────┬───────────┬───────────┘
│ YES │ NO (inflight полон)
▼ ▼
┌──────────────┐ ┌─────────────────────┐
│router_send_one│ │ router_enqueue_send │ → send_q (FIFO)
│ (прямая) │ │ (в очередь) │ send_blocked = 1
└──────┬────────┘ └─────────┬───────────┘ send_resume_timer (50ms)
│ │
▼ │
┌────────────────────┐ │
│ tx_seq++ │ │
│ hdr.dst = dst │ │
│ hdr.src = node_id │ │
│ hdr.seq = tx_seq │ │
│ hdr.svc_id = svc │ │
│ │ │
│ etcp_send(conn, e) │ │
└────────────────────┘ │
│
┌────────────────────────────┘
│ (при получении ACK от remote)
│ router_enqueue_send(rconn) │ всегда через send_q (FIFO)
│ if send_q >= 64 && !force: │ → вернуть -1 (backpressure)
│ return -1 │
└──────┬───────────────────────┘
│ копия в send_q + router_send_kick()
▼
┌─────────────────────────────────────────┐
│ router_send_kick(rconn) │
│ closed/no_route → return │
│ send_q пуст → return │
│ inflight >= limit → send_blocked=1, ret │
│ waiter уже ждёт → return │
│ нет conn → set_no_route │
│ queue_waiter_wait(send_input_q, │
│ &send_waiter, drain_cb) │
└──────┬───────────────────────────────────┘
│ send_input_q опустел (round-robin: FIFO waiter-ов)
▼
┌────────────────────┐
│ router_drain_send_q│ вызывается из ACK handler
│ │ или по send_resume_timer
│ while inflight < 256│
┌──────────────────────────────┐
│ router_send_drain_cb(rconn) │ шлём ОДИН пакет
│ e = send_q.pop() │
│ router_send_one(e)│
│ │
│ если send_q пуст: │
│ send_blocked = 0 │
└────────────────────┘
│ router_send_one_flags(e) │
│ → tx_seq++ │
│ → etcp_send(conn, entry) │
│ если send_q ещё непуст: │
│ inflight полон → blocked=1│
│ иначе → queue_waiter_wait │ (в хвост → round-robin)
│ если пуст → send_blocked=0 │
└──────────────────────────────┘
Повторный запуск drain:
• приход ACK → router_handle_ack → kick
• маршрут появился → router_no_route_retry_cb → kick
• рост inflight-лимита → router_update_inflight_limit → kick
• страховка → router_send_watchdog_cb (500ms) → DEBUG_ERROR + kick
```
## Receive Path (приём данных)
@ -119,7 +121,7 @@ SVC_ROUTE_HDR (22 байта)
│ │ = hdr.seq││ |d| < ROUTER_MAX_INFLIGHT? │
│ │ (только ││ │
│ │ вперёд) ││ Дубликат? │
│ │ drain_q ││ rx_seq - seq > 0 ИЛИ │
│ │ kick ││ rx_seq - seq > 0 ИЛИ │
│ └──────────┘│ seq уже в recv_q? │
│ └──────┬─────────────┬──────────────────┘
│ │ OK │ out-of-bounds / dup
@ -226,7 +228,8 @@ SVC_ROUTE_HDR (22 байта)
│ DEBUG_WARN "stale ACK" │
│ │
│ if send_blocked: │
│ router_drain_send_q(rconn) │
│ send_blocked = 0 │
│ router_send_kick(rconn) │
└────────────────────────────────────────┘
```
@ -241,12 +244,12 @@ SVC_ROUTE_HDR (22 байта)
При блокировке:
данные → send_q (FIFO, без дропов)
send_blocked = 1
send_resume_timer = 50ms retry
Разблокировка:
приход ACK от remote → tx_acked обновляется
→ router_drain_send_q() разгребает send_q
или send_resume_timer (50ms) → router_drain_send_q()
→ router_send_kick() → waiter → drain_cb разгребает send_q
или маршрут появился / inflight-лимит вырос → kick
страховка: router_send_watchdog_cb (500ms) → force kick
send_q полностью разобран → send_blocked = 0
```
@ -304,8 +307,8 @@ SVC_ROUTE_HDR (22 байта)
│ idle_ack_timer ─── 500ms ─────▶ router_idle_ack_cb │
│ (досылка последнего ACK) │
│ │
│ send_resume_timer ─ 50ms ─────▶ router_send_resume │
│ (retry когда send_q не пуст) │
│ watchdog_timer ── 500ms ─────▶ router_send_watchdog│
│ (страховка инварианта drain, force kick) │
│ │
└─────────────────────────────────────────────────────┘
```

271
src/routing_layer/etcp_router.c

@ -22,6 +22,10 @@ in ---> {Q} -> [src, etcp] --> ... --> [dst,etcp] -> {asm_q} -> {buf q} ---> out
#include "../transport_layer/secure_channel.h"
#include <string.h>
// Флаги метаданных элемента send_q (data[0]): bit0=encrypted, bit1=signed
#define SENDQ_FLAG_ENCRYPTED 0x01
#define SENDQ_FLAG_SIGNED 0x02
// ====================================================================
// Внутренние forward declarations
// ====================================================================
@ -31,11 +35,14 @@ static void router_send_ack(struct ETCP_ROUTER_CONN* rconn);
static void router_schedule_ack(struct ETCP_ROUTER_CONN* rconn);
static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn);
static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, const uint8_t* payload, size_t payload_len);
static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len);
static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, uint8_t flag_bits, int is_signed);
static void router_send_to(struct UTUN_INSTANCE* inst, uint64_t group_id, 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 struct ETCP_CONN* router_send_conn(struct ETCP_ROUTER_CONN* rconn);
static void router_send_kick(struct ETCP_ROUTER_CONN* rconn);
static void router_send_drain_cb(struct ll_queue* q, void* arg);
static void router_send_watchdog_cb(void* arg);
static void router_send_watchdog_arm(struct ETCP_ROUTER_CONN* rconn);
static void router_send_watchdog_disarm(struct ETCP_ROUTER_CONN* rconn);
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);
@ -194,11 +201,23 @@ static struct secure_channel* e2e_get_ctx(struct UTUN_INSTANCE* inst, uint64_t p
}
// ====================================================================
// Отправка: router_send_one + send_q
// Отправка: router_send_one_flags + send_q (дренится через waiter на send_input_q)
// ====================================================================
static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len) {
return router_send_one_flags(rconn, payload, pl_len, 0, 0);
// Найти ETCP_CONN для отправки (учитывая indirect-посредников). NULL = нет маршрута.
static struct ETCP_CONN* router_send_conn(struct ETCP_ROUTER_CONN* rconn) {
struct UTUN_INSTANCE* inst = rconn->inst;
struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, rconn->group_id);
if (group) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, 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++) {
struct ETCP_CONN* c = topo_group_find_conn_for_node(group, nq->conn_mgr_intermediaries[i]);
if (c) return c;
}
}
}
return topo_group_find_conn_for_node(group, rconn->remote_node_id);
}
static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, uint8_t flag_bits, int is_signed) {
@ -228,18 +247,7 @@ static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t*
hdr->timestamp = get_current_timestamp();
if (pl_len > 0) memcpy(dgram + SVC_ROUTE_HDR_SIZE, payload, pl_len);
struct ETCP_CONN* conn = NULL;
struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, rconn->group_id);
if (group) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, 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 = topo_group_find_conn_for_node(group, nq->conn_mgr_intermediaries[i]);
if (conn) break;
}
}
}
if (!conn) conn = topo_group_find_conn_for_node(group, rconn->remote_node_id);
struct ETCP_CONN* conn = router_send_conn(rconn);
if (!conn) {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "no route to %016llx svc_id=%u",
(unsigned long long)rconn->remote_node_id, rconn->svc_id);
@ -374,13 +382,16 @@ static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
struct UTUN_INSTANCE* inst = rconn->inst;
rconn->closed = 1;
router_send_close_to_service(rconn);
// Отправляем CLOSE удалённой стороне
struct ETCP_CONN* conn = topo_group_find_conn_for_node(topo_groups_find(inst->topo_groups, rconn->group_id), rconn->remote_node_id);
if (conn) router_send_one_flags(rconn, NULL, 0, ROUTER_FLAG_CLOSE, 0);
// Отправляем CLOSE удалённой стороне + снимаем waiter отправки
struct ETCP_CONN* conn = router_send_conn(rconn);
if (conn) {
if (conn->send_input_q) queue_waiter_cancel(conn->send_input_q, &rconn->send_waiter);
router_send_one_flags(rconn, NULL, 0, ROUTER_FLAG_CLOSE, 0);
}
// Таймеры
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->watchdog_timer) { uasync_cancel_timeout(inst->ua, rconn->watchdog_timer); rconn->watchdog_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;
@ -415,49 +426,102 @@ static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
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));
while (1) {
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) break;
// Пометить отсутствие маршрута и взвести таймер повторных проверок (идемпотентно).
static void router_set_no_route(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->no_route) return;
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");
}
// Попытаться продвинуть отправку: зарегистрировать waiter на send_input_q, если есть что слать.
// Инвариант: send_q непуст ∧ !closed ∧ !no_route ∧ inflight свободен ⇒ waiter зарегистрирован.
static void router_send_kick(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->closed || rconn->no_route) return;
if (queue_entry_count(rconn->send_q) == 0) return;
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
rconn->send_blocked = 1;
return;
}
if (rconn->send_waiter.internal) return;
struct ETCP_CONN* conn = router_send_conn(rconn);
if (!conn) { router_set_no_route(rconn); return; }
queue_waiter_wait(conn->send_input_q, &rconn->send_waiter, router_send_drain_cb, rconn);
}
// Waiter-коллбэк: send_input_q опустел → шлём один пакет, при необходимости снова встаём в хвост (round-robin).
static void router_send_drain_cb(struct ll_queue* q, void* arg) {
(void)q;
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
if (rconn->closed || rconn->no_route) return;
if (queue_entry_count(rconn->send_q) == 0) { rconn->send_blocked = 0; router_send_watchdog_disarm(rconn); return; }
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
rconn->send_blocked = 1;
return;
}
struct ll_entry* e = queue_data_get(rconn->send_q);
if (!e) { rconn->send_blocked = 0; break; }
u_check(e, "drain:after_get", "router");
int is_enc = (e->size > 0) ? e->data[0] : 0;
uint8_t fb = is_enc ? ROUTER_FLAG_ENCRYPTED : 0;
int s_err = router_send_one_flags(rconn, e->dgram, e->len, fb, 0);
u_check(e, "drain:after_send", "router");
if (!e) { rconn->send_blocked = 0; router_send_watchdog_disarm(rconn); return; }
uint8_t m = (e->size > 0) ? e->data[0] : 0;
uint8_t fb = (m & SENDQ_FLAG_ENCRYPTED) ? ROUTER_FLAG_ENCRYPTED : 0;
int is_signed = (m & SENDQ_FLAG_SIGNED) ? 1 : 0;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "SEND_Q_DRAIN: svc_id=%u send_q=%d inflight=%d sq_in=%d",
rconn->svc_id, queue_entry_count(rconn->send_q),
(int32_t)(rconn->tx_seq - rconn->tx_acked), rconn->send_q->count);
router_send_one_flags(rconn, e->dgram, e->len, fb, is_signed);
queue_dgram_free(e); queue_entry_free(e);
if (s_err != 0 && (int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) break;
}
if (rconn->send_blocked && !rconn->send_resume_timer)
rconn->send_resume_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_SEND_RESUME_TB, rconn, router_send_resume_cb, "router_send_resume");
if (!rconn->send_blocked && rconn->send_resume_timer) {
uasync_cancel_timeout(rconn->inst->ua, rconn->send_resume_timer);
rconn->send_resume_timer = NULL;
if (queue_entry_count(rconn->send_q) == 0) { rconn->send_blocked = 0; router_send_watchdog_disarm(rconn); return; }
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
rconn->send_blocked = 1;
return;
}
struct ETCP_CONN* conn = router_send_conn(rconn);
if (!conn) { router_set_no_route(rconn); return; }
queue_waiter_wait(conn->send_input_q, &rconn->send_waiter, router_send_drain_cb, rconn);
}
static void router_send_resume_cb(void* arg) {
// DEBUG-watchdog: медленная (500мс) проверка инварианта drain. Если send_q непуст,
// канал свободен, но waiter не зарегистрирован — инвариант нарушен: логируем и форсируем kick.
static void router_send_watchdog_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
rconn->watchdog_timer = NULL;
if (rconn->closed) return;
rconn->send_resume_timer = NULL;
router_drain_send_q(rconn);
if (queue_entry_count(rconn->send_q) == 0) return;
if (!rconn->no_route
&& (int32_t)(rconn->tx_seq - rconn->tx_acked) < (int32_t)router_effective_max_inflight(rconn)
&& !rconn->send_waiter.internal && !rconn->send_waiter.call_soon_id) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router watchdog: send_q stalled svc_id=%u send_q=%d inflight=%d blocked=%d — force kick",
rconn->svc_id, queue_entry_count(rconn->send_q),
(int32_t)(rconn->tx_seq - rconn->tx_acked), rconn->send_blocked);
router_send_kick(rconn);
}
if (queue_entry_count(rconn->send_q) > 0)
rconn->watchdog_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_SEND_WATCHDOG_TB, rconn, router_send_watchdog_cb, "router_send_watchdog");
}
// Взвести watchdog пока send_q непуст (страховка от пропущенного re-arm).
static void router_send_watchdog_arm(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->watchdog_timer || queue_entry_count(rconn->send_q) == 0) return;
rconn->watchdog_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_SEND_WATCHDOG_TB, rconn, router_send_watchdog_cb, "router_send_watchdog");
}
// Снять watchdog когда send_q опустел.
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; }
}
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 = topo_group_find_conn_for_node(topo_groups_find(rconn->inst->topo_groups, rconn->group_id), rconn->remote_node_id);;
if (c) {
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router: route appeared for %016llx svc_id=%u, draining send_q=%d",
if (router_send_conn(rconn)) {
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router: route appeared for %016llx svc_id=%u, resuming 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);
router_send_kick(rconn);
return;
}
rconn->no_route_timer = uasync_set_timeout(rconn->inst->ua,
@ -484,18 +548,7 @@ static void router_retransmit_one(struct ETCP_ROUTER_CONN* rconn, struct ROUTER_
if (inf->is_encrypted) hdr->flags |= ROUTER_FLAG_ENCRYPTED;
hdr->timestamp = get_current_timestamp();
struct ETCP_CONN* conn = NULL;
struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, rconn->group_id);
if (group) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, 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 = topo_group_find_conn_for_node(group, nq->conn_mgr_intermediaries[i]);
if (conn) break;
}
}
}
if (!conn) conn = topo_group_find_conn_for_node(group, rconn->remote_node_id);
struct ETCP_CONN* conn = router_send_conn(rconn);
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);
@ -566,7 +619,8 @@ static void router_retrans_timer_cb(void* arg) {
}
}
static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, int force, int is_encrypted) {
// Флаги метаданных элемента send_q (data[0]): bit0=encrypted, bit1=signed
static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, int force, int is_encrypted, int is_signed) {
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",
rconn->svc_id, queue_entry_count(rconn->send_q));
@ -577,16 +631,14 @@ static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* pa
qe->dgram = u_malloc(pl_len);
if (!qe->dgram) { queue_entry_free(qe); return -1; }
qe->len = pl_len;
qe->data[0] = is_encrypted ? 1 : 0;
qe->data[0] = (is_encrypted ? SENDQ_FLAG_ENCRYPTED : 0) | (is_signed ? SENDQ_FLAG_SIGNED : 0);
if (pl_len > 0) memcpy(qe->dgram, payload, pl_len);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "router_send_q: queued svc_id=%u inflight=%d send_q=%d enc=%d",
rconn->svc_id, (int32_t)(rconn->tx_seq - rconn->tx_acked), queue_entry_count(rconn->send_q), is_encrypted);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "router_send_q: queued svc_id=%u inflight=%d send_q=%d enc=%d signed=%d",
rconn->svc_id, (int32_t)(rconn->tx_seq - rconn->tx_acked), queue_entry_count(rconn->send_q), is_encrypted, is_signed);
queue_data_put(rconn->send_q, qe);
u_check(qe, "enqueue_send:after_put", "router");
rconn->send_blocked = 1;
if (!rconn->send_resume_timer)
rconn->send_resume_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_SEND_RESUME_TB, rconn, router_send_resume_cb, "router_send_resume");
router_send_kick(rconn);
router_send_watchdog_arm(rconn);
return 0;
}
@ -609,7 +661,8 @@ static void router_update_inflight_limit(struct ETCP_ROUTER_CONN* rconn) {
if (!rconn->send_blocked) rconn->send_blocked = 1;
}
if (new_limit > old && (rconn->send_blocked || queue_entry_count(rconn->send_q) > 0)) {
router_drain_send_q(rconn);
rconn->send_blocked = 0;
router_send_kick(rconn);
}
}
@ -961,7 +1014,7 @@ static void router_handle_ack(struct UTUN_INSTANCE* inst, struct ETCP_ROUTER_CON
rconn->c_stale_ack++;
}
rconn->last_dgram_ts = get_current_timestamp();
if (rconn->send_blocked) router_drain_send_q(rconn);
if (rconn->send_blocked) { rconn->send_blocked = 0; router_send_kick(rconn); }
}
static void router_forward_transit(struct UTUN_INSTANCE* inst, struct ll_entry* entry,
@ -1226,32 +1279,10 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_
return 0;
}
struct ETCP_CONN* etcp_conn = topo_group_find_conn_for_node(topo_groups_find(inst->topo_groups, group_id), dst_node_id);
if (!etcp_conn) {
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, group_id, 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, 0);
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;
}
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, group_id, dst_node_id, svc_id);
if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
int eq_ret = router_enqueue_send(rconn, entry->dgram + 1, payload_len, force, 0);
queue_dgram_free(entry); queue_entry_free(entry);
return eq_ret;
}
int ret = router_send_one(rconn, entry->dgram + 1, payload_len);
int ret = router_enqueue_send(rconn, entry->dgram + 1, payload_len, force, 0, 0);
queue_dgram_free(entry); queue_entry_free(entry);
return ret;
}
@ -1288,31 +1319,10 @@ int etcp_route_send_encrypted(struct UTUN_INSTANCE* inst, uint64_t group_id, uin
}
queue_dgram_free(entry); queue_entry_free(entry);
struct ETCP_CONN* etcp_conn = topo_group_find_conn_for_node(topo_groups_find(inst->topo_groups, group_id), dst_node_id);
if (!etcp_conn) {
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, group_id, dst_node_id, svc_id);
if (!rconn) { u_free(enc_buf); return -1; }
int eq_ret = router_enqueue_send(rconn, enc_buf, enc_len, force, 1);
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");
}
u_free(enc_buf);
return eq_ret;
}
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, group_id, dst_node_id, svc_id);
if (!rconn) { u_free(enc_buf); return -1; }
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
int eq_ret = router_enqueue_send(rconn, enc_buf, enc_len, force, 1);
u_free(enc_buf);
return eq_ret;
}
int ret = router_send_one_flags(rconn, enc_buf, enc_len, ROUTER_FLAG_ENCRYPTED, 0);
int ret = router_enqueue_send(rconn, enc_buf, enc_len, force, 1, 0);
u_free(enc_buf);
return ret;
}
@ -1365,9 +1375,14 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t group_id, uin
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->watchdog_timer) { uasync_cancel_timeout(inst->ua, rconn->watchdog_timer); rconn->watchdog_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 ETCP_CONN* c = router_send_conn(rconn);
if (c && c->send_input_q) queue_waiter_cancel(c->send_input_q, &rconn->send_waiter);
memset(&rconn->send_waiter, 0, sizeof(rconn->send_waiter));
}
struct ll_entry* f;
if (rconn->send_q) {
@ -1513,7 +1528,7 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst,
rconn->ack_timer = NULL;
rconn->idle_ack_timer = NULL;
rconn->send_blocked = 0;
rconn->send_resume_timer = NULL;
rconn->watchdog_timer = NULL;
rconn->sess_id = 0;
rconn->start_sent = 0;
rconn->peer_sync_done = 0;
@ -1560,10 +1575,7 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn,
(void*)rconn, (const void*)data, len);
return -1;
}
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
return router_enqueue_send(rconn, data, len, 0, 0);
}
return router_send_one(rconn, data, len);
return router_enqueue_send(rconn, data, len, 0, 0, 0);
}
int etcp_router_conn_send_signed(struct ETCP_ROUTER_CONN* rconn,
@ -1573,11 +1585,7 @@ int etcp_router_conn_send_signed(struct ETCP_ROUTER_CONN* rconn,
(void*)rconn, (const void*)data, len);
return -1;
}
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= (int32_t)router_effective_max_inflight(rconn)) {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router_conn_send_signed: inflight full svc_id=%u", rconn->svc_id);
return -1;
}
return router_send_one_flags(rconn, data, len, 0, 1);
return router_enqueue_send(rconn, data, len, 0, 0, 1);
}
void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) {
@ -1604,6 +1612,7 @@ void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t rem
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
if (rconn->remote_node_id == remote_node_id && !rconn->closed) {
if (rconn->retrans_timer) { uasync_cancel_timeout(inst->ua, rconn->retrans_timer); rconn->retrans_timer = NULL; }
if (rconn->watchdog_timer) { uasync_cancel_timeout(inst->ua, rconn->watchdog_timer); rconn->watchdog_timer = NULL; }
if (rconn->inflight_q) {
struct ll_entry* f;
while ((f = queue_data_get(rconn->inflight_q)) != NULL) {
@ -1612,6 +1621,12 @@ void etcp_router_pause_retrans_for_node(struct UTUN_INSTANCE* inst, uint64_t rem
u_free(inf);
}
}
{
struct ETCP_CONN* c = router_send_conn(rconn);
if (c && c->send_input_q) queue_waiter_cancel(c->send_input_q, &rconn->send_waiter);
memset(&rconn->send_waiter, 0, sizeof(rconn->send_waiter));
}
rconn->send_blocked = 0;
rconn->no_ack_count = 0;
rconn->last_ack_changed_tb = get_time_tb();
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router: paused retrans for remote=%016llx svc_id=%u",

7
src/routing_layer/etcp_router.h

@ -102,8 +102,9 @@ struct ETCP_ROUTER_CONN {
void* idle_ack_timer; // idle таймер (дослать последний ack)
struct ll_queue* send_q; // очередь ожидающих отправки (inflight полон)
void* send_resume_timer; // таймер возобновления отправки
uint8_t send_blocked; // 1 = inflight полон, ждём ack/таймер
struct queue_waiter_handle send_waiter; // waiter на send_input_q (по 1 пакету, round-robin)
void* watchdog_timer; // DEBUG-watchdog: медленная проверка инварианта drain
uint8_t send_blocked; // 1 = inflight полон, ждём ack
uint8_t no_route; // 1 = нет BGP-маршрута, передача приостановлена
void* no_route_timer; // таймер 20ms проверки появления маршрута
uint8_t no_ack_count; // счётчик последовательных ретрансмиссий без ACK
@ -141,8 +142,8 @@ struct ETCP_ROUTER_CONN {
#define ROUTER_MAX_INFLIGHT 256 // макс пакетов в полёте (для контроля inflight)
#define ROUTER_ACK_INTERVAL_TB 100 // интервал ACK: 10ms в timebase (0.1ms)
#define ROUTER_ACK_IDLE_TB 5000 // idle таймаут: 500ms
#define ROUTER_SEND_RESUME_TB 500 // retry интервал send_q: 50ms
#define ROUTER_MAX_SEND_Q_PACKETS 64 // порог backpressure на send_q
#define ROUTER_SEND_WATCHDOG_TB 5000 // DEBUG-watchdog: 500ms проверка инварианта drain
// Ретрансмиты
#define ROUTER_RETRANS_TIMEOUT_TB 3000 // 300ms в timebase (0.1ms)

6
src/routing_layer/etcp_router_doc.md

@ -70,6 +70,8 @@ etcp_router_conn_send_signed(rconn, data, len); // добавляет Ed25519-
- `tx_seq`, `rx_seq`, `tx_acked` — seq-нумерация для порядка и inflight
- `recv_q` — reorder-очередь с хеш-индексом по seq (восстановление порядка)
- `send_q` — очередь ожидания при полном inflight (backpressure)
- `send_waiter` — waiter на send_input_q (по 1 пакету, round-robin между сервисами)
- `watchdog_timer` — медленная проверка инварианта drain (страховка от заклинивания)
- `inflight_q` — копии отправленных пакетов для ретрансмита (хеш по seq)
- `incoming_q` — FIFO между сетевым приёмом и recv_q (защита от гонок)
- `rtt`, `rtt_jitter`, `minrtt` — измерения задержки
@ -92,7 +94,8 @@ etcp_router_conn_send_signed(rconn, data, len); // добавляет Ed25519-
router_try_assembly → deliver (rx_seq++)
Отправка: router_enqueue_send → send_q (FIFO, backpressure)
router_drain_send_q → router_send_one → etcp_send + inflight_q (копия)
router_send_kick → waiter на send_input_q → router_send_drain_cb (по 1 пакету, round-robin)
→ router_send_one_flags → etcp_send + inflight_q (копия)
↓
router_track_inflight_state + retrans_schedule
ACK: периодический (10ms) + idle (500ms) → router_send_ack(rx_seq)
@ -107,6 +110,7 @@ ACK: периодический (10ms) + idle (500ms) → router_send_ack(rx_seq
| `ROUTER_RETRANS_TIMEOUT_TB` | 3000 | Таймаут ретрансмита 300ms |
| `ROUTER_NO_ACK_MAX_RETRANS` | 17 | Макс. ретрансмитов без ACK (≈5s) |
| `ROUTER_MAX_SEND_Q_PACKETS` | 64 | Порог backpressure send_q |
| `ROUTER_SEND_WATCHDOG_TB` | 5000 | Watchdog проверка инварианта drain (500ms) |
### Функции
| Функция | Описание |

10
tests/test_etcp_router_unit.c

@ -573,12 +573,12 @@ static int test_send_q(void) {
int ret = etcp_router_conn_send(c, data, sizeof(data));
if (ret != 0) FAIL("conn_send should return 0 (queued in send_q)");
if (c->send_blocked == 0) FAIL("send_blocked not set");
if (c->send_resume_timer == NULL) FAIL("send_resume_timer not set");
if (c->watchdog_timer == NULL) FAIL("watchdog_timer not set");
if (queue_entry_count(c->send_q) != 1) FAIL("send_q count != 1");
if (c->tx_seq != ROUTER_MAX_INFLIGHT) FAIL("tx_seq advanced while queued");
// Cancel timers manually
if (c->send_resume_timer) { uasync_cancel_timeout(ua, c->send_resume_timer); c->send_resume_timer = NULL; }
if (c->watchdog_timer) { uasync_cancel_timeout(ua, c->watchdog_timer); c->watchdog_timer = NULL; }
etcp_router_conn_close(c);
TEARDOWN();
@ -633,7 +633,7 @@ static int test_timestamp(void) {
}
static int test_conn_send_no_bgp(void) {
TEST("conn_send — no BGP returns error");
TEST("conn_send — no BGP buffers with no_route (lossless)");
struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb;
SETUP();
struct ETCP_ROUTER_CONN* c = etcp_router_conn_get(&inst, TOPO_GROUP_UTUN, TEST_REMOTE_NODE, TEST_SVC_ID);
@ -641,7 +641,9 @@ static int test_conn_send_no_bgp(void) {
uint8_t data[4] = { 1, 2, 3, 4 };
int ret = etcp_router_conn_send(c, data, sizeof(data));
if (ret == 0) FAIL("conn_send should fail without BGP");
if (ret != 0) FAIL("conn_send should buffer (return 0) without BGP");
if (c->no_route == 0) FAIL("no_route not set");
if (queue_entry_count(c->send_q) != 1) FAIL("send_q count != 1");
etcp_router_conn_close(c);
TEARDOWN();

Loading…
Cancel
Save