Browse Source

fix: add retry_timer fallback in tcp_proxy_server for backpressure stall

waiter-based backpressure (commit ba14e31e) could starve HTTPS CONNECT
worker under concurrent 4-worker load — send_q never drained below threshold
fast enough for the waiter to fire before curl timed out (exit=28).

Add retry_timer (500ms, force=1) alongside etcp_router_on_send_ready
waiter: waiter gives instant wake when send_q drains, timer guarantees
delivery via force=1 if waiter hasn't fired in time.

Also improve socks_proxy DEBUG diagnostics in backpressure branches.
chatgui v0.1.0-rc2
Evgeny 3 months ago
parent
commit
55fa5c7759
  1. 4
      src/proxy/socks_proxy.c
  2. 29
      src/proxy/tcp_proxy_server.c
  3. 1
      src/proxy/tcp_proxy_server.h

4
src/proxy/socks_proxy.c

@ -386,7 +386,7 @@ static void on_read_cb(struct ll_queue* q, void* arg) {
return;
}
int ret = send_data(c, e->dgram, e->len);
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS PROXY SEND sid=%08x len=%u", c->stream_id, e->len);
DEBUG_TRACE(DEBUG_CATEGORY_TRAFFIC, "SOCKS PROXY SEND sid=%08x len=%u ret=%d is_http=%d", c->stream_id, e->len, ret, c->is_http);
if (ret == 0) {
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
queue_resume_callback(q);
@ -396,6 +396,7 @@ static void on_read_cb(struct ll_queue* q, void* arg) {
else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_proxy: tx_buf malloc=%u failed sid=%08x — drop", e->len, c->stream_id); }
memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e);
etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_RT_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCKS PROXY BP: sid=%08x tx_buf=%u waiter_reg", c->stream_id, c->tx_len);
}
}
@ -404,6 +405,7 @@ static void tx_waiter_cb(struct ll_queue* q, void* arg) {
struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg;
if (!c->tx_buf || c->rem_closed || c->close_sent) return;
int ret = send_data(c, c->tx_buf, c->tx_len);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCKS PROXY WAKE: sid=%08x tx_buf=%u ret=%d is_http=%d", c->stream_id, c->tx_len, ret, c->is_http);
if (ret == 0) {
u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0;
queue_resume_callback(c->tc->read_queue);

29
src/proxy/tcp_proxy_server.c

@ -33,6 +33,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg);
static void pause_resume_cb(struct ll_queue* q, void* arg);
static void close_retry_cb(void* arg);
static void diag_timer_cb(void* arg);
static void retry_timer_cb(void* arg);
static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
uint32_t sid, const uint8_t* data, size_t len, int force);
static void send_close(struct tcp_proxy_server_conn* rc);
@ -214,8 +215,10 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
if (rc->tx_buf) {
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 0);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "TPS TXBUF RETRY: sid=%08x len=%u ret=%d", rc->stream_id, rc->tx_len, ret);
if (ret == 0) {
u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0;
if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_timer = NULL; }
if (rc->dst_fin_deferred && !q->head) {
rc->dst_fin_deferred = 0;
send_fin(rc);
@ -224,6 +227,8 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
} else {
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY,
&rc->pause_waiter, pause_resume_cb, rc);
if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry");
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "TPS TXBUF BP AGAIN: sid=%08x waiter_reg", rc->stream_id);
return;
}
}
@ -261,6 +266,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
(int)rc->tc->sock, rc->stream_id, rc->tx_len);
etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY,
&rc->pause_waiter, pause_resume_cb, rc);
if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry");
}
}
@ -268,12 +274,32 @@ static void pause_resume_cb(struct ll_queue* q, void* arg) {
(void)q;
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
if (!rc->tc || rc->tc->sock == SOCKET_INVALID || rc->cli_closed) return;
DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "TPS RESUME sid=%08x tx_buf=%s rq=%d",
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "TPS RESUME sid=%08x tx_buf=%s rq=%d",
rc->stream_id, rc->tx_buf ? "y" : "n",
rc->tc->read_queue ? rc->tc->read_queue->count : 0);
queue_resume_callback(rc->tc->read_queue);
}
static void retry_timer_cb(void* arg) {
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
rc->retry_timer = NULL;
if (!rc->tx_buf || !rc->tc || rc->tc->sock == SOCKET_INVALID || rc->cli_closed) return;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 1);
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "TPS RETRY TIMER: sid=%08x len=%u force=1 ret=%d", rc->stream_id, rc->tx_len, ret);
if (ret == 0) {
u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0;
if (rc->dst_fin_deferred && rc->tc->read_queue && !rc->tc->read_queue->head) {
rc->dst_fin_deferred = 0;
send_fin(rc);
return;
}
queue_resume_callback(rc->tc->read_queue);
} else {
rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry");
}
}
// ====================================================================
// Управление жизненным циклом коннекта
// ====================================================================
@ -307,6 +333,7 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
}
if (rc->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY, &rc->pause_waiter);
if (rc->tx_buf) { u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; }
if (rc->retry_timer) { uasync_cancel_timeout(rc->ua, rc->retry_timer); rc->retry_timer = NULL; }
if (rc->close_timer) { uasync_cancel_timeout(rc->ua, rc->close_timer); rc->close_timer = NULL; }
if (rc->diag_timer) { uasync_cancel_timeout(rc->ua, rc->diag_timer); rc->diag_timer = NULL; }
if (rc->tc) { tcp_conn_destroy(rc->tc); rc->tc = NULL; }

1
src/proxy/tcp_proxy_server.h

@ -40,6 +40,7 @@ struct tcp_proxy_server_conn {
uint8_t freed; // 1 = уже освобождён, защита от double-free
struct queue_waiter_handle pause_waiter;
void* retry_timer; // fallback timer (500ms) for force=1 retry on stall
uint8_t* tx_buf; // буфер при backpressure (retry в pause_resume_cb)
uint16_t tx_len;
uint32_t bytes_sent; // байт записано в destination (реальный TCP)

Loading…
Cancel
Save