Browse Source

etcp_router: fix inflight underflow by splitting rx_acked into tx_acked + last_sent_ack_seq

rx_acked was used for two conflicting purposes:
1. remote ack of our sends (inflight = tx_seq - rx_acked)
2. our last sent ACK seq (dedup: rx_seq != rx_acked)

When the ack timer fired and set rx_acked = rx_seq, it overwrote
the inflight-tracking value. If rx_seq > tx_seq, the computation
tx_seq - rx_acked underflowed (e.g. 18 - 24 = 0xFFFFFFFA),
permanently blocking router_drain_send_q and causing send_q to
grow indefinitely.

Fix:
- Split rx_acked into tx_acked (remote ack, for inflight) and
  last_sent_ack_seq (our ACK, for dedup)
- Incoming ACK handler only advances tx_acked forward (stale guard)
- Use int32_t cast on all inflight comparisons to handle stale states
- Add etcp_router architecture diagram (doc/etcp_router_arch.md)
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
7f5325ac24
  1. 324
      doc/etcp_router_arch.md
  2. 35
      src/etcp_router.c
  3. 3
      src/etcp_router.h
  4. 13
      tests/test_etcp_router_unit.c

324
doc/etcp_router_arch.md

@ -0,0 +1,324 @@
# ETCP Router Architecture
## Структуры данных
```
ETCP_ROUTER_CONN (один на пару remote_node_id + svc_id)
┌─────────────────────────────────────────────────────┐
│ ll_entry (хеш-индекс: remote_node_id[8] + svc_id[1]) │
│ remote_node_id: u64 │
│ svc_id: u8 │
│ last_dgram_ts: u16 (timestamp последней датаграммы) │
│ │
│ tx_seq: u32 следующий seq для отправки │
│ rx_seq: u32 ожидаемый seq для сборки │
│ tx_acked: u32 сколько наших пакетов подтвердил remote │
│ last_sent_ack_seq: u32 последний отправленный ACK │
│ │
│ recv_q: ll_queue* очередь reorder (хеш по seq[4]) │
│ send_q: ll_queue* очередь при переполнении inflight │
│ │
│ ack_timer: void* периодический 100ms │
│ idle_ack_timer: void* idle таймаут 500ms │
│ send_resume_timer: void* retry send_q 50ms │
│ send_blocked: u8 inflight полон │
└─────────────────────────────────────────────────────┘
SVC_ROUTE_HDR (22 байта)
┌──────┬────────────────┬────────────────┬───────┬────────┬──────────┐
│ cmd │ dst_node_id │ src_node_id │ seq │ svc_id │ payload │
│ u8 │ u64 │ u64 │ u32 │ u8 │ ... │
└──────┴────────────────┴────────────────┴───────┴────────┴──────────┘
seq в data-пакетах = tx_seq (счетчик отправленных)
seq в ACK-пакетах = rx_seq (ожидаемый следующий)
CLOSE-флаг = 0x80000000 в seq, RST-флаг = 0x40000000
```
## Send Path (отправка данных)
```
etcp_route_send(inst, dst, entry)
etcp_router_conn_send(rconn, data, len)
│
▼
┌─────────────────────┐
│ rconn = find/get │ поиск по (dst, svc_id)
│ авто-создание если │ или etcp_router_conn_get()
│ ещё нет │
└────────┬────────────┘
│
▼
┌──────────────────────────────┐
│ 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_drain_send_q│ вызывается из ACK handler
│ │ или по send_resume_timer
│ while inflight < 256│
│ e = send_q.pop() │
│ router_send_one(e)│
│ │
│ если send_q пуст: │
│ send_blocked = 0 │
└────────────────────┘
```
## Receive Path (приём данных)
```
etcp_router_recv_cb(conn, entry)
│
▼
┌──────────────────────┐
│ dst == node_id? │
└──┬───────────────┬───┘
│ YES │ NO (транзит)
▼ ▼
┌────────────────────┐ ┌──────────────────────────┐
│ Мы — целевая нода │ │ Транзит: найти next hop │
└────────┬───────────┘ │ route_bgp_find_conn_for_ │
│ │ node(bgp, hdr.dst) │
▼ │ etcp_send(next, entry) │
┌────────────────────┐ └──────────────────────────┘
│ pl_len == 0? │
└──┬───────────┬─────┘
│ YES │ NO (данные)
▼ ▼
┌─────────┐ ┌─────────────────────────┐
│ CLOSE/ │ │ rconn = find/get │
│ RST? │ │ (авто-создание если нет) │
└──┬──┬───┘ └───────────┬─────────────┘
│ │ │
│ │ NO (ACK) ▼
│ ▼ ┌───────────────────────────────────────┐
│ ┌──────────┐│ Проверка границ seq: │
│ │ tx_acked ││ d = seq - rx_seq │
│ │ = hdr.seq││ |d| < ROUTER_MAX_INFLIGHT? │
│ │ (только ││ │
│ │ вперёд) ││ Дубликат? │
│ │ drain_q ││ rx_seq - seq > 0 ИЛИ │
│ └──────────┘│ seq уже в recv_q? │
│ └──────┬─────────────┬──────────────────┘
│ │ OK │ out-of-bounds / dup
│ ▼ ▼
│ ┌──────────────┐ ┌─────────┐
│ │ recv_q.put() │ │ DROP │
│ │ хеш по seq │ └─────────┘
│ └──────┬───────┘
│ │
│ ▼
│ ┌──────────────────────┐
│ │ seq == rx_seq? │
│ └──┬───────────────┬───┘
│ │ YES │ NO
│ ▼ │
│ ┌──────────────┐ │
│ │ try_assembly │ │
│ │ while next в │ │
│ │ recv_q: │ │
│ │ deliver() │ │
│ │ rx_seq++ │ │
│ └──────┬───────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────┴──┐
│ │ router_deliver() │
│ │ вызывает callback │
│ │ сервиса (svc_id) │
│ │ entry: [svc_id|payload] │
│ └──────────────┬──────────┘
│ │
▼ ▼
┌────────────────────────────────────┐
│ router_schedule_ack(rconn) │
│ запускает ack_timer (100ms) если │
│ ещё не запущен │
└────────────────────────────────────┘
```
```
CLOSE / RST ветка:
┌────────────────────────────────────────┐
│ pl_len == 0 && seq & CLOSE/RST │
│ → router_close_and_notify(rconn) │
│ ├─ router_send_close_to_service() │
│ │ (callback с NULL conn) │
│ ├─ отправить CLOSE удалённой стороне│
│ ├─ отмена всех таймеров │
│ ├─ очистка recv_q + send_q │
│ └─ удаление из inst->router_conns │
└────────────────────────────────────────┘
```
## ACK Mechanism
```
┌─────────────────────────────────────────┐
│ router_schedule_ack(rconn) │
│ (вызывается при получении данных) │
└────────────────┬────────────────────────┘
│
┌────────────────▼────────────────────────┐
│ если idle_ack_timer активен — отменить │
│ если ack_timer не активен — запустить │
│ ROUTER_ACK_INTERVAL_TB = 1000 (100ms) │
└────────────────┬────────────────────────┘
│
▼ (через 100ms)
┌─────────────────────────────────────────┐
│ router_ack_timer_cb() │
│ │
│ if rx_seq != last_sent_ack_seq: │
│ router_send_ack(rconn) │
│ last_sent_ack_seq = rx_seq │
│ перезапустить ack_timer │
│ else: │
│ запустить idle_ack_timer (500ms) │
└────────────────┬────────────────────────┘
│
▼ (через 500ms)
┌─────────────────────────────────────────┐
│ router_idle_ack_timer_cb() │
│ │
│ if rx_seq != last_sent_ack_seq: │
│ router_send_ack(rconn) (досылка) │
│ last_sent_ack_seq = rx_seq │
└─────────────────────────────────────────┘
router_send_ack(rconn):
hdr.seq = rx_seq
hdr.svc_id = svc_id
pl_len = 0 → это ACK-пакет
etcp_send(conn, entry)
```
```
Обработка входящего ACK:
┌────────────────────────────────────────┐
│ pl_len == 0 && нет CLOSE/RST флагов │
│ │
│ if seq >= tx_acked: // только вперёд │
│ tx_acked = seq // обновить │
│ else: │
│ DEBUG_WARN "stale ACK" │
│ │
│ if send_blocked: │
│ router_drain_send_q(rconn) │
└────────────────────────────────────────┘
```
## Inflight Control
```
ROUTER_MAX_INFLIGHT = 256
Отправка возможна когда: tx_seq - tx_acked < 256
Отправка заблокирована: tx_seq - tx_acked >= 256
При блокировке:
данные → 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()
send_q полностью разобран → send_blocked = 0
```
## Полная схема потоков
```
┌──────────┐
│ Сервис │ (routing, proxy, NAT, ...)
│ (svc_id)│
└────┬─────┘
│
etcp_route_send() │ callback(svc_id, entry)
▼ ▲
┌─────────┐ send ┌─────────────────────────────────┐ recv ┌─────────┐
│ TUN │─────────▶│ ETCP ROUTER │─────────▶│ TUN │
│ (выход) │ │ │ │ (вход) │
└─────────┘ │ ┌──────────┐ ┌──────────┐ │ └─────────┘
│ │ send_q │ │ recv_q │ │
│ │ (FIFO) │ │ (хеш seq)│ │
│ └────┬─────┘ └────┬─────┘ │
│ │ │ │
│ ▼ ▼ │
│ ┌─────────────────────────┐ │
│ │ Inflight / Assembly │ │
│ │ tx_seq - tx_acked │ │
│ │ rx_seq → deliver │ │
│ └───────────┬─────────────┘ │
│ │ │
└──────────────┼──────────────────┘
│
┌──────▼──────┐
│ BGP │
│ (нахождение│
│ next hop) │
└──────┬──────┘
│
┌──────▼──────┐
│ ETCP conn │──▶ сеть
└─────────────┘
──▶ send (исходящие данные)
──▶ recv (входящие данные, callback сервису)
──▶ транзит (через BGP к следующему hop)
```
## Таймеры (один rconn)
```
┌─────────────────────────────────────────────────────┐
│ │
│ ack_timer ──────── 100ms ─────▶ router_ack_timer_cb│
│ (периодический) │
│ │
│ idle_ack_timer ─── 500ms ─────▶ router_idle_ack_cb │
│ (досылка последнего ACK) │
│ │
│ send_resume_timer ─ 50ms ─────▶ router_send_resume │
│ (retry когда send_q не пуст) │
│ │
└─────────────────────────────────────────────────────┘
```
## Флаги seq
```
┌──────────────────────────────────────────────────┐
│ CLOSE: 0x80000000 нормальное закрытие conn │
│ RST: 0x40000000 conn не найден (reset) │
│ │
│ CLOSE/RST-пакеты: seq=0 | флаг, pl_len=0 │
│ Data-пакеты: seq=tx_seq, pl_len>0 │
│ ACK-пакеты: seq=rx_seq, pl_len=0 │
└──────────────────────────────────────────────────┘
```

35
src/etcp_router.c

@ -78,9 +78,9 @@ static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t*
queue_entry_free(entry); queue_dgram_free(entry); if (!seq_flags) rconn->tx_seq--; return -1;
}
rconn->last_dgram_ts = get_current_timestamp();
DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "router_send: seq=%u → %016llx svc_id=%u len=%zu inflight=%u",
DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "router_send: seq=%u → %016llx svc_id=%u len=%zu inflight=%d",
seq, (unsigned long long)rconn->remote_node_id, rconn->svc_id, pl_len,
rconn->tx_seq - rconn->rx_acked);
(int32_t)(rconn->tx_seq - rconn->tx_acked));
return etcp_send(conn, entry);
}
@ -150,7 +150,7 @@ static void router_close_and_notify(struct ETCP_ROUTER_CONN* rconn) {
static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn) {
while (1) {
if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) break;
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= ROUTER_MAX_INFLIGHT) break;
struct ll_entry* e = queue_data_get(rconn->send_q);
if (!e) { rconn->send_blocked = 0; break; }
router_send_one(rconn, e->dgram, e->len);
@ -178,8 +178,8 @@ static void router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* p
if (!qe->dgram) { queue_entry_free(qe); return; }
qe->len = pl_len;
if (pl_len > 0) memcpy(qe->dgram, payload, pl_len);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "router_send_q: queued svc_id=%u inflight=%u send_q=%d",
rconn->svc_id, rconn->tx_seq - rconn->rx_acked, queue_entry_count(rconn->send_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCPROUTE, "router_send_q: queued svc_id=%u inflight=%d send_q=%d",
rconn->svc_id, (int32_t)(rconn->tx_seq - rconn->tx_acked), queue_entry_count(rconn->send_q));
queue_data_put(rconn->send_q, qe);
rconn->send_blocked = 1;
if (!rconn->send_resume_timer)
@ -230,11 +230,11 @@ static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) {
static void router_ack_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
rconn->ack_timer = NULL;
if (rconn->rx_seq != rconn->rx_acked) {
if (rconn->rx_seq != rconn->last_sent_ack_seq) {
router_send_ack(rconn);
rconn->rx_acked = rconn->rx_seq;
rconn->last_sent_ack_seq = rconn->rx_seq;
}
if (rconn->rx_seq != rconn->rx_acked)
if (rconn->rx_seq != rconn->last_sent_ack_seq)
rconn->ack_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_ACK_INTERVAL_TB, rconn, router_ack_timer_cb, "router_ack");
else
@ -245,9 +245,9 @@ static void router_ack_timer_cb(void* arg) {
static void router_idle_ack_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
rconn->idle_ack_timer = NULL;
if (rconn->rx_seq != rconn->rx_acked) {
if (rconn->rx_seq != rconn->last_sent_ack_seq) {
router_send_ack(rconn);
rconn->rx_acked = rconn->rx_seq;
rconn->last_sent_ack_seq = rconn->rx_seq;
}
}
@ -320,9 +320,13 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
}
if (pl_len == 0) {
// ACK-пакет: обновляем rx_acked, drain send_q, выходим
if (rconn) {
rconn->rx_acked = hdr->seq;
if ((int32_t)(hdr->seq - rconn->tx_acked) >= 0) {
rconn->tx_acked = hdr->seq;
} else {
DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router: stale ACK seq=%u tx_acked=%u from %016llx, ignoring",
hdr->seq, rconn->tx_acked, (unsigned long long)hdr->src_node_id);
}
rconn->last_dgram_ts = get_current_timestamp();
if (rconn->send_blocked) router_drain_send_q(rconn);
}
@ -497,7 +501,7 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, dst_node_id, svc_id);
if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) {
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= ROUTER_MAX_INFLIGHT) {
router_enqueue_send(rconn, entry->dgram + 1, payload_len);
queue_dgram_free(entry); queue_entry_free(entry);
return 0;
@ -550,7 +554,8 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst,
rconn->last_dgram_ts = get_current_timestamp();
rconn->tx_seq = 0;
rconn->rx_seq = 0;
rconn->rx_acked = 0;
rconn->tx_acked = 0;
rconn->last_sent_ack_seq = 0;
rconn->inst = inst;
rconn->ack_timer = NULL;
rconn->idle_ack_timer = NULL;
@ -573,7 +578,7 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn,
(void*)rconn, (const void*)data, len);
return -1;
}
if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) {
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= ROUTER_MAX_INFLIGHT) {
router_enqueue_send(rconn, data, len);
return 0;
}

3
src/etcp_router.h

@ -32,7 +32,8 @@ struct ETCP_ROUTER_CONN {
uint32_t tx_seq; // следующий seq для отправки
uint32_t rx_seq; // ожидаемый seq для сборки (next expected)
uint32_t rx_acked; // последний отправленный ACK (= rx_seq на момент отправки)
uint32_t tx_acked; // сколько наших пакетов подтвердил remote (для inflight)
uint32_t last_sent_ack_seq; // последний отправленный ACK (= rx_seq на момент отправки)
struct UTUN_INSTANCE* inst;
struct ll_queue* recv_q; // reorder очередь: hash по seq (4 байта, offset 0)

13
tests/test_etcp_router_unit.c

@ -185,7 +185,8 @@ static int test_conn_get_create(void) {
if (c->svc_id != TEST_SVC_ID) FAIL("wrong svc_id");
if (c->tx_seq != 0) FAIL("tx_seq != 0");
if (c->rx_seq != 0) FAIL("rx_seq != 0");
if (c->rx_acked != 0) FAIL("rx_acked != 0");
if (c->tx_acked != 0) FAIL("tx_acked != 0");
if (c->last_sent_ack_seq != 0) FAIL("last_sent_ack_seq != 0");
if (c->recv_q == NULL) FAIL("recv_q is NULL");
if (c->last_dgram_ts == 0) FAIL("last_dgram_ts not set");
if (c->ack_timer != NULL) FAIL("ack_timer set at creation");
@ -402,7 +403,7 @@ static int test_data_integrity(void) {
}
static int test_ack_receive(void) {
TEST("ACK receive — rx_acked updated");
TEST("ACK receive — tx_acked updated");
struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb;
SETUP();
struct ETCP_ROUTER_CONN* c = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID);
@ -412,7 +413,7 @@ static int test_ack_receive(void) {
// Send ACK: is_ack=1, seq=7 means remote has rx_seq=7
inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 7, NULL, 0, 1);
if (c->rx_acked != 7) FAIL("rx_acked not updated from ACK");
if (c->tx_acked != 7) FAIL("tx_acked not updated from ACK");
if (rx.delivered != 0) FAIL("ACK delivered to handler (should not)");
if (c->last_dgram_ts == 0) FAIL("last_dgram_ts not updated on ACK");
@ -436,7 +437,7 @@ static int test_ack_timer(void) {
inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0);
if (c->ack_timer == NULL) FAIL("ack_timer not scheduled after data");
if (c->rx_acked != 0) FAIL("rx_acked changed before timer fired");
if (c->last_sent_ack_seq != 0) FAIL("last_sent_ack_seq changed before timer fired");
// Cancel timer manually (it would fire in real event loop)
if (c->ack_timer) uasync_cancel_timeout(ua, c->ack_timer);
@ -455,9 +456,9 @@ static int test_send_q(void) {
struct ETCP_ROUTER_CONN* c = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID);
if (!c) FAIL("conn_get failed");
// Simulate full inflight: set rx_acked far behind tx_seq
// Simulate full inflight: set tx_acked far behind tx_seq
c->tx_seq = ROUTER_MAX_INFLIGHT;
c->rx_acked = 0;
c->tx_acked = 0;
if (c->send_q == NULL) FAIL("send_q is NULL");
// Send via conn_send — should queue, not drop

Loading…
Cancel
Save