From ac3cdd77390f4b4e3e282b04acacaa92d4293cfb Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sun, 31 May 2026 21:59:46 +0300 Subject: [PATCH] =?UTF-8?q?etcp=5Frouter:=20simplified=20TCP=20=E2=80=94?= =?UTF-8?q?=20reorder=20recv=5Fq,=20dedup,=20ACK=20timers,=20inflight=20co?= =?UTF-8?q?ntrol,=20per-(remote,svc)=20conn=20state?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - New SVC_ROUTE_HDR (packed struct, 22 bytes): cmd+dst+src+seq+svc_id - ETCP_ROUTER_CONN: state per (remote_node_id, svc_id), hash-indexed in router_conns - Reorder via recv_q (hash by seq), dedup with 32-bit circular compare - Periodic ACK (100ms), idle ACK (500ms), inflight limit via tx_seq - rx_acked - Legacy mode: no conn → direct delivery without reorder - etcp_router_conn_get/send/close API for seq-managed connections - Unit test test_etcp_router_unit: 18 tests, 3ms, no ETCP/sockets/BGP --- src/etcp_router.c | 426 ++++++++++++++++++++++---- src/etcp_router.h | 71 ++++- src/utun_instance.h | 3 +- tests/Makefile.am | 4 + tests/test_etcp_router_unit.c | 560 ++++++++++++++++++++++++++++++++++ 5 files changed, 996 insertions(+), 68 deletions(-) create mode 100644 tests/test_etcp_router_unit.c diff --git a/src/etcp_router.c b/src/etcp_router.c index f11e5bd6..fce3d870 100644 --- a/src/etcp_router.c +++ b/src/etcp_router.c @@ -1,5 +1,7 @@ // etcp_router.c — Сервисный слой маршрутизации поверх ETCP -// Маршрутизирует сервисные пакеты до целевой ноды, на промежуточных нодах ретранслирует +// Упрощённый TCP: восстановление порядка (recv_q), дедупликация, без переповторов +// ACK: периодическая отправка rx_seq (не чаще 100ms), idle-таймер для последнего seq +// Inflight контроль: tx_seq - rx_acked < ROUTER_MAX_INFLIGHT #include "etcp_router.h" #include "etcp.h" #include "utun_instance.h" @@ -7,56 +9,250 @@ #include "../lib/debug_config.h" #include "../lib/mem.h" #include "../lib/ll_queue.h" +#include "../lib/u_async.h" #include -// Обработчик ETCP_ID_SVC_ROUTE — вызывается etcp_int_recv на каждом узле -static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!conn || !entry) { - if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } +// ==================================================================== +// Внутренние forward declarations +// ==================================================================== +static void router_ack_timer_cb(void* arg); +static void router_idle_ack_timer_cb(void* arg); +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); +static void router_deliver(struct ETCP_ROUTER_CONN* rconn, + const uint8_t* payload, size_t payload_len); + +// ==================================================================== +// Управление ROUTER_CONN +// ==================================================================== + +static struct ETCP_ROUTER_CONN* router_conn_find(struct UTUN_INSTANCE* inst, + uint64_t remote_node_id, uint8_t svc_id) { + if (!inst || !inst->router_conns) return NULL; + uint8_t key[9]; + memcpy(key, &remote_node_id, 8); + key[8] = svc_id; + return (struct ETCP_ROUTER_CONN*)queue_find_data_by_index(inst->router_conns, key); +} + +// ==================================================================== +// ACK timers +// ==================================================================== + +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); + rconn->idle_ack_timer = NULL; + } + if (!rconn->ack_timer) + rconn->ack_timer = uasync_set_timeout(rconn->inst->ua, + ROUTER_ACK_INTERVAL_TB, rconn, router_ack_timer_cb, "router_ack"); +} + +static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) { + struct UTUN_INSTANCE* inst = rconn->inst; + struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, rconn->remote_node_id); + if (!conn) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "router_ack: no route to %016llx svc_id=%u", + (unsigned long long)rconn->remote_node_id, rconn->svc_id); return; } - struct UTUN_INSTANCE* inst = conn->instance; + struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)u_malloc(SVC_ROUTE_HDR_SIZE); + if (!hdr) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_ack: u_malloc failed"); return; } + hdr->cmd = ETCP_ID_SVC_ROUTE; + hdr->dst_node_id = rconn->remote_node_id; + hdr->src_node_id = inst->node_id; + hdr->seq = rconn->rx_seq; + hdr->svc_id = rconn->svc_id; + + struct ll_entry* entry = queue_entry_new(0); + if (!entry) { u_free(hdr); return; } + entry->dgram = (uint8_t*)hdr; + entry->len = SVC_ROUTE_HDR_SIZE; + + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "router_ack: → %016llx svc_id=%u rx_seq=%u", + (unsigned long long)rconn->remote_node_id, rconn->svc_id, rconn->rx_seq); + etcp_send(conn, entry); +} + +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) { + router_send_ack(rconn); + rconn->rx_acked = rconn->rx_seq; + } + // перезапускаем если за время ожидания пришли новые данные + if (rconn->rx_seq != rconn->rx_acked) + rconn->ack_timer = uasync_set_timeout(rconn->inst->ua, + ROUTER_ACK_INTERVAL_TB, rconn, router_ack_timer_cb, "router_ack"); + else + rconn->idle_ack_timer = uasync_set_timeout(rconn->inst->ua, + ROUTER_ACK_IDLE_TB, rconn, router_idle_ack_timer_cb, "router_idle_ack"); +} + +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) { + router_send_ack(rconn); + rconn->rx_acked = rconn->rx_seq; + } +} + +// ==================================================================== +// Reorder / assembly (аналог etcp_output_try_assembly) +// ==================================================================== + +static void router_deliver(struct ETCP_ROUTER_CONN* rconn, + const uint8_t* payload, size_t payload_len) { + struct UTUN_INSTANCE* inst = rconn->inst; + etcp_recv_fn cb = inst->router_bindings.callbacks[rconn->svc_id]; + if (!cb) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "router_deliver: no handler for svc_id=%u", rconn->svc_id); + return; + } + size_t entry_len = 1 + payload_len; + struct ll_entry* svc_entry = queue_entry_new(0); + if (!svc_entry) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_deliver: queue_entry_new failed"); return; } + svc_entry->len = entry_len; + svc_entry->dgram = u_malloc(entry_len); + if (!svc_entry->dgram) { queue_entry_free(svc_entry); return; } + svc_entry->dgram[0] = rconn->svc_id; + if (payload_len > 0) memcpy(svc_entry->dgram + 1, payload, payload_len); + + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "router_deliver: svc_id=%u len=%zu from remote=%016llx", + rconn->svc_id, payload_len, (unsigned long long)rconn->remote_node_id); + cb(NULL, svc_entry); +} + +static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn) { + uint32_t next = rconn->rx_seq; + while (1) { + struct ll_entry* entry = queue_find_data_by_index(rconn->recv_q, &next); + if (!entry) break; + queue_remove_data(rconn->recv_q, entry); + router_deliver(rconn, entry->dgram, entry->len); + rconn->rx_seq = next + 1; + next++; + queue_dgram_free(entry); + queue_entry_free(entry); + } +} + +// ==================================================================== +// etcp_router_recv_cb — обработчик ETCP_ID_SVC_ROUTE +// ==================================================================== + +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; if (!inst || entry->len < SVC_ROUTE_HDR_SIZE) { - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: invalid packet inst=%p len=%zu min=%d", - (void*)inst, entry ? entry->len : 0, SVC_ROUTE_HDR_SIZE); - queue_dgram_free(entry); queue_entry_free(entry); + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: invalid packet inst=%p len=%zu min=%u", + (void*)inst, entry->len, (unsigned)SVC_ROUTE_HDR_SIZE); + if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - uint8_t svc_id = entry->dgram[1]; - uint64_t dst_node_id = 0; memcpy(&dst_node_id, entry->dgram + 2, 8); - if (dst_node_id == inst->node_id) { - if (!inst->router_bindings.callbacks[svc_id]) { - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no handler for svc_id=%u dst=%016llx self=%016llx", - svc_id, (unsigned long long)dst_node_id, (unsigned long long)inst->node_id); + struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)entry->dgram; + size_t pl_len = entry->len - SVC_ROUTE_HDR_SIZE; + uint8_t* pl = entry->dgram + SVC_ROUTE_HDR_SIZE; + + if (hdr->dst_node_id == inst->node_id) { + // ========== Мы — целевая нода ========== + struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, hdr->src_node_id, hdr->svc_id); + + if (pl_len == 0) { + // ACK-пакет: обновляем rx_acked и выходим + if (rconn) { + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "router_ack: recv from %016llx svc_id=%u rx_seq=%u", + (unsigned long long)hdr->src_node_id, hdr->svc_id, hdr->seq); + rconn->rx_acked = hdr->seq; + rconn->last_dgram_ts = get_current_timestamp(); + } + queue_dgram_free(entry); queue_entry_free(entry); + return; + } + + // Data-пакет + if (!rconn) { + // Нет seq-состояния — legacy-доставка напрямую + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: legacy deliver svc_id=%u from %016llx len=%zu", + hdr->svc_id, (unsigned long long)hdr->src_node_id, pl_len); + etcp_recv_fn cb = inst->router_bindings.callbacks[hdr->svc_id]; + if (cb) { + struct ll_entry* svc_entry = queue_entry_new(0); + if (!svc_entry) { queue_dgram_free(entry); queue_entry_free(entry); return; } + svc_entry->len = 1 + pl_len; + svc_entry->dgram = u_malloc(svc_entry->len); + if (!svc_entry->dgram) { queue_entry_free(svc_entry); queue_dgram_free(entry); queue_entry_free(entry); return; } + svc_entry->dgram[0] = hdr->svc_id; + if (pl_len > 0) memcpy(svc_entry->dgram + 1, pl, pl_len); + queue_dgram_free(entry); queue_entry_free(entry); + cb(conn, svc_entry); + } else { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no handler for svc_id=%u", hdr->svc_id); + queue_dgram_free(entry); queue_entry_free(entry); + } + return; + } + + // Seq-состояние есть — reorder + rconn->last_dgram_ts = get_current_timestamp(); + uint32_t seq = hdr->seq; + + // Проверка границ (32-bit circular, аналог ETCP MAX_INFLIGHT_SIZE) + { + int32_t d = (int32_t)(seq - rconn->rx_seq); + if (d > ROUTER_MAX_INFLIGHT || d < -ROUTER_MAX_INFLIGHT) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router: seq=%u out of bounds, rx_seq=%u (d=%d), dropping", + seq, rconn->rx_seq, d); + queue_dgram_free(entry); queue_entry_free(entry); + return; + } + } + + // Дубликат (seq уже доставлен или в recv_q) + if (((int32_t)(rconn->rx_seq - seq) > 0) || + queue_find_data_by_index(rconn->recv_q, &seq)) { + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "router: dup seq=%u rx_seq=%u, dropping", seq, rconn->rx_seq); queue_dgram_free(entry); queue_entry_free(entry); return; } - size_t payload_len = entry->len - SVC_ROUTE_HDR_SIZE; - struct ll_entry* svc_entry = queue_entry_new(0); - if (!svc_entry) { queue_dgram_free(entry); queue_entry_free(entry); return; } - svc_entry->len = 1 + payload_len; - svc_entry->dgram = u_malloc(svc_entry->len); - if (!svc_entry->dgram) { queue_entry_free(svc_entry); queue_dgram_free(entry); queue_entry_free(entry); return; } - svc_entry->dgram[0] = svc_id; - if (payload_len > 0) memcpy(svc_entry->dgram + 1, entry->dgram + SVC_ROUTE_HDR_SIZE, payload_len); - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: delivering svc_id=%u len=%zu to handler %p src=%016llx", - svc_id, payload_len, (void*)inst->router_bindings.callbacks[svc_id], - (unsigned long long)(*(uint64_t*)(entry->dgram + 10))); + + // Кладём в recv_q + struct ll_entry* qe = queue_entry_new(4); // 4 байта data[] для seq + if (!qe) { queue_dgram_free(entry); queue_entry_free(entry); return; } + *(uint32_t*)qe->data = seq; + qe->dgram = u_malloc(pl_len); + if (!qe->dgram) { queue_entry_free(qe); queue_dgram_free(entry); queue_entry_free(entry); return; } + qe->len = pl_len; + memcpy(qe->dgram, pl, pl_len); + + queue_data_put_with_index(rconn->recv_q, qe); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "router: queued seq=%u rx_seq=%u recv_q=%d from %016llx svc_id=%u", + seq, rconn->rx_seq, queue_entry_count(rconn->recv_q), + (unsigned long long)hdr->src_node_id, rconn->svc_id); + + // Если пришёл ожидаемый — собираем contiguous цепочку + if (seq == rconn->rx_seq) + router_try_assembly(rconn); + + // Запускаем/продлеваем ACK таймер + router_schedule_ack(rconn); + queue_dgram_free(entry); queue_entry_free(entry); - inst->router_bindings.callbacks[svc_id](conn, svc_entry); - if (svc_id == ETCP_ID_ICMP_PROXY) - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router: ICMP_PROXY reply delivered to handler self=%016llx", - (unsigned long long)inst->node_id); } else { - struct ETCP_CONN* next = route_bgp_find_conn_for_node(inst->bgp, dst_node_id); + // ========== Транзит ========== + struct ETCP_CONN* next = route_bgp_find_conn_for_node(inst->bgp, hdr->dst_node_id); if (!next) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no route to %016llx svc_id=%u, dropping", - (unsigned long long)dst_node_id, svc_id); + (unsigned long long)hdr->dst_node_id, hdr->svc_id); queue_dgram_free(entry); queue_entry_free(entry); return; } DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: forwarding svc_id=%u → %016llx via %s", - svc_id, (unsigned long long)dst_node_id, next->log_name); + hdr->svc_id, (unsigned long long)hdr->dst_node_id, next->log_name); etcp_send(next, entry); } } @@ -68,8 +264,18 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) int etcp_router_init(struct UTUN_INSTANCE* inst) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_init: NULL instance"); return -1; } memset(&inst->router_bindings, 0, sizeof(inst->router_bindings)); + inst->router_conns = queue_new(inst->ua, + ROUTER_CONN_HASH_SIZE, 0, 9, "router_conns"); + if (!inst->router_conns) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_init: queue_new(router_conns) failed"); + return -1; + } int ret = etcp_bind(inst, ETCP_ID_SVC_ROUTE, etcp_router_recv_cb); - if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_init: etcp_bind failed, ret=%d", ret); return -1; } + if (ret != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_init: etcp_bind failed, ret=%d", ret); + queue_free(inst->router_conns); inst->router_conns = NULL; + return -1; + } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router initialized for node %016llx", (unsigned long long)inst->node_id); return 0; } @@ -77,14 +283,34 @@ int etcp_router_init(struct UTUN_INSTANCE* inst) { void etcp_router_destroy(struct UTUN_INSTANCE* inst) { if (!inst) return; etcp_unbind(inst, ETCP_ID_SVC_ROUTE); + // Очищаем router_conns + if (inst->router_conns) { + struct ll_entry* entry; + while ((entry = queue_data_get(inst->router_conns)) != NULL) { + struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry; + if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->ack_timer); + if (rconn->idle_ack_timer) uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); + if (rconn->recv_q) { + struct ll_entry* f; + while ((f = queue_data_get(rconn->recv_q)) != NULL) { + queue_dgram_free(f); queue_entry_free(f); + } + queue_free(rconn->recv_q); + } + queue_entry_free(entry); + } + queue_free(inst->router_conns); + inst->router_conns = NULL; + } memset(&inst->router_bindings, 0, sizeof(inst->router_bindings)); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router destroyed for node %016llx", (unsigned long long)inst->node_id); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router destroyed"); } int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_bind: NULL instance"); return -1; } if (!callback) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_bind: NULL callback for svc_id=%u", svc_id); return -1; } - if (inst->router_bindings.callbacks[svc_id]) DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router_bind: overwriting svc_id=%u", svc_id); + if (inst->router_bindings.callbacks[svc_id]) + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router_bind: overwriting svc_id=%u", svc_id); inst->router_bindings.callbacks[svc_id] = callback; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router_bind: svc_id=%u → cb=%p", svc_id, (void*)callback); return 0; @@ -92,7 +318,10 @@ int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn ca int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: NULL instance"); return -1; } - if (!inst->router_bindings.callbacks[svc_id]) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: svc_id=%u not bound", svc_id); return -1; } + if (!inst->router_bindings.callbacks[svc_id]) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: svc_id=%u not bound", svc_id); + return -1; + } inst->router_bindings.callbacks[svc_id] = NULL; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router_unbind: svc_id=%u", svc_id); return 0; @@ -105,14 +334,11 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_ return -1; } if (!entry->dgram || entry->len < 1) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: empty entry"); - queue_dgram_free(entry); queue_entry_free(entry); - return -1; + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: empty entry"); queue_dgram_free(entry); queue_entry_free(entry); return -1; } uint8_t svc_id = entry->dgram[0]; size_t payload_len = entry->len - 1; - // Loopback — dispatch прямо локально if (dst_node_id == inst->node_id) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_route_send: loopback svc_id=%u len=%zu", svc_id, payload_len); if (inst->router_bindings.callbacks[svc_id]) @@ -121,14 +347,16 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_ return 0; } - // Упаковываем в routing header size_t total_len = SVC_ROUTE_HDR_SIZE + payload_len; uint8_t* new_dgram = u_malloc(total_len); if (!new_dgram) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } - new_dgram[0] = ETCP_ID_SVC_ROUTE; - new_dgram[1] = svc_id; - memcpy(new_dgram + 2, &dst_node_id, 8); - memcpy(new_dgram + 10, &inst->node_id, 8); + + struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)new_dgram; + hdr->cmd = ETCP_ID_SVC_ROUTE; + hdr->dst_node_id = dst_node_id; + hdr->src_node_id = inst->node_id; + hdr->seq = 0; // legacy: без seq-трекинга + hdr->svc_id = svc_id; if (payload_len > 0) memcpy(new_dgram + SVC_ROUTE_HDR_SIZE, entry->dgram + 1, payload_len); struct ll_entry* new_entry = queue_entry_new(0); @@ -145,12 +373,8 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_ queue_entry_free(new_entry); queue_dgram_free(new_entry); return -1; } - if (svc_id == ETCP_ID_ICMP_PROXY) - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router: ICMP_PROXY reply sent %zu bytes → %016llx via %s", - total_len, (unsigned long long)dst_node_id, conn->log_name); - else - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_route_send: svc_id=%u → %016llx via %s len=%zu", - svc_id, (unsigned long long)dst_node_id, conn->log_name, total_len); + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_route_send: svc_id=%u → %016llx via %s len=%zu", + svc_id, (unsigned long long)dst_node_id, conn->log_name, total_len); return etcp_send(conn, new_entry); } @@ -160,3 +384,101 @@ int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id) { if (!conn || !conn->normalizer || !conn->normalizer->input) return -1; return queue_entry_count(conn->normalizer->input); } + +// ==================================================================== +// Seq-connection API +// ==================================================================== + +struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, + uint64_t remote_node_id, uint8_t svc_id) { + if (!inst || !inst->router_conns) return NULL; + struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, remote_node_id, svc_id); + if (rconn) return rconn; + + size_t data_size = sizeof(struct ETCP_ROUTER_CONN) - sizeof(struct ll_entry); + rconn = (struct ETCP_ROUTER_CONN*)queue_entry_new(data_size); + if (!rconn) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_conn_get: queue_entry_new failed"); return NULL; } + rconn->remote_node_id = remote_node_id; + rconn->svc_id = svc_id; + rconn->last_dgram_ts = get_current_timestamp(); + rconn->tx_seq = 0; + rconn->rx_seq = 0; + rconn->rx_acked = 0; + rconn->inst = inst; + rconn->ack_timer = NULL; + rconn->idle_ack_timer = NULL; + rconn->recv_q = queue_new(inst->ua, ROUTER_RECVQ_HASH_SIZE, 0, 4, "router_recv_q"); + if (!rconn->recv_q) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_conn_get: queue_new(recv_q) failed"); + queue_entry_free(&rconn->ll); + return NULL; + } + queue_data_put_with_index(inst->router_conns, &rconn->ll); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: new conn remote=%016llx svc_id=%u", + (unsigned long long)remote_node_id, svc_id); + return rconn; +} + +int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn, + const uint8_t* data, size_t len) { + if (!rconn || !data || len == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_conn_send: invalid args rconn=%p data=%p len=%zu", + (void*)rconn, (const void*)data, len); + return -1; + } + // Inflight контроль + uint32_t inflight = rconn->tx_seq - rconn->rx_acked; + if (inflight >= ROUTER_MAX_INFLIGHT) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "router_conn_send: inflight=%u >= max=%u, drop", + inflight, ROUTER_MAX_INFLIGHT); + return -1; + } + + struct UTUN_INSTANCE* inst = rconn->inst; + uint32_t seq = rconn->tx_seq++; + + size_t total_len = SVC_ROUTE_HDR_SIZE + len; + uint8_t* dgram = u_malloc(total_len); + if (!dgram) { rconn->tx_seq--; return -1; } + + struct SVC_ROUTE_HDR* hdr = (struct SVC_ROUTE_HDR*)dgram; + hdr->cmd = ETCP_ID_SVC_ROUTE; + hdr->dst_node_id = rconn->remote_node_id; + hdr->src_node_id = inst->node_id; + hdr->seq = seq; + hdr->svc_id = rconn->svc_id; + memcpy(dgram + SVC_ROUTE_HDR_SIZE, data, len); + + struct ll_entry* entry = queue_entry_new(0); + if (!entry) { u_free(dgram); rconn->tx_seq--; return -1; } + entry->dgram = dgram; + entry->len = total_len; + + struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, rconn->remote_node_id); + if (!conn) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "router_conn_send: no route to %016llx svc_id=%u", + (unsigned long long)rconn->remote_node_id, rconn->svc_id); + queue_entry_free(entry); queue_dgram_free(entry); rconn->tx_seq--; return -1; + } + rconn->last_dgram_ts = get_current_timestamp(); + DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "router_conn_send: seq=%u → %016llx svc_id=%u len=%zu inflight=%u", + seq, (unsigned long long)rconn->remote_node_id, rconn->svc_id, len, inflight + 1); + return etcp_send(conn, entry); +} + +void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) { + if (!rconn) return; + struct UTUN_INSTANCE* inst = rconn->inst; + if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->ack_timer); + if (rconn->idle_ack_timer) uasync_cancel_timeout(inst->ua, rconn->idle_ack_timer); + if (rconn->recv_q) { + struct ll_entry* f; + while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); } + queue_free(rconn->recv_q); + } + if (inst->router_conns) + queue_remove_data(inst->router_conns, &rconn->ll); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: closed remote=%016llx svc_id=%u", + (unsigned long long)rconn->remote_node_id, rconn->svc_id); + queue_entry_free(&rconn->ll); +} diff --git a/src/etcp_router.h b/src/etcp_router.h index dff7c6c3..56fc8308 100644 --- a/src/etcp_router.h +++ b/src/etcp_router.h @@ -1,41 +1,82 @@ // etcp_router.h — Сервисный слой маршрутизации поверх ETCP -// Позволяет отправлять сервисные пакеты конкретной ноде по node_id с многошаговой маршрутизацией -// Формат transit-пакета: [ETCP_ID_SVC_ROUTE:1] [svc_id:1] [dst_node_id:8] [src_node_id:8] [payload...] +// Маршрутизирует сервисные пакеты до целевой ноды, на промежуточных нодах ретранслирует +// Упрощённый TCP поверх ETCP: восстановление порядка, дедупликация, без переповторов +// +// Формат SVC_ROUTE пакета: +// [cmd:1] [dst_node_id:8] [src_node_id:8] [seq:4] [svc_id:1] [payload...] +// ACK-пакет: тот же заголовок, seq=rx_seq, payload_len=0 #ifndef ETCP_ROUTER_H #define ETCP_ROUTER_H #include #include "etcp_api.h" -#define SVC_ROUTE_HDR_SIZE 18 // ETCP_ID_SVC_ROUTE(1) + svc_id(1) + dst_node_id(8) + src_node_id(8) +#pragma pack(push, 1) +struct SVC_ROUTE_HDR { + uint8_t cmd; // ETCP_ID_SVC_ROUTE (0x03) + uint64_t dst_node_id; + uint64_t src_node_id; + uint32_t seq; // data: tx_seq; ACK: rx_seq (ожидаемый seq) + uint8_t svc_id; // идентификатор сервиса +}; +#pragma pack(pop) +#define SVC_ROUTE_HDR_SIZE sizeof(struct SVC_ROUTE_HDR) // 22 #define SVC_ROUTE_MAX_BINDINGS 256 +// Состояние одного логического подключения (remote_node_id + svc_id) +struct ETCP_ROUTER_CONN { + struct ll_entry ll; // data[0..7]=remote_node_id, data[8]=svc_id — хеш-индекс + uint64_t remote_node_id; // = data[0..7] + uint8_t svc_id; // = data[8] + uint16_t last_dgram_ts; // последний timestamp датаграммы (rx или tx, timebase 0.1ms) + + uint32_t tx_seq; // следующий seq для отправки + uint32_t rx_seq; // ожидаемый seq для сборки (next expected) + uint32_t rx_acked; // последний отправленный ACK (= rx_seq на момент отправки) + + struct UTUN_INSTANCE* inst; + struct ll_queue* recv_q; // reorder очередь: hash по seq (4 байта, offset 0) + void* ack_timer; // периодический таймер (100ms) + void* idle_ack_timer; // idle таймер (дослать последний ack) +}; + +#define ROUTER_CONN_HASH_SIZE 256 +#define ROUTER_RECVQ_HASH_SIZE 1024 +#define ROUTER_MAX_INFLIGHT 256 // макс пакетов в полёте (для контроля inflight) +#define ROUTER_ACK_INTERVAL_TB 1000 // интервал ACK: 100ms в timebase (0.1ms) +#define ROUTER_ACK_IDLE_TB 5000 // idle таймаут: 500ms + // Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS) struct ETCP_ROUTER_BINDINGS { etcp_recv_fn callbacks[SVC_ROUTE_MAX_BINDINGS]; }; -// Инициализация: etcp_bind(ETCP_ID_SVC_ROUTE) + очистка bindings +// Инициализация: etcp_bind(ETCP_ID_SVC_ROUTE) + очистка bindings + создание router_conns int etcp_router_init(struct UTUN_INSTANCE* inst); -// Деинициализация: etcp_unbind(ETCP_ID_SVC_ROUTE) +// Деинициализация void etcp_router_destroy(struct UTUN_INSTANCE* inst); -// Зарегистрировать обработчик сервиса. -// Когда пакет достигает целевой ноды, вызывается callback(conn, entry) -// entry содержит: [svc_id:1] [payload...] — как отправлено через etcp_route_send +// Зарегистрировать обработчик сервиса (legacy, без seq) int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn callback); - -// Отписаться от сервиса int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id); -// Отправить сервисный пакет узлу по node_id через overlay. -// entry формат: [svc_id:1] [payload...] -// Функция заворачивает в routing header, находит next hop через BGP, отправляет. -// Принимает ownership entry. +// Отправить сервисный пакет (legacy, seq=0, без контроля inflight) int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry); -// Возвращает количество пакетов в очереди normalizer->input для узла node_id, -1 если нет соединения +// Найти/создать состояние seq-подключения по (remote_node_id, svc_id) +struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, + uint64_t remote_node_id, uint8_t svc_id); + +// Отправить данные с авто-seq и контролем inflight +// data: payload без svc_id, flags: битовые флаги (зарезервировано) +int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn, + const uint8_t* data, size_t len); + +// Закрыть seq-подключение +void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn); + +// Возвращает количество пакетов в очереди normalizer->input для узла node_id int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id); #endif // ETCP_ROUTER_H diff --git a/src/utun_instance.h b/src/utun_instance.h index 3984797b..b6986c7f 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -100,8 +100,9 @@ struct UTUN_INSTANCE { // TCP proxy (optional, NULL if not enabled) struct tcp_proxy* tcp_proxy; - // etcp_router bindings (per-instance service routing) + // etcp_router bindings и seq-connections (per-instance service routing) struct ETCP_ROUTER_BINDINGS router_bindings; + struct ll_queue* router_conns; // Remote proxy (exit node) struct remote_proxy_ctx remote_proxy; diff --git a/tests/Makefile.am b/tests/Makefile.am index 38bb0e3a..d64cfeb1 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -34,6 +34,7 @@ check_PROGRAMS = \ test_tcp_proxy \ test_lwip_tcp \ test_etcp_router \ + test_etcp_router_unit \ test_remote_proxy \ test_udp_proxy \ test_icmp_proxy \ @@ -210,6 +211,9 @@ test_lwip_tcp_LDADD = \ test_etcp_router_SOURCES = test_etcp_router.c test_etcp_router_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_router_unit_SOURCES = test_etcp_router_unit.c +test_etcp_router_unit_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_remote_proxy_SOURCES = test_remote_proxy.c test_remote_proxy_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_router_unit.c b/tests/test_etcp_router_unit.c new file mode 100644 index 00000000..c6385844 --- /dev/null +++ b/tests/test_etcp_router_unit.c @@ -0,0 +1,560 @@ +// test_etcp_router_unit.c — Юнит-тест etcp_router: reorder, seq, ACK, conn lifecycle +// Без ETCP-соединений, без сокетов, без BGP. +// Инжектит SVC_ROUTE пакеты напрямую через etcp_router_recv_cb из api_bindings. +#include +#include +#include +#include +#include +#include + +#include "../src/etcp.h" +#include "../src/etcp_api.h" +#include "../src/etcp_router.h" +#include "../src/utun_instance.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define TEST_SVC_ID 0x42 +#define TEST_SVC_ID2 0x43 +#define TEST_REMOTE_NODE 0xAAAA000000000001ULL +#define TEST_REMOTE_NODE2 0xAAAA000000000002ULL + +// ======================== Test handler state ======================== +static struct { + int delivered; + int errors; + uint32_t expected_seq; + uint32_t last_seq; + uint8_t marker; +} rx; + +static void test_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { + (void)conn; + if (!entry || !entry->dgram || entry->len < 2) { + if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } + rx.errors++; + return; + } + if (entry->dgram[1] != rx.marker) { rx.errors++; } else { rx.delivered++; rx.last_seq = rx.expected_seq; rx.expected_seq++; } + queue_entry_free(entry); queue_dgram_free(entry); +} + +static void rx_reset(uint8_t marker) { + memset(&rx, 0, sizeof(rx)); + rx.marker = marker; +} + +// ======================== Inject helper ======================== +static struct ETCP_CONN fake_conn; + +static void inject(etcp_recv_fn recv_cb, struct UTUN_INSTANCE* inst, + uint64_t src, uint8_t svc_id, uint32_t seq, + const uint8_t* pl, size_t pl_len, int is_ack) { + struct SVC_ROUTE_HDR hdr; + memset(&hdr, 0, sizeof(hdr)); + hdr.cmd = ETCP_ID_SVC_ROUTE; + hdr.dst_node_id = inst->node_id; + hdr.src_node_id = src; + hdr.seq = seq; + hdr.svc_id = svc_id; + + size_t total = SVC_ROUTE_HDR_SIZE + (is_ack ? 0 : pl_len); + struct ll_entry* e = queue_entry_new(0); + if (!e) { printf(" inject: queue_entry_new failed\n"); return; } + e->dgram = u_malloc(total); + if (!e->dgram) { queue_entry_free(e); return; } + memcpy(e->dgram, &hdr, SVC_ROUTE_HDR_SIZE); + if (!is_ack && pl_len > 0) + memcpy(e->dgram + SVC_ROUTE_HDR_SIZE, pl, pl_len); + e->len = total; + + recv_cb(&fake_conn, e); +} + +// ======================== Setup ======================== +#define SETUP() do { \ + ua = uasync_create(); \ + if (!ua) { printf(" SETUP: uasync_create failed\n"); return -1; } \ + memset(&inst, 0, sizeof(inst)); \ + inst.ua = ua; \ + inst.node_id = 0xBBBB000000000001ULL; \ + debug_config_init(); \ + debug_set_level(DEBUG_LEVEL_ERROR); \ + debug_set_categories(DEBUG_CATEGORY_ALL); \ + memset(&fake_conn, 0, sizeof(fake_conn)); \ + fake_conn.instance = &inst; \ + TESTASSERT(etcp_router_init(&inst) == 0); \ + TESTASSERT(etcp_router_bind(&inst, TEST_SVC_ID, test_handler) == 0); \ + recv_cb = inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE]; \ + TESTASSERT(recv_cb != NULL); \ +} while(0) + +#define TEARDOWN() do { \ + etcp_router_destroy(&inst); \ + uasync_destroy(ua, 0); \ + ua = NULL; \ +} while(0) + +// ======================== Macros ======================== +static int test_total = 0, test_passed = 0, test_failed = 0; + +#define TESTASSERT(cond) do { \ + if (!(cond)) { printf(" FAIL: %s:%d — %s\n", __FILE__, __LINE__, #cond); return -1; } \ +} while(0) + +#define TEST(name) do { \ + test_total++; \ + printf("TEST %d: %-50s ", test_total, name); fflush(stdout); \ +} while(0) + +#define PASS() do { puts("PASS"); test_passed++; } while(0) +#define FAIL(msg) do { printf("FAIL: %s\n", msg); test_failed++; return; } while(0) + +// ======================== Tests ======================== + +static int test_init_destroy(void) { + TEST("init/destroy"); + struct UASYNC* ua = uasync_create(); + if (!ua) FAIL("uasync_create"); + struct UTUN_INSTANCE inst; + memset(&inst, 0, sizeof(inst)); + inst.ua = ua; + inst.node_id = 1; + debug_config_init(); + debug_set_level(DEBUG_LEVEL_ERROR); + debug_set_categories(DEBUG_CATEGORY_ALL); + + memset(&fake_conn, 0, sizeof(fake_conn)); + fake_conn.instance = &inst; + + if (etcp_router_init(&inst) != 0) FAIL("init"); + if (inst.router_conns == NULL) FAIL("router_conns NULL after init"); + if (inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE] == NULL) FAIL("recv_cb not bound"); + + etcp_router_destroy(&inst); + if (inst.router_conns != NULL) FAIL("router_conns not NULL after destroy"); + if (inst.api_bindings.callbacks[ETCP_ID_SVC_ROUTE] != NULL) FAIL("recv_cb still bound"); + + uasync_destroy(ua, 0); + PASS(); + return 0; +} + +static int test_bind_unbind(void) { + TEST("bind/unbind API"); + struct UASYNC* ua = uasync_create(); + struct UTUN_INSTANCE inst; + memset(&inst, 0, sizeof(inst)); + inst.ua = ua; + inst.node_id = 1; + debug_config_init(); + debug_set_level(DEBUG_LEVEL_ERROR); + debug_set_categories(DEBUG_CATEGORY_ALL); + + if (etcp_router_init(&inst) != 0) FAIL("init"); + + // bind + if (etcp_router_bind(&inst, 0xEE, test_handler) != 0) FAIL("bind"); + if (inst.router_bindings.callbacks[0xEE] != test_handler) FAIL("cb not set"); + // rebind (overwrite) + if (etcp_router_bind(&inst, 0xEE, test_handler) != 0) FAIL("rebind"); + // unbind + if (etcp_router_unbind(&inst, 0xEE) != 0) FAIL("unbind"); + if (inst.router_bindings.callbacks[0xEE] != NULL) FAIL("cb not cleared"); + // double unbind + if (etcp_router_unbind(&inst, 0xEE) == 0) FAIL("double unbind should fail"); + // null args + if (etcp_router_bind(NULL, 0, test_handler) == 0) FAIL("bind null inst"); + etcp_router_destroy(&inst); + uasync_destroy(ua, 0); + PASS(); + return 0; +} + +static int test_conn_get_create(void) { + TEST("conn_get — create"); + 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); + if (!c) FAIL("conn_get returned NULL"); + if (c->remote_node_id != TEST_REMOTE_NODE) FAIL("wrong remote_node_id"); + 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->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"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_conn_get_reuse(void) { + TEST("conn_get — reuse / different"); + struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb; + SETUP(); + + struct ETCP_ROUTER_CONN* c1 = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID); + if (!c1) FAIL("c1 NULL"); + struct ETCP_ROUTER_CONN* c2 = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID); + if (c1 != c2) FAIL("reuse: different pointers"); + + struct ETCP_ROUTER_CONN* c3 = etcp_router_conn_get(&inst, TEST_REMOTE_NODE2, TEST_SVC_ID); + if (c3 == c1) FAIL("different remote, same pointer"); + struct ETCP_ROUTER_CONN* c4 = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID2); + if (c4 == c1) FAIL("different svc, same pointer"); + + etcp_router_conn_close(c1); + etcp_router_conn_close(c3); + etcp_router_conn_close(c4); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_conn_close(void) { + TEST("conn_close — new after close"); + struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb; + SETUP(); + + struct ETCP_ROUTER_CONN* c1 = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID); + if (!c1) FAIL("c1 NULL"); + c1->tx_seq = 99; // mark + etcp_router_conn_close(c1); + struct ETCP_ROUTER_CONN* c2 = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID); + if (!c2) FAIL("c2 NULL"); + if (c2->tx_seq != 0 || c2->rx_seq != 0) FAIL("not fresh after close+reopen"); + + etcp_router_conn_close(c2); + TEARDOWN(); + PASS(); + return 0; +} + +// ======================== Seq/reorder tests ======================== + +static int test_in_order(void) { + TEST("in-order delivery (seq 0,1,2)"); + 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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0xAA); + uint8_t data[] = { TEST_SVC_ID, 0xAA, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 1, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 2, data + 1, 2, 0); + + if (rx.delivered != 3) FAIL("delivered != 3"); + if (rx.errors != 0) FAIL("errors != 0"); + if (c->rx_seq != 3) FAIL("rx_seq != 3"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_out_of_order(void) { + TEST("out-of-order (0,2,3 then 1 → assembly to 3)"); + 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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0xBB); + uint8_t data[] = { TEST_SVC_ID, 0xBB, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); + if (rx.delivered != 1) FAIL("after seq 0: delivered != 1"); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 2, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 3, data + 1, 2, 0); + if (rx.delivered != 1) FAIL("after seq 2,3: still delivered != 1 (gap at 1)"); + + // Fill the gap + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 1, data + 1, 2, 0); + if (rx.delivered != 4) FAIL("after seq 1 fill: delivered != 4"); + if (rx.errors != 0) FAIL("errors != 0"); + if (c->rx_seq != 4) FAIL("rx_seq != 4"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_duplicate_same(void) { + TEST("duplicate — same seq injected twice"); + 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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0xCC); + uint8_t data[] = { TEST_SVC_ID, 0xCC, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 1, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 1, data + 1, 2, 0); // dup + + if (rx.delivered != 2) FAIL("delivered != 2 (dup dropped)"); + if (c->rx_seq != 2) FAIL("rx_seq != 2"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_duplicate_past(void) { + TEST("duplicate — past seq (already delivered)"); + 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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0xDD); + uint8_t data[] = { TEST_SVC_ID, 0xDD, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 1, data + 1, 2, 0); + if (rx.delivered != 2) FAIL("delivered != 2"); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); // past + + if (rx.delivered != 2) FAIL("delivered != 2 (past dup dropped)"); + if (c->rx_seq != 2) FAIL("rx_seq != 2"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_out_of_bounds(void) { + TEST("out of bounds — seq too far ahead"); + 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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0xEE); + uint8_t data[] = { TEST_SVC_ID, 0xEE, 0 }; + uint32_t bad_seq = ROUTER_MAX_INFLIGHT + 2; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, bad_seq, data + 1, 2, 0); + + if (rx.delivered != 0) FAIL("out-of-bounds delivered"); + if (c->rx_seq != 0) FAIL("rx_seq changed on bounds drop"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_circular_wrap(void) { + TEST("32-bit circular wrap-around"); + 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); + if (!c) FAIL("conn_get failed"); + + c->rx_seq = 0xFFFFFFFE; + rx_reset(0x11); + rx.expected_seq = 0xFFFFFFFE; + uint8_t data[] = { TEST_SVC_ID, 0x11, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0xFFFFFFFE, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0xFFFFFFFF, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 1, data + 1, 2, 0); + + if (rx.delivered != 4) FAIL("delivered != 4 (wrap)"); + if (c->rx_seq != 2) FAIL("rx_seq != 2 after wrap (wrapped)"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_data_integrity(void) { + TEST("data integrity — payload preserved"); + 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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0x77); + uint8_t payload[10] = { 0x77, 0xA1, 0xB2, 0xC3, 0xD4, 0xE5, 0xF6, 0x07, 0x18, 0x29 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, payload, 10, 0); + + if (rx.delivered != 1) FAIL("delivered != 1"); + if (rx.errors != 0) FAIL("data mismatch"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_ack_receive(void) { + TEST("ACK receive — rx_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); + if (!c) FAIL("conn_get failed"); + + rx_reset(0x00); + // 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 (rx.delivered != 0) FAIL("ACK delivered to handler (should not)"); + if (c->last_dgram_ts == 0) FAIL("last_dgram_ts not updated on ACK"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_ack_timer(void) { + TEST("ACK timer scheduled on data receive"); + 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); + if (!c) FAIL("conn_get failed"); + + if (c->ack_timer != NULL) FAIL("ack_timer set before any data"); + + rx_reset(0x55); + uint8_t data[] = { TEST_SVC_ID, 0x55, 0 }; + 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"); + + // Cancel timer manually (it would fire in real event loop) + if (c->ack_timer) uasync_cancel_timeout(ua, c->ack_timer); + c->ack_timer = NULL; + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_legacy_no_conn(void) { + TEST("legacy — no seq-conn, direct delivery"); + struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb; + SETUP(); + // Register handler but DON'T create conn for svc_id=0x50 + if (etcp_router_bind(&inst, 0x50, test_handler) != 0) FAIL("bind legacy"); + + rx_reset(0x33); + uint8_t data[] = { 0x50, 0x33, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, 0x50, 0, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, 0x50, 1, data + 1, 2, 0); + inject(recv_cb, &inst, TEST_REMOTE_NODE, 0x50, 2, data + 1, 2, 0); + + if (rx.delivered != 3) FAIL("legacy: delivered != 3"); + if (rx.errors != 0) FAIL("legacy: errors != 0"); + + etcp_router_unbind(&inst, 0x50); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_loopback(void) { + TEST("loopback — etcp_route_send to self"); + struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb; + SETUP(); + + rx_reset(0x66); + uint8_t buf[4] = { TEST_SVC_ID, 0x66, 0x00, 0x00 }; + struct ll_entry* e = queue_entry_new(0); + if (!e) FAIL("queue_entry_new"); + e->dgram = u_malloc(4); + if (!e->dgram) { queue_entry_free(e); FAIL("u_malloc"); } + memcpy(e->dgram, buf, 4); + e->len = 4; + + if (etcp_route_send(&inst, inst.node_id, e) != 0) FAIL("loopback send failed"); + if (rx.delivered != 1) FAIL("loopback not delivered"); + if (rx.errors != 0) FAIL("loopback data mismatch"); + + TEARDOWN(); + PASS(); + return 0; +} + +static int test_timestamp(void) { + TEST("last_dgram_ts updated on data"); + 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); + if (!c) FAIL("conn_get failed"); + + uint16_t before = c->last_dgram_ts; + + rx_reset(0x99); + uint8_t data[] = { TEST_SVC_ID, 0x99, 0 }; + inject(recv_cb, &inst, TEST_REMOTE_NODE, TEST_SVC_ID, 0, data + 1, 2, 0); + + if (c->last_dgram_ts == 0) FAIL("last_dgram_ts zero after data"); + // Cancel ack_timer + if (c->ack_timer) { uasync_cancel_timeout(ua, c->ack_timer); c->ack_timer = NULL; } + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +static int test_conn_send_no_bgp(void) { + TEST("conn_send — no BGP returns error"); + 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); + if (!c) FAIL("conn_get failed"); + + 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"); + + etcp_router_conn_close(c); + TEARDOWN(); + PASS(); + return 0; +} + +// ======================== Main ======================== +int main(void) { + // Run all tests + test_init_destroy(); + test_bind_unbind(); + test_conn_get_create(); + test_conn_get_reuse(); + test_conn_close(); + test_in_order(); + test_out_of_order(); + test_duplicate_same(); + test_duplicate_past(); + test_out_of_bounds(); + test_circular_wrap(); + test_data_integrity(); + test_ack_receive(); + test_ack_timer(); + test_legacy_no_conn(); + test_loopback(); + test_timestamp(); + test_conn_send_no_bgp(); + + printf("\n=== Results: %d/%d passed, %d failed ===\n", test_passed, test_total, test_failed); + return test_failed > 0 ? 1 : 0; +}