diff --git a/doc/etcp_router_arch.md b/doc/etcp_router_arch.md new file mode 100644 index 00000000..246b6606 --- /dev/null +++ b/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 │ + └──────────────────────────────────────────────────┘ +``` diff --git a/src/etcp_router.c b/src/etcp_router.c index 88b3f07e..555c7796 100644 --- a/src/etcp_router.c +++ b/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; } diff --git a/src/etcp_router.h b/src/etcp_router.h index 78cf6552..5e69fed1 100644 --- a/src/etcp_router.h +++ b/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) diff --git a/tests/test_etcp_router_unit.c b/tests/test_etcp_router_unit.c index b3d49779..bcd902a7 100644 --- a/tests/test_etcp_router_unit.c +++ b/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