Browse Source

etcp_router: simplified TCP — reorder recv_q, dedup, ACK timers, inflight control, per-(remote,svc) conn state

- 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
congestion
Evgeny 4 months ago
parent
commit
ac3cdd7739
  1. 426
      src/etcp_router.c
  2. 71
      src/etcp_router.h
  3. 3
      src/utun_instance.h
  4. 4
      tests/Makefile.am
  5. 560
      tests/test_etcp_router_unit.c

426
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 <string.h>
// Обработчик 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);
}

71
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 <stdint.h>
#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

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

4
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)

560
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 <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <unistd.h>
#include <sys/time.h>
#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;
}
Loading…
Cancel
Save