Browse Source

etcp_router: remove legacy, auto-conn for all services, send_q inflight control, remove tcp_proxy internal seq

- 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)
congestion
Evgeny 4 months ago
parent
commit
21cfe2d381
  1. 244
      src/etcp_router.c
  2. 9
      src/etcp_router.h
  3. 37
      src/remote_proxy.c
  4. 10
      src/remote_proxy.h
  5. 42
      src/tcp_proxy.c
  6. 3
      src/tcp_proxy.h
  7. 34
      tests/test_etcp_router_unit.c

244
src/etcp_router.c

@ -1,7 +1,7 @@
// etcp_router.c — Сервисный слой маршрутизации поверх ETCP // etcp_router.c — Сервисный слой маршрутизации поверх ETCP
// Упрощённый TCP: восстановление порядка (recv_q), дедупликация, без переповторов // Упрощённый TCP: восстановление порядка (recv_q), дедупликация, без переповторов
// ACK: периодическая отправка rx_seq (не чаще 100ms), idle-таймер для последнего seq // ACK: периодическая отправка rx_seq (не чаще 100ms), idle-таймер для последнего seq
// Inflight контроль: tx_seq - rx_acked < ROUTER_MAX_INFLIGHT // Inflight контроль через send_q: при переполнении — очередь + retry-таймер, без дропов
#include "etcp_router.h" #include "etcp_router.h"
#include "etcp.h" #include "etcp.h"
#include "utun_instance.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_idle_ack_timer_cb(void* arg);
static void router_send_ack(struct ETCP_ROUTER_CONN* rconn); static void router_send_ack(struct ETCP_ROUTER_CONN* rconn);
static void router_schedule_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_try_assembly(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn);
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);
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 // Управление 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); 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 // ACK timers
// ==================================================================== // ====================================================================
@ -83,7 +160,6 @@ static void router_ack_timer_cb(void* arg) {
router_send_ack(rconn); router_send_ack(rconn);
rconn->rx_acked = rconn->rx_seq; rconn->rx_acked = rconn->rx_seq;
} }
// перезапускаем если за время ожидания пришли новые данные
if (rconn->rx_seq != rconn->rx_acked) if (rconn->rx_seq != rconn->rx_acked)
rconn->ack_timer = uasync_set_timeout(rconn->inst->ua, rconn->ack_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_ACK_INTERVAL_TB, rconn, router_ack_timer_cb, "router_ack"); 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) // 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) { const uint8_t* payload, size_t payload_len) {
struct UTUN_INSTANCE* inst = rconn->inst; struct UTUN_INSTANCE* inst = rconn->inst;
etcp_recv_fn cb = inst->router_bindings.callbacks[rconn->svc_id]; 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", 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); 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; uint32_t next = rconn->rx_seq;
while (1) { while (1) {
struct ll_entry* entry = queue_find_data_by_index(rconn->recv_q, &next); struct ll_entry* entry = queue_find_data_by_index(rconn->recv_q, &next);
if (!entry) break; if (!entry) break;
queue_remove_data(rconn->recv_q, entry); 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; rconn->rx_seq = next + 1;
next++; next++;
queue_dgram_free(entry); 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) { 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); struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, hdr->src_node_id, hdr->svc_id);
if (pl_len == 0) { if (pl_len == 0) {
// ACK-пакет: обновляем rx_acked и выходим // ACK-пакет: обновляем rx_acked, drain send_q, выходим
if (rconn) { 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->rx_acked = hdr->seq;
rconn->last_dgram_ts = get_current_timestamp(); rconn->last_dgram_ts = get_current_timestamp();
if (rconn->send_blocked) router_drain_send_q(rconn);
} }
queue_dgram_free(entry); queue_entry_free(entry); queue_dgram_free(entry); queue_entry_free(entry);
return; return;
} }
// Data-пакет // Data-пакет — авто-создаём conn если ещё нет
if (!rconn) { if (!rconn) {
// Нет seq-состояния — legacy-доставка напрямую rconn = etcp_router_conn_get(inst, hdr->src_node_id, hdr->svc_id);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: legacy deliver svc_id=%u from %016llx len=%zu", if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return; }
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 // Seq-состояние есть — reorder
@ -221,7 +280,7 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
} }
// Кладём в recv_q // Кладём в 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; } if (!qe) { queue_dgram_free(entry); queue_entry_free(entry); return; }
*(uint32_t*)qe->data = seq; *(uint32_t*)qe->data = seq;
qe->dgram = u_malloc(pl_len); 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), seq, rconn->rx_seq, queue_entry_count(rconn->recv_q),
(unsigned long long)hdr->src_node_id, rconn->svc_id); (unsigned long long)hdr->src_node_id, rconn->svc_id);
// Если пришёл ожидаемый — собираем contiguous цепочку
if (seq == rconn->rx_seq) if (seq == rconn->rx_seq)
router_try_assembly(rconn); router_try_assembly(rconn, conn);
// Запускаем/продлеваем ACK таймер
router_schedule_ack(rconn); router_schedule_ack(rconn);
queue_dgram_free(entry); queue_entry_free(entry); queue_dgram_free(entry); queue_entry_free(entry);
} else { } else {
// ========== Транзит ========== // ========== Транзит ==========
@ -283,20 +339,23 @@ int etcp_router_init(struct UTUN_INSTANCE* inst) {
void etcp_router_destroy(struct UTUN_INSTANCE* inst) { void etcp_router_destroy(struct UTUN_INSTANCE* inst) {
if (!inst) return; if (!inst) return;
etcp_unbind(inst, ETCP_ID_SVC_ROUTE); etcp_unbind(inst, ETCP_ID_SVC_ROUTE);
// Очищаем router_conns
if (inst->router_conns) { if (inst->router_conns) {
struct ll_entry* entry; struct ll_entry* entry;
while ((entry = queue_data_get(inst->router_conns)) != NULL) { while ((entry = queue_data_get(inst->router_conns)) != NULL) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry; struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry;
if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->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->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) { if (rconn->recv_q) {
struct ll_entry* f; struct ll_entry* f;
while ((f = queue_data_get(rconn->recv_q)) != NULL) { while ((f = queue_data_get(rconn->recv_q)) != NULL) { queue_dgram_free(f); queue_entry_free(f); }
queue_dgram_free(f); queue_entry_free(f);
}
queue_free(rconn->recv_q); 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_entry_free(entry);
} }
queue_free(inst->router_conns); 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) { int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry) {
if (!inst || !entry) { if (!inst || !entry) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: NULL inst=%p entry=%p", (void*)inst, (void*)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); } if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return -1;
return -1;
} }
if (!entry->dgram || entry->len < 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;
@ -347,35 +405,18 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_
return 0; return 0;
} }
size_t total_len = SVC_ROUTE_HDR_SIZE + payload_len; struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, dst_node_id, svc_id);
uint8_t* new_dgram = u_malloc(total_len); if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
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 ll_entry* new_entry = queue_entry_new(0); if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) {
if (!new_entry) { u_free(new_dgram); queue_dgram_free(entry); queue_entry_free(entry); return -1; } router_enqueue_send(rconn, entry->dgram + 1, payload_len);
new_entry->dgram = new_dgram; queue_dgram_free(entry); queue_entry_free(entry);
new_entry->len = total_len; return 0;
}
int ret = router_send_one(rconn, entry->dgram + 1, payload_len);
queue_dgram_free(entry); queue_entry_free(entry); queue_dgram_free(entry); queue_entry_free(entry);
return ret;
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);
} }
int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id) { 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->inst = inst;
rconn->ack_timer = NULL; rconn->ack_timer = NULL;
rconn->idle_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"); rconn->recv_q = queue_new(inst->ua, ROUTER_RECVQ_HASH_SIZE, 0, 4, "router_recv_q");
if (!rconn->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; }
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "router_conn_get: queue_new(recv_q) failed"); rconn->send_q = queue_new(inst->ua, 0, 0, 0, "router_send_q");
queue_entry_free(&rconn->ll); 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; }
return NULL;
}
queue_data_put_with_index(inst->router_conns, &rconn->ll); queue_data_put_with_index(inst->router_conns, &rconn->ll);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: new conn remote=%016llx svc_id=%u", DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: new conn remote=%016llx svc_id=%u",
(unsigned long long)remote_node_id, svc_id); (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); (void*)rconn, (const void*)data, len);
return -1; return -1;
} }
// Inflight контроль if (rconn->tx_seq - rconn->rx_acked >= ROUTER_MAX_INFLIGHT) {
uint32_t inflight = rconn->tx_seq - rconn->rx_acked; router_enqueue_send(rconn, data, len);
if (inflight >= ROUTER_MAX_INFLIGHT) { return 0;
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(); return router_send_one(rconn, data, len);
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) { void etcp_router_conn_close(struct ETCP_ROUTER_CONN* rconn) {
if (!rconn) return; if (!rconn) return;
struct UTUN_INSTANCE* inst = rconn->inst; struct UTUN_INSTANCE* inst = rconn->inst;
if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->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->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) { if (rconn->recv_q) {
struct ll_entry* f; 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); 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) if (inst->router_conns)
queue_remove_data(inst->router_conns, &rconn->ll); queue_remove_data(inst->router_conns, &rconn->ll);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: closed remote=%016llx svc_id=%u", DEBUG_INFO(DEBUG_CATEGORY_ETCP, "router_conn: closed remote=%016llx svc_id=%u",

9
src/etcp_router.h

@ -38,6 +38,10 @@ struct ETCP_ROUTER_CONN {
struct ll_queue* recv_q; // reorder очередь: hash по seq (4 байта, offset 0) struct ll_queue* recv_q; // reorder очередь: hash по seq (4 байта, offset 0)
void* ack_timer; // периодический таймер (100ms) void* ack_timer; // периодический таймер (100ms)
void* idle_ack_timer; // idle таймер (дослать последний ack) 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 #define ROUTER_CONN_HASH_SIZE 256
@ -45,6 +49,7 @@ struct ETCP_ROUTER_CONN {
#define ROUTER_MAX_INFLIGHT 256 // макс пакетов в полёте (для контроля inflight) #define ROUTER_MAX_INFLIGHT 256 // макс пакетов в полёте (для контроля inflight)
#define ROUTER_ACK_INTERVAL_TB 1000 // интервал ACK: 100ms в timebase (0.1ms) #define ROUTER_ACK_INTERVAL_TB 1000 // интервал ACK: 100ms в timebase (0.1ms)
#define ROUTER_ACK_IDLE_TB 5000 // idle таймаут: 500ms #define ROUTER_ACK_IDLE_TB 5000 // idle таймаут: 500ms
#define ROUTER_SEND_RESUME_TB 500 // retry интервал send_q: 50ms
// Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS) // Bindings для сервисов внутри etcp_router (аналогично ETCP_BINDINGS)
struct ETCP_ROUTER_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); 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_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); 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); int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry);
// Найти/создать состояние seq-подключения по (remote_node_id, svc_id) // Найти/создать состояние seq-подключения по (remote_node_id, svc_id)

37
src/remote_proxy.c

@ -39,7 +39,7 @@ void rp_conn_free(struct remote_proxy_conn* rc);
#define RP_NORMALIZER_Q_THRESHOLD 64 #define RP_NORMALIZER_Q_THRESHOLD 64
static int rp_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, 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); 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; } 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); 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[0] = ETCP_ID_TCP_PROXY;
e->dgram[1] = subcmd; e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4); memcpy(e->dgram + 2, &sid, 4);
memcpy(e->dgram + 6, &seq, 2);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len; e->len = TCP_PROXY_HDR_SIZE + len;
return etcp_route_send(inst, dst, e); 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) { uint16_t local_port, uint8_t status) {
uint8_t buf[3]; uint8_t buf[3];
memcpy(buf, &local_port, 2); buf[2] = status; 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) { 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) { 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); 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); rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->pause_buf, rc->pause_len);
rc->send_seq++;
u_free(rc->pause_buf); rc->pause_buf = NULL; rc->pause_len = 0; 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); 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; 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)); 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); rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, buf, (size_t)n);
rc->send_seq++;
} else if (n == 0) { } else if (n == 0) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x total=%d cli_closed=%d sock_closed=%d", 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); (int)rc->sock, rc->stream_id, rp_conn_total(rc), rc->cli_closed, rc->sock_closed);
rc->sock_closed = 1; rc->sock_closed = 1;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; 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); 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->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; }
if (rc->cli_closed) rp_conn_free(rc); 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)); (int)rc->sock, rc->stream_id, errno, strerror(errno), rp_conn_total(rc));
rc->error = 1; rc->error = 1;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; 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); 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); 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); DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:ERR fd=%d sid=%08x", (int)rc->sock, rc->stream_id);
rc->error = 1; rc->error = 1;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; 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); 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); 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); struct remote_proxy_conn* rc = remote_proxy_find_conn(ctx, stream_id);
if (!rc || rc->sock == SOCKET_INVALID) { if (!rc || rc->sock == SOCKET_INVALID) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "RP handle_data drop: нет соединения для sid=%08x, шлём CLOSE", stream_id); 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; queue_dgram_free(entry); queue_entry_free(entry); return -1;
} }
if (rc->error) { 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; size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
if (rc->connected) { if (rc->connected) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:SEND fd=%d sid=%08x len=%zu total=%d", DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:SEND fd=%d sid=%08x len=%zu total=%d",

10
src/remote_proxy.h

@ -20,9 +20,9 @@ struct remote_proxy_ctx;
#define TCP_PROXY_CONNECTED_OK 0 #define TCP_PROXY_CONNECTED_OK 0
#define TCP_PROXY_CONNECTED_REFUSED 1 #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_HDR_SIZE 6 // svc_id(1)+subcmd(1)+stream_id(4)
#define TCP_PROXY_CONNECT_HDR_SIZE 14 // HDR_SIZE + dest_ip(4)+dest_port(2) #define TCP_PROXY_CONNECT_HDR_SIZE 12 // 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_CONNECTED_HDR_SIZE 9 // HDR_SIZE + local_port(2)+status(1)
struct remote_proxy_conn { struct remote_proxy_conn {
struct remote_proxy_conn* next; struct remote_proxy_conn* next;
@ -35,10 +35,6 @@ struct remote_proxy_conn {
uint8_t dest_ip[4]; uint8_t dest_ip[4];
uint16_t dest_port; 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 cli_closed; // client sent CLOSE
uint8_t sock_closed; // socket EOF received uint8_t sock_closed; // socket EOF received
uint8_t error; uint8_t error;

42
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 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 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 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); 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 // Протокол: сборка и отправка сообщений прокси через ETCP
// ==================================================================== // ====================================================================
static int proxy_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, 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); 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; } 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); 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[0] = ETCP_ID_TCP_PROXY;
e->dgram[1] = subcmd; e->dgram[1] = subcmd;
memcpy(e->dgram + 2, &sid, 4); memcpy(e->dgram + 2, &sid, 4);
memcpy(e->dgram + 6, &seq, 2);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len); if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len; e->len = TCP_PROXY_HDR_SIZE + len;
return etcp_route_send(inst, dst, e); 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], 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); ntohs(pc->dest_port), (unsigned long long)pc->proxy->via_node_id);
return proxy_send_msg(pc->proxy->inst, 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) { 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); DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY SEND sid=%08x len=%u", pc->stream_id, len);
int ret = proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id,
TCP_PROXY_SUBCMD_DATA, pc->stream_id, pc->send_seq, data, len); TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len);
if (ret == 0) pc->send_seq++;
return ret;
} }
static int send_close(struct proxy_conn* pc) { static int send_close(struct proxy_conn* pc) {
return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, 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) { static int send_error(struct proxy_conn* pc) {
return proxy_send_msg(pc->proxy->inst, pc->proxy->via_node_id, 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); struct proxy_conn* pc = find_pc_by_stream(p, stream_id);
if (!pc) { if (!pc) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — нет соединения, шлём CLOSE", stream_id); 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; queue_dgram_free(entry); queue_entry_free(entry); return;
} }
if (pc->rem_closed || pc->error) { 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; 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) { 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); 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); if (e) queue_data_put(pc->to_lwip, e);

3
src/tcp_proxy.h

@ -25,9 +25,6 @@ struct proxy_conn {
struct tcp_pcb* pcb; struct tcp_pcb* pcb;
uint32_t stream_id; 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 (управление потоком) struct ll_queue* to_lwip; // DATA от exit → lwIP (управление потоком)

34
tests/test_etcp_router_unit.c

@ -448,23 +448,31 @@ static int test_ack_timer(void) {
return 0; return 0;
} }
static int test_legacy_no_conn(void) { static int test_send_q(void) {
TEST("legacy — no seq-conn, direct delivery"); TEST("send_q — inflight full queues, not lost");
struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb; struct UASYNC* ua; struct UTUN_INSTANCE inst; etcp_recv_fn recv_cb;
SETUP(); SETUP();
// Register handler but DON'T create conn for svc_id=0x50 struct ETCP_ROUTER_CONN* c = etcp_router_conn_get(&inst, TEST_REMOTE_NODE, TEST_SVC_ID);
if (etcp_router_bind(&inst, 0x50, test_handler) != 0) FAIL("bind legacy"); 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); // Send via conn_send — should queue, not drop
uint8_t data[] = { 0x50, 0x33, 0 }; uint8_t data[4] = { 1, 2, 3, 4 };
inject(recv_cb, &inst, TEST_REMOTE_NODE, 0x50, 0, data + 1, 2, 0); int ret = etcp_router_conn_send(c, data, sizeof(data));
inject(recv_cb, &inst, TEST_REMOTE_NODE, 0x50, 1, data + 1, 2, 0); if (ret != 0) FAIL("conn_send should return 0 (queued in send_q)");
inject(recv_cb, &inst, TEST_REMOTE_NODE, 0x50, 2, data + 1, 2, 0); 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"); // Cancel timers manually
if (rx.errors != 0) FAIL("legacy: errors != 0"); 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(); TEARDOWN();
PASS(); PASS();
return 0; return 0;
@ -550,7 +558,7 @@ int main(void) {
test_data_integrity(); test_data_integrity();
test_ack_receive(); test_ack_receive();
test_ack_timer(); test_ack_timer();
test_legacy_no_conn(); test_send_q();
test_loopback(); test_loopback();
test_timestamp(); test_timestamp();
test_conn_send_no_bgp(); test_conn_send_no_bgp();

Loading…
Cancel
Save