diff --git a/lib/u_async.h b/lib/u_async.h index fbfc30fb..2e7a003a 100644 --- a/lib/u_async.h +++ b/lib/u_async.h @@ -86,6 +86,11 @@ uint64_t get_time_tb(void); uint64_t get_time_us(void); // 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); err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id); diff --git a/src/etcp_router.c b/src/etcp_router.c index 555c7796..4e5426ff 100644 --- a/src/etcp_router.c +++ b/src/etcp_router.c @@ -191,14 +191,10 @@ static void router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* p // 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_ack_do_send(struct ETCP_ROUTER_CONN* rconn) { + router_send_ack(rconn); + rconn->last_sent_ack_seq = rconn->rx_seq; + rconn->last_ack_sent_tb = get_time_tb(); } 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); } +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) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)arg; 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) - rconn->ack_timer = uasync_set_timeout(rconn->inst->ua, - ROUTER_ACK_INTERVAL_TB, rconn, router_ack_timer_cb, "router_ack"); - else + 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, + ROUTER_ACK_INTERVAL_TB, rconn, router_ack_timer_cb, "router_ack"); + } else if (!rconn->idle_ack_timer) { 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->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) + router_ack_do_send(rconn); } // ==================================================================== @@ -380,7 +391,7 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) if (seq == rconn->rx_seq) 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); } else { // ========== Транзит ========== @@ -428,9 +439,9 @@ void etcp_router_destroy(struct UTUN_INSTANCE* inst) { while ((entry = queue_data_get(inst->router_conns)) != NULL) { struct ETCP_ROUTER_CONN* rconn = (struct ETCP_ROUTER_CONN*)entry; router_send_close_to_service(rconn); - 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->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); rconn->idle_ack_timer = NULL; } + if (rconn->send_resume_timer) { uasync_cancel_timeout(inst->ua, rconn->send_resume_timer); rconn->send_resume_timer = NULL; } 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); } @@ -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); } +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 // ==================================================================== @@ -556,6 +574,8 @@ struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst, rconn->rx_seq = 0; rconn->tx_acked = 0; rconn->last_sent_ack_seq = 0; + rconn->consumer_ack = 0; + rconn->last_ack_sent_tb = 0; rconn->inst = inst; rconn->ack_timer = NULL; rconn->idle_ack_timer = NULL; diff --git a/src/etcp_router.h b/src/etcp_router.h index 5e69fed1..7645a1f3 100644 --- a/src/etcp_router.h +++ b/src/etcp_router.h @@ -34,10 +34,12 @@ struct ETCP_ROUTER_CONN { uint32_t rx_seq; // ожидаемый seq для сборки (next expected) uint32_t tx_acked; // сколько наших пакетов подтвердил remote (для inflight) 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 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) struct ll_queue* send_q; // очередь ожидающих отправки (inflight полон) @@ -48,7 +50,7 @@ struct ETCP_ROUTER_CONN { #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_INTERVAL_TB 100 // интервал ACK: 10ms в timebase (0.1ms) #define ROUTER_ACK_IDLE_TB 5000 // idle таймаут: 500ms #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, 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 diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index e36316f5..9ffb2965 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/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); 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; 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", diff --git a/tests/test_etcp_router_unit.c b/tests/test_etcp_router_unit.c index bcd902a7..939ca8c7 100644 --- a/tests/test_etcp_router_unit.c +++ b/tests/test_etcp_router_unit.c @@ -436,8 +436,8 @@ static int test_ack_timer(void) { 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->last_sent_ack_seq != 0) FAIL("last_sent_ack_seq changed before timer fired"); + if (c->ack_timer == NULL && c->last_sent_ack_seq != 1) FAIL("ack not sent (no timer, no immediate)"); + 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) if (c->ack_timer) uasync_cancel_timeout(ua, c->ack_timer);