From 21cfe2d381b066662a12adfcd2904d0edf190502 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sun, 31 May 2026 22:55:23 +0300 Subject: [PATCH] etcp_router: remove legacy, auto-conn for all services, send_q inflight control, remove tcp_proxy internal seq MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - All SVC_ROUTE traffic now auto-creates ETCP_ROUTER_CONN, no legacy path - etcp_route_send: auto-conn + seq + inflight check → send_q if full - Inflight control via send_q: queue instead of drop, drain on ACK, retry timer 50ms - router_send_one(): unified SVC_ROUTE wrapper with seq assignment - router_deliver/try_assembly: pass real ETCP_CONN* to handlers (NAT needs it) - tcp_proxy: remove internal 16-bit seq (send_seq/recv_last_seq/recv_init), shrink HDR 8→6 - remote_proxy: same seq removal, HDR/constants adjusted - conn init/close/destroy: manage send_q, send_resume_timer lifecycle - test_etcp_router_unit: −legacy, +send_q inflight test (18 tests) --- src/etcp_router.c | 244 ++++++++++++++++++---------------- src/etcp_router.h | 9 +- src/remote_proxy.c | 37 ++---- src/remote_proxy.h | 10 +- src/tcp_proxy.c | 42 ++---- src/tcp_proxy.h | 3 - tests/test_etcp_router_unit.c | 34 +++-- 7 files changed, 178 insertions(+), 201 deletions(-) diff --git a/src/etcp_router.c b/src/etcp_router.c index fce3d870..569a2224 100644 --- a/src/etcp_router.c +++ b/src/etcp_router.c @@ -1,7 +1,7 @@ // etcp_router.c — Сервисный слой маршрутизации поверх ETCP // Упрощённый TCP: восстановление порядка (recv_q), дедупликация, без переповторов // ACK: периодическая отправка rx_seq (не чаще 100ms), idle-таймер для последнего seq -// Inflight контроль: tx_seq - rx_acked < ROUTER_MAX_INFLIGHT +// Inflight контроль через send_q: при переполнении — очередь + retry-таймер, без дропов #include "etcp_router.h" #include "etcp.h" #include "utun_instance.h" @@ -19,9 +19,11 @@ 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); +static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn); +static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, const uint8_t* payload, size_t payload_len); +static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len); +static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn); +static void router_send_resume_cb(void* arg); // ==================================================================== // Управление ROUTER_CONN @@ -36,6 +38,81 @@ static struct ETCP_ROUTER_CONN* router_conn_find(struct UTUN_INSTANCE* inst, return (struct ETCP_ROUTER_CONN*)queue_find_data_by_index(inst->router_conns, key); } +// ==================================================================== +// Отправка: router_send_one + send_q +// ==================================================================== + +static int router_send_one(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len) { + struct UTUN_INSTANCE* inst = rconn->inst; + uint32_t seq = rconn->tx_seq++; + size_t total_len = SVC_ROUTE_HDR_SIZE + pl_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; + if (pl_len > 0) memcpy(dgram + SVC_ROUTE_HDR_SIZE, payload, pl_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_send_one: 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_send: seq=%u → %016llx svc_id=%u len=%zu inflight=%u", + seq, (unsigned long long)rconn->remote_node_id, rconn->svc_id, pl_len, + rconn->tx_seq - rconn->rx_acked); + return etcp_send(conn, entry); +} + +static void router_drain_send_q(struct ETCP_ROUTER_CONN* rconn) { + while (1) { + if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) break; + struct ll_entry* e = queue_data_get(rconn->send_q); + if (!e) { rconn->send_blocked = 0; break; } + router_send_one(rconn, e->dgram, e->len); + queue_dgram_free(e); queue_entry_free(e); + } + if (rconn->send_blocked && !rconn->send_resume_timer) + rconn->send_resume_timer = uasync_set_timeout(rconn->inst->ua, + ROUTER_SEND_RESUME_TB, rconn, router_send_resume_cb, "router_send_resume"); + if (!rconn->send_blocked && rconn->send_resume_timer) { + uasync_cancel_timeout(rconn->inst->ua, rconn->send_resume_timer); + rconn->send_resume_timer = NULL; + } +} + +static void router_send_resume_cb(void* arg) { + struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; + rconn->send_resume_timer = NULL; + router_drain_send_q(rconn); +} + +static void router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len) { + struct ll_entry* qe = queue_entry_new(0); + if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_enqueue_send: queue_entry_new failed"); return; } + qe->dgram = u_malloc(pl_len); + if (!qe->dgram) { queue_entry_free(qe); return; } + qe->len = pl_len; + if (pl_len > 0) memcpy(qe->dgram, payload, pl_len); + queue_data_put(rconn->send_q, qe); + rconn->send_blocked = 1; + if (!rconn->send_resume_timer) + rconn->send_resume_timer = uasync_set_timeout(rconn->inst->ua, + ROUTER_SEND_RESUME_TB, rconn, router_send_resume_cb, "router_send_resume"); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "router_send_q: queued svc_id=%u inflight=%u send_q=%d", + rconn->svc_id, rconn->tx_seq - rconn->rx_acked, queue_entry_count(rconn->send_q)); +} + // ==================================================================== // ACK timers // ==================================================================== @@ -83,7 +160,6 @@ static void router_ack_timer_cb(void* arg) { 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"); @@ -105,7 +181,7 @@ static void router_idle_ack_timer_cb(void* arg) { // Reorder / assembly (аналог etcp_output_try_assembly) // ==================================================================== -static void router_deliver(struct ETCP_ROUTER_CONN* rconn, +static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, 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]; @@ -124,16 +200,16 @@ static void router_deliver(struct ETCP_ROUTER_CONN* rconn, 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); + cb(conn, svc_entry); } -static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn) { +static void router_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn) { 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); + router_deliver(rconn, conn, entry->dgram, entry->len); rconn->rx_seq = next + 1; next++; queue_dgram_free(entry); @@ -160,41 +236,24 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) if (hdr->dst_node_id == inst->node_id) { // ========== Мы — целевая нода ========== + // Авто-создаём conn при первом пакете struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, hdr->src_node_id, hdr->svc_id); if (pl_len == 0) { - // ACK-пакет: обновляем rx_acked и выходим + // ACK-пакет: обновляем rx_acked, drain send_q, выходим 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(); + if (rconn->send_blocked) router_drain_send_q(rconn); } queue_dgram_free(entry); queue_entry_free(entry); return; } - // Data-пакет + // Data-пакет — авто-создаём conn если ещё нет 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; + rconn = etcp_router_conn_get(inst, hdr->src_node_id, hdr->svc_id); + if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return; } } // Seq-состояние есть — reorder @@ -221,7 +280,7 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) } // Кладём в recv_q - struct ll_entry* qe = queue_entry_new(4); // 4 байта data[] для seq + struct ll_entry* qe = queue_entry_new(4); if (!qe) { queue_dgram_free(entry); queue_entry_free(entry); return; } *(uint32_t*)qe->data = seq; qe->dgram = u_malloc(pl_len); @@ -234,13 +293,10 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) 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); + router_try_assembly(rconn, conn); - // Запускаем/продлеваем ACK таймер router_schedule_ack(rconn); - queue_dgram_free(entry); queue_entry_free(entry); } else { // ========== Транзит ========== @@ -283,20 +339,23 @@ 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->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->send_resume_timer) uasync_cancel_timeout(inst->ua, rconn->send_resume_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); - } + while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); } queue_free(rconn->recv_q); } + if (rconn->send_q) { + struct ll_entry* f; + while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); } + queue_free(rconn->send_q); + } queue_entry_free(entry); } queue_free(inst->router_conns); @@ -330,8 +389,7 @@ int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) { int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry) { if (!inst || !entry) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: NULL inst=%p entry=%p", (void*)inst, (void*)entry); - if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } - return -1; + if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } 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; @@ -347,35 +405,18 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_ return 0; } - 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; } - - 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 ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, dst_node_id, svc_id); + if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } - struct ll_entry* new_entry = queue_entry_new(0); - if (!new_entry) { u_free(new_dgram); queue_dgram_free(entry); queue_entry_free(entry); return -1; } - new_entry->dgram = new_dgram; - new_entry->len = total_len; + if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) { + router_enqueue_send(rconn, entry->dgram + 1, payload_len); + queue_dgram_free(entry); queue_entry_free(entry); + return 0; + } + int ret = router_send_one(rconn, entry->dgram + 1, payload_len); queue_dgram_free(entry); queue_entry_free(entry); - - struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, dst_node_id); - if (!conn) { - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_route_send: no route to %016llx svc_id=%u, dropping", - (unsigned long long)dst_node_id, svc_id); - queue_entry_free(new_entry); queue_dgram_free(new_entry); - return -1; - } - 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); + return ret; } int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id) { @@ -407,12 +448,12 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, rconn->inst = inst; rconn->ack_timer = NULL; rconn->idle_ack_timer = NULL; + rconn->send_blocked = 0; + rconn->send_resume_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; - } + if (!rconn->recv_q) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_conn_get: queue_new(recv_q) failed"); queue_entry_free(&rconn->ll); return NULL; } + rconn->send_q = queue_new(inst->ua, 0, 0, 0, "router_send_q"); + if (!rconn->send_q) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_conn_get: queue_new(send_q) failed"); queue_free(rconn->recv_q); 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); @@ -426,56 +467,29 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn, (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; + if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) { + router_enqueue_send(rconn, data, len); + return 0; } - 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); + return router_send_one(rconn, data, len); } 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->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->send_resume_timer) uasync_cancel_timeout(inst->ua, rconn->send_resume_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 (rconn->send_q) { + struct ll_entry* f; + while ((f = queue_data_get(rconn->send_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); } + queue_free(rconn->send_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", diff --git a/src/etcp_router.h b/src/etcp_router.h index 56fc8308..91bb8406 100644 --- a/src/etcp_router.h +++ b/src/etcp_router.h @@ -38,6 +38,10 @@ struct ETCP_ROUTER_CONN { struct ll_queue* recv_q; // reorder очередь: hash по seq (4 байта, offset 0) void* ack_timer; // периодический таймер (100ms) void* idle_ack_timer; // idle таймер (дослать последний ack) + + struct ll_queue* send_q; // очередь ожидающих отправки (inflight полон) + void* send_resume_timer; // таймер возобновления отправки + uint8_t send_blocked; // 1 = inflight полон, ждём ack/таймер }; #define ROUTER_CONN_HASH_SIZE 256 @@ -45,6 +49,7 @@ struct ETCP_ROUTER_CONN { #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 +#define ROUTER_SEND_RESUME_TB 500 // retry интервал send_q: 50ms // Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS) struct ETCP_ROUTER_BINDINGS { @@ -57,11 +62,11 @@ int etcp_router_init(struct UTUN_INSTANCE* inst); // Деинициализация void etcp_router_destroy(struct UTUN_INSTANCE* inst); -// Зарегистрировать обработчик сервиса (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); -// Отправить сервисный пакет (legacy, seq=0, без контроля inflight) +// Отправить сервисный пакет (авто-conn, seq, inflight-контроль через send_q) int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry); // Найти/создать состояние seq-подключения по (remote_node_id, svc_id) diff --git a/src/remote_proxy.c b/src/remote_proxy.c index c8ebbab5..d17aef2a 100644 --- a/src/remote_proxy.c +++ b/src/remote_proxy.c @@ -39,7 +39,7 @@ void rp_conn_free(struct remote_proxy_conn* rc); #define RP_NORMALIZER_Q_THRESHOLD 64 static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, - uint32_t sid, uint16_t seq, const uint8_t* data, size_t len) { + uint32_t sid, const uint8_t* data, size_t len) { struct ll_entry* e = queue_entry_new(0); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "rp_send_msg: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; } e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); @@ -47,7 +47,6 @@ static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = subcmd; memcpy(e->dgram + 2, &sid, 4); - memcpy(e->dgram + 6, &seq, 2); if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); e->len = TCP_PROXY_HDR_SIZE + len; return etcp_route_send(inst, dst, e); @@ -57,7 +56,7 @@ static int rp_send_connected(struct UTUN_INSTANCE* inst, uint64_t dst, uint32_t uint16_t local_port, uint8_t status) { uint8_t buf[3]; memcpy(buf, &local_port, 2); buf[2] = status; - return rp_send_msg(inst, dst, TCP_PROXY_SUBCMD_CONNECTED, sid, 0, buf, 3); + return rp_send_msg(inst, dst, TCP_PROXY_SUBCMD_CONNECTED, sid, buf, 3); } static int rp_conn_total(struct remote_proxy_conn* rc) { @@ -142,8 +141,7 @@ static void rp_sock_pause_cb(void* arg) { } if (rc->pause_buf && rc->pause_len > 0) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE_FLUSH fd=%d sid=%08x len=%zu qcnt=%d", (int)rc->sock, rc->stream_id, rc->pause_len, qcnt); - rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->send_seq, rc->pause_buf, rc->pause_len); - rc->send_seq++; + rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->pause_buf, rc->pause_len); u_free(rc->pause_buf); rc->pause_buf = NULL; rc->pause_len = 0; } rc->read_id = uasync_add_socket_t(rc->ua, rc->sock, rp_sock_read_cb, rp_sock_write_cb, rp_sock_error_cb, rc); @@ -173,14 +171,13 @@ static void rp_sock_read_cb(socket_t sock, void* arg) { return; } DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%zd total=%d", (int)rc->sock, rc->stream_id, n, rp_conn_total(rc)); - rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->send_seq, buf, (size_t)n); - rc->send_seq++; + rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, buf, (size_t)n); } else if (n == 0) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x total=%d cli_closed=%d sock_closed=%d", (int)rc->sock, rc->stream_id, rp_conn_total(rc), rc->cli_closed, rc->sock_closed); rc->sock_closed = 1; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; - if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, 0, NULL, 0); + if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0); else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x — inst is NULL, can't send CLOSE", (int)rc->sock, rc->stream_id); if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; } if (rc->cli_closed) rp_conn_free(rc); @@ -189,7 +186,7 @@ static void rp_sock_read_cb(socket_t sock, void* arg) { (int)rc->sock, rc->stream_id, errno, strerror(errno), rp_conn_total(rc)); rc->error = 1; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; - if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, 0, NULL, 0); + if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0); else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:ERROR fd=%d sid=%08x — inst is NULL, can't send ERROR", (int)rc->sock, rc->stream_id); rp_conn_free(rc); } @@ -230,7 +227,7 @@ static void rp_sock_error_cb(socket_t sock, void* arg) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:ERR fd=%d sid=%08x", (int)rc->sock, rc->stream_id); rc->error = 1; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; - if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, 0, NULL, 0); + if (inst) rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0); else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:ERR fd=%d sid=%08x — inst is NULL", (int)rc->sock, rc->stream_id); rp_conn_free(rc); } @@ -310,29 +307,11 @@ int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn, struct remote_proxy_conn* rc = remote_proxy_find_conn(ctx, stream_id); if (!rc || rc->sock == SOCKET_INVALID) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "RP handle_data drop: нет соединения для sid=%08x, шлём CLOSE", stream_id); - if (conn) rp_send_msg(inst, conn->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, stream_id, 0, NULL, 0); + if (conn) rp_send_msg(inst, conn->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, stream_id, NULL, 0); queue_dgram_free(entry); queue_entry_free(entry); return -1; } if (rc->error) { queue_dgram_free(entry); queue_entry_free(entry); return -1; } - uint16_t seq; memcpy(&seq, entry->dgram + 6, 2); - if (!rc->recv_init) { - rc->recv_last_seq = seq; rc->recv_init = 1; - } else { - uint16_t delta = seq - rc->recv_last_seq; - if (delta == 0) { - DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "remote_proxy: duplicate DATA seq=%u sid=%08x", seq, stream_id); - queue_dgram_free(entry); queue_entry_free(entry); return 0; - } - if (delta != 1) { - DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: seq gap seq=%u last=%u sid=%08x", seq, rc->recv_last_seq, stream_id); - rc->error = 1; - rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, 0, NULL, 0); - rp_conn_free(rc); - queue_dgram_free(entry); queue_entry_free(entry); return -1; - } - rc->recv_last_seq = seq; - } size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; if (rc->connected) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:SEND fd=%d sid=%08x len=%zu total=%d", diff --git a/src/remote_proxy.h b/src/remote_proxy.h index ae1a9497..ae89bc26 100644 --- a/src/remote_proxy.h +++ b/src/remote_proxy.h @@ -20,9 +20,9 @@ struct remote_proxy_ctx; #define TCP_PROXY_CONNECTED_OK 0 #define TCP_PROXY_CONNECTED_REFUSED 1 -#define TCP_PROXY_HDR_SIZE 8 // svc_id(1)+subcmd(1)+stream_id(4)+seq(2) -#define TCP_PROXY_CONNECT_HDR_SIZE 14 // HDR_SIZE + dest_ip(4)+dest_port(2) -#define TCP_PROXY_CONNECTED_HDR_SIZE 11 // HDR_SIZE + local_port(2)+status(1) +#define TCP_PROXY_HDR_SIZE 6 // svc_id(1)+subcmd(1)+stream_id(4) +#define TCP_PROXY_CONNECT_HDR_SIZE 12 // HDR_SIZE + dest_ip(4)+dest_port(2) +#define TCP_PROXY_CONNECTED_HDR_SIZE 9 // HDR_SIZE + local_port(2)+status(1) struct remote_proxy_conn { struct remote_proxy_conn* next; @@ -35,10 +35,6 @@ struct remote_proxy_conn { uint8_t dest_ip[4]; uint16_t dest_port; - uint16_t send_seq; // 16-bit circular - uint16_t recv_last_seq; - uint8_t recv_init; - uint8_t cli_closed; // client sent CLOSE uint8_t sock_closed; // socket EOF received uint8_t error; diff --git a/src/tcp_proxy.c b/src/tcp_proxy.c index 59f1ed89..26b73483 100644 --- a/src/tcp_proxy.c +++ b/src/tcp_proxy.c @@ -41,7 +41,7 @@ static err_t tcp_output_cb(void *arg, struct pbuf *p, uint32_t src_ip, uint32_t static void proxy_feed_from_transport(struct proxy_conn *pc); static struct ll_entry* entry_from_data(struct memory_pool* pool, const uint8_t* data, uint16_t len); static void proxy_conn_free(struct proxy_conn *pc); -static int proxy_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, uint16_t seq, const uint8_t* data, size_t len); +static int proxy_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len); static int send_data(struct proxy_conn* pc, const uint8_t* data, uint16_t len); // ==================================================================== @@ -63,7 +63,7 @@ static struct ll_entry* entry_from_data(struct memory_pool* pool, const uint8_t* // Протокол: сборка и отправка сообщений прокси через ETCP // ==================================================================== static int proxy_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, - uint32_t sid, uint16_t seq, const uint8_t* data, size_t len) { + uint32_t sid, const uint8_t* data, size_t len) { struct ll_entry* e = queue_entry_new(0); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "proxy_send_msg: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; } e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len); @@ -71,7 +71,6 @@ static int proxy_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subc e->dgram[0] = ETCP_ID_TCP_PROXY; e->dgram[1] = subcmd; memcpy(e->dgram + 2, &sid, 4); - memcpy(e->dgram + 6, &seq, 2); if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); e->len = TCP_PROXY_HDR_SIZE + len; return etcp_route_send(inst, dst, e); @@ -84,25 +83,23 @@ static int send_connect(struct proxy_conn* pc) { pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port), (unsigned long long)pc->proxy->via_node_id); return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_CONNECT, pc->stream_id, 0, buf, 6); + TCP_PROXY_SUBCMD_CONNECT, pc->stream_id, buf, 6); } static int send_data(struct proxy_conn* pc, const uint8_t* data, uint16_t len) { - DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY SEND sid=%08x seq=%u len=%u", pc->stream_id, pc->send_seq, len); - int ret = proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_DATA, pc->stream_id, pc->send_seq, data, len); - if (ret == 0) pc->send_seq++; - return ret; + DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY SEND sid=%08x len=%u", pc->stream_id, len); + return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, + TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len); } static int send_close(struct proxy_conn* pc) { return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_CLOSE, pc->stream_id, 0, NULL, 0); + TCP_PROXY_SUBCMD_CLOSE, pc->stream_id, NULL, 0); } static int send_error(struct proxy_conn* pc) { return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_ERROR, pc->stream_id, 0, NULL, 0); + TCP_PROXY_SUBCMD_ERROR, pc->stream_id, NULL, 0); } // ==================================================================== @@ -410,31 +407,12 @@ static void handle_data(struct tcp_proxy* p, struct ETCP_CONN* conn, uint32_t st struct proxy_conn* pc = find_pc_by_stream(p, stream_id); if (!pc) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — нет соединения, шлём CLOSE", stream_id); - if (conn) proxy_send_msg(p->inst, conn->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, stream_id, 0, NULL, 0); + if (conn) proxy_send_msg(p->inst, conn->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, stream_id, NULL, 0); queue_dgram_free(entry); queue_entry_free(entry); return; } if (pc->rem_closed || pc->error) { queue_dgram_free(entry); queue_entry_free(entry); return; } - uint16_t seq; memcpy(&seq, entry->dgram + 6, 2); - - if (!pc->recv_init) { - pc->recv_last_seq = seq; pc->recv_init = 1; - } else { - uint16_t delta = seq - pc->recv_last_seq; - if (delta == 0) { - DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy: duplicate DATA seq=%u sid=%08x", seq, stream_id); - queue_dgram_free(entry); queue_entry_free(entry); return; - } - if (delta != 1) { - DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: seq gap seq=%u last=%u sid=%08x", seq, pc->recv_last_seq, stream_id); - pc->error = 1; - if (send_error(pc) < 0) DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY handle_data send_error failed sid=%08x", pc->stream_id); - pc->close_sent = 1; - queue_dgram_free(entry); queue_entry_free(entry); return; - } - pc->recv_last_seq = seq; - } size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; - DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY DATA <- sid=%08x seq=%u len=%zu", stream_id, seq, data_len); + DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY DATA <- sid=%08x len=%zu", stream_id, data_len); if (data_len > 0) { struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len); if (e) queue_data_put(pc->to_lwip, e); diff --git a/src/tcp_proxy.h b/src/tcp_proxy.h index 1a464e35..00b6bd8b 100644 --- a/src/tcp_proxy.h +++ b/src/tcp_proxy.h @@ -25,9 +25,6 @@ struct proxy_conn { struct tcp_pcb* pcb; uint32_t stream_id; - uint16_t send_seq; // 16-битный циклический, только для DATA - uint16_t recv_last_seq; // последний принятый seq для обнаружения пропусков - uint8_t recv_init; // первый DATA принят struct ll_queue* to_lwip; // DATA от exit → lwIP (управление потоком) diff --git a/tests/test_etcp_router_unit.c b/tests/test_etcp_router_unit.c index c6385844..8b1ee3d4 100644 --- a/tests/test_etcp_router_unit.c +++ b/tests/test_etcp_router_unit.c @@ -448,23 +448,31 @@ static int test_ack_timer(void) { return 0; } -static int test_legacy_no_conn(void) { - TEST("legacy — no seq-conn, direct delivery"); +static int test_send_q(void) { + TEST("send_q — inflight full queues, not lost"); 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"); + struct ETCP_ROUTER_CONN* c = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID); + if (!c) FAIL("conn_get failed"); + + // Simulate full inflight: set rx_acked far behind tx_seq + c->tx_seq = ROUTER_MAX_INFLIGHT; + c->rx_acked = 0; + if (c->send_q == NULL) FAIL("send_q is NULL"); - 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); + // Send via conn_send — should queue, not drop + 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 return 0 (queued in send_q)"); + if (c->send_blocked == 0) FAIL("send_blocked not set"); + if (c->send_resume_timer == NULL) FAIL("send_resume_timer not set"); + if (queue_entry_count(c->send_q) != 1) FAIL("send_q count != 1"); + if (c->tx_seq != ROUTER_MAX_INFLIGHT) FAIL("tx_seq advanced while queued"); - if (rx.delivered != 3) FAIL("legacy: delivered != 3"); - if (rx.errors != 0) FAIL("legacy: errors != 0"); + // Cancel timers manually + if (c->send_resume_timer) { uasync_cancel_timeout(ua, c->send_resume_timer); c->send_resume_timer = NULL; } - etcp_router_unbind(&inst, 0x50); + etcp_router_conn_close(c); TEARDOWN(); PASS(); return 0; @@ -550,7 +558,7 @@ int main(void) { test_data_integrity(); test_ack_receive(); test_ack_timer(); - test_legacy_no_conn(); + test_send_q(); test_loopback(); test_timestamp(); test_conn_send_no_bgp();