diff --git a/lib/tcp_io.c b/lib/tcp_io.c index 6705ab68..2c7550fd 100644 --- a/lib/tcp_io.c +++ b/lib/tcp_io.c @@ -105,6 +105,8 @@ void tcp_conn_destroy(struct tcp_conn* tc) { } queue_waiter_cancel(tc->read_queue, &tc->read_waiter); queue_set_empty_callback(tc->read_queue, NULL, NULL); + queue_set_callback(tc->read_queue, NULL, NULL); + queue_set_callback(tc->write_queue, NULL, NULL); struct ll_entry* e; while ((e = queue_data_get(tc->read_queue)) != NULL) { diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 36596283..7ba90b62 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -107,10 +107,14 @@ static void write_on_get_cb(struct ll_queue* q, void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; if (!inst) return; - DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:ACK fd=%d sid=%08x write_q=%d", - rc->tc ? (int)rc->tc->sock : -1, rc->stream_id, - queue_entry_count(q)); - etcp_router_consumer_ack(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY); + rc->ack_batch++; + if (rc->ack_batch >= 8) { + rc->ack_batch = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:ACK fd=%d sid=%08x write_q=%d", + rc->tc ? (int)rc->tc->sock : -1, rc->stream_id, + queue_entry_count(q)); + etcp_router_consumer_ack(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY); + } } static void on_fin_cb(struct tcp_conn* tc, void* arg) { diff --git a/src/proxy/tcp_proxy_server.h b/src/proxy/tcp_proxy_server.h index 416430aa..c5efd5dd 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/src/proxy/tcp_proxy_server.h @@ -36,6 +36,7 @@ struct tcp_proxy_server_conn { void* close_timer; // таймер повтора CLOSE/ERROR int close_backoff; // backoff: 50..5000 tb (5ms..500ms) void* diag_timer; // 1-секундный таймер диагностики + uint8_t ack_batch; // счётчик для throttled consumer_ack (каждые 8 записей) struct queue_waiter_handle pause_waiter; struct UASYNC* ua;