diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index 5ef3d1ee..33f459b0 100644 --- a/src/proxy/socks_proxy.c +++ b/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); diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 2e5556d8..f91bcbf6 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/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; } diff --git a/src/proxy/tcp_proxy_server.h b/src/proxy/tcp_proxy_server.h index 2193eeb5..71e7b497 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/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)