Browse Source

etcp_router: consumer-driven ACK with interval throttle (10ms)

- consumer_ack flag: ACK sent only on consumption, not assembly
- etcp_router_consumer_ack() — called from tcp_proxy_client feed_from_transport
- last_ack_sent_tb: interval-based throttle — send immediately if >=10ms passed,
  otherwise timer for remaining time
- timer rules: NULL handle after cancel/fire, NULL check before start
- ROUTER_ACK_INTERVAL_TB 1000→100 (100ms→10ms)
- test_etcp_router_unit: updated test 14 for new immediate-send behavior
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
ac3fe79db8
  1. 5
      lib/u_async.h
  2. 62
      src/etcp_router.c
  3. 9
      src/etcp_router.h
  4. 1
      src/proxy/tcp_proxy_client.c
  5. 4
      tests/test_etcp_router_unit.c

5
lib/u_async.h

@ -86,6 +86,11 @@ uint64_t get_time_tb(void);
uint64_t get_time_us(void); uint64_t get_time_us(void);
// Timeouts, timebase = 0.1 mS // Timeouts, timebase = 0.1 mS
// Как работать с таймаутами если таймаут может быть использован несколько раз:
// 1. заводим дескриптор таймаута и обнуляем
// 2. перед активацией таймаута проверяем дескриптор (t_id) на null
// 3. при cancel или срабатывании дескриптор обнуляем
// это обеспечит отсутствие утечек и задвоений таймаутов
void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* user_arg, timeout_callback_t callback, const char* name); void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* user_arg, timeout_callback_t callback, const char* name);
err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id); err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id);

62
src/etcp_router.c

@ -191,14 +191,10 @@ static void router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* p
// ACK timers // ACK timers
// ==================================================================== // ====================================================================
static void router_schedule_ack(struct ETCP_ROUTER_CONN* rconn) { static void router_ack_do_send(struct ETCP_ROUTER_CONN* rconn) {
if (rconn->idle_ack_timer) { router_send_ack(rconn);
uasync_cancel_timeout(rconn->inst->ua, rconn->idle_ack_timer); rconn->last_sent_ack_seq = rconn->rx_seq;
rconn->idle_ack_timer = NULL; rconn->last_ack_sent_tb = get_time_tb();
}
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) { static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) {
@ -227,28 +223,43 @@ static void router_send_ack(struct ETCP_ROUTER_CONN* rconn) {
etcp_send(conn, entry); etcp_send(conn, entry);
} }
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->rx_seq == rconn->last_sent_ack_seq) return;
uint64_t now = get_time_tb();
uint64_t elapsed = now - rconn->last_ack_sent_tb;
if (elapsed >= ROUTER_ACK_INTERVAL_TB) {
router_ack_do_send(rconn);
return;
}
if (!rconn->ack_timer)
rconn->ack_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_ACK_INTERVAL_TB - (uint32_t)elapsed, rconn, router_ack_timer_cb, "router_ack");
}
static void router_ack_timer_cb(void* arg) { static void router_ack_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
rconn->ack_timer = NULL; rconn->ack_timer = NULL;
if (rconn->rx_seq != rconn->last_sent_ack_seq) {
router_send_ack(rconn);
rconn->last_sent_ack_seq = rconn->rx_seq;
}
if (rconn->rx_seq != rconn->last_sent_ack_seq) if (rconn->rx_seq != rconn->last_sent_ack_seq)
router_ack_do_send(rconn);
if (rconn->rx_seq != rconn->last_sent_ack_seq) {
if (!rconn->ack_timer)
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");
else } else if (!rconn->idle_ack_timer) {
rconn->idle_ack_timer = uasync_set_timeout(rconn->inst->ua, rconn->idle_ack_timer = uasync_set_timeout(rconn->inst->ua,
ROUTER_ACK_IDLE_TB, rconn, router_idle_ack_timer_cb, "router_idle_ack"); ROUTER_ACK_IDLE_TB, rconn, router_idle_ack_timer_cb, "router_idle_ack");
}
} }
static void router_idle_ack_timer_cb(void* arg) { static void router_idle_ack_timer_cb(void* arg) {
struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg;
rconn->idle_ack_timer = NULL; rconn->idle_ack_timer = NULL;
if (rconn->rx_seq != rconn->last_sent_ack_seq) { if (rconn->rx_seq != rconn->last_sent_ack_seq)
router_send_ack(rconn); router_ack_do_send(rconn);
rconn->last_sent_ack_seq = rconn->rx_seq;
}
} }
// ==================================================================== // ====================================================================
@ -380,7 +391,7 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (seq == rconn->rx_seq) if (seq == rconn->rx_seq)
router_try_assembly(rconn, conn); router_try_assembly(rconn, conn);
router_schedule_ack(rconn); if (!rconn->consumer_ack) router_schedule_ack(rconn);
queue_dgram_free(entry); queue_entry_free(entry); queue_dgram_free(entry); queue_entry_free(entry);
} else { } else {
// ========== Транзит ========== // ========== Транзит ==========
@ -428,9 +439,9 @@ void etcp_router_destroy(struct UTUN_INSTANCE* inst) {
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;
router_send_close_to_service(rconn); router_send_close_to_service(rconn);
if (rconn->ack_timer) uasync_cancel_timeout(inst->ua, rconn->ack_timer); if (rconn->ack_timer) { uasync_cancel_timeout(inst->ua, rconn->ack_timer); rconn->ack_timer = NULL; }
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); rconn->idle_ack_timer = NULL; }
if (rconn->send_resume_timer) uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; }
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); }
@ -536,6 +547,13 @@ void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id
queue_waiter_cancel(conn->normalizer->input, h); queue_waiter_cancel(conn->normalizer->input, h);
} }
void etcp_router_consumer_ack(struct UTUN_INSTANCE* inst, uint64_t remote_node_id, uint8_t svc_id) {
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, remote_node_id, svc_id);
if (!rconn) return;
rconn->consumer_ack = 1;
router_schedule_ack(rconn);
}
// ==================================================================== // ====================================================================
// Seq-connection API // Seq-connection API
// ==================================================================== // ====================================================================
@ -556,6 +574,8 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst,
rconn->rx_seq = 0; rconn->rx_seq = 0;
rconn->tx_acked = 0; rconn->tx_acked = 0;
rconn->last_sent_ack_seq = 0; rconn->last_sent_ack_seq = 0;
rconn->consumer_ack = 0;
rconn->last_ack_sent_tb = 0;
rconn->inst = inst; rconn->inst = inst;
rconn->ack_timer = NULL; rconn->ack_timer = NULL;
rconn->idle_ack_timer = NULL; rconn->idle_ack_timer = NULL;

9
src/etcp_router.h

@ -34,10 +34,12 @@ struct ETCP_ROUTER_CONN {
uint32_t rx_seq; // ожидаемый seq для сборки (next expected) uint32_t rx_seq; // ожидаемый seq для сборки (next expected)
uint32_t tx_acked; // сколько наших пакетов подтвердил remote (для inflight) uint32_t tx_acked; // сколько наших пакетов подтвердил remote (для inflight)
uint32_t last_sent_ack_seq; // последний отправленный ACK (= rx_seq на момент отправки) uint32_t last_sent_ack_seq; // последний отправленный ACK (= rx_seq на момент отправки)
uint8_t consumer_ack; // 1 = ACK отправляется только при потреблении, не при сборке
uint64_t last_ack_sent_tb; // время отправки последнего ACK (timebase 0.1ms)
struct UTUN_INSTANCE* inst; struct UTUN_INSTANCE* inst;
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; // периодический таймер (10ms при consumer_ack)
void* idle_ack_timer; // idle таймер (дослать последний ack) void* idle_ack_timer; // idle таймер (дослать последний ack)
struct ll_queue* send_q; // очередь ожидающих отправки (inflight полон) struct ll_queue* send_q; // очередь ожидающих отправки (inflight полон)
@ -48,7 +50,7 @@ struct ETCP_ROUTER_CONN {
#define ROUTER_CONN_HASH_SIZE 256 #define ROUTER_CONN_HASH_SIZE 256
#define ROUTER_RECVQ_HASH_SIZE 1024 #define ROUTER_RECVQ_HASH_SIZE 1024
#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 100 // интервал ACK: 10ms в 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 #define ROUTER_SEND_RESUME_TB 500 // retry интервал send_q: 50ms
@ -99,4 +101,7 @@ void etcp_router_waiter_register(struct UTUN_INSTANCE* inst, uint64_t peer_node_
void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
struct queue_waiter_handle* h); struct queue_waiter_handle* h);
// Уведомить router о потреблении данных потребителем (для consumer-driven ACK)
void etcp_router_consumer_ack(struct UTUN_INSTANCE* inst, uint64_t remote_node_id, uint8_t svc_id);
#endif // ETCP_ROUTER_H #endif // ETCP_ROUTER_H

1
src/proxy/tcp_proxy_client.c

@ -186,6 +186,7 @@ static void tcp_proxy_client_feed_from_transport(struct tcp_proxy_client_conn *p
} }
queue_resume_callback(pc->to_lwip); queue_resume_callback(pc->to_lwip);
if (sent_any) { if (sent_any) {
etcp_router_consumer_ack(pc->proxy->inst, pc->proxy->via_node_id, ETCP_ID_TCP_PROXY);
uint32_t unsent = 0; struct tcp_seg* s; uint32_t unsent = 0; struct tcp_seg* s;
for (s = pc->pcb->unsent; s; s = s->next) unsent++; for (s = pc->pcb->unsent; s; s = s->next) unsent++;
DEBUG_DEBUG(DEBUG_CATEGORY_TRAFFIC, "PROXY FEED exit->client sid=%08x fed=%u q=%u->%u snd_wnd=%u cwnd=%u unsent=%u", DEBUG_DEBUG(DEBUG_CATEGORY_TRAFFIC, "PROXY FEED exit->client sid=%08x fed=%u q=%u->%u snd_wnd=%u cwnd=%u unsent=%u",

4
tests/test_etcp_router_unit.c

@ -436,8 +436,8 @@ static int test_ack_timer(void) {
uint8_t data[] = { TEST_SVC_ID, 0x55, 0 }; uint8_t data[] = { TEST_SVC_ID, 0x55, 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, 0, data + 1, 2, 0);
if (c->ack_timer == NULL) FAIL("ack_timer not scheduled after data"); if (c->ack_timer == NULL && c->last_sent_ack_seq != 1) FAIL("ack not sent (no timer, no immediate)");
if (c->last_sent_ack_seq != 0) FAIL("last_sent_ack_seq changed before timer fired"); if (c->ack_timer != NULL && c->last_sent_ack_seq != 0) FAIL("last_sent_ack_seq changed before timer fired");
// Cancel timer manually (it would fire in real event loop) // Cancel timer manually (it would fire in real event loop)
if (c->ack_timer) uasync_cancel_timeout(ua, c->ack_timer); if (c->ack_timer) uasync_cancel_timeout(ua, c->ack_timer);

Loading…
Cancel
Save