diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index a9afc902..8935f759 100644 --- a/src/proxy/socks_proxy.c +++ b/src/proxy/socks_proxy.c @@ -331,7 +331,7 @@ static void process_http_request(struct socks_proxy_conn* c) { u_free(pkt); } else { c->tx_buf = pkt; c->tx_len = total; - etcp_router_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c); + etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c); } c->buf_len = 0; @@ -395,7 +395,7 @@ static void on_read_cb(struct ll_queue* q, void* arg) { if (c->tx_buf) { memcpy(c->tx_buf, e->dgram, e->len); c->tx_len = e->len; } 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_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c); + etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c); } } @@ -408,7 +408,7 @@ static void tx_waiter_cb(struct ll_queue* q, void* arg) { u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; queue_resume_callback(c->tc->read_queue); } else { - etcp_router_waiter_register(c->inst, c->via_node_id, &c->tx_waiter, tx_waiter_cb, c); + etcp_router_on_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter, tx_waiter_cb, c); } } @@ -522,7 +522,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: REM_CLOSED sid=%08x", stream_id); c->rem_closed = 1; - etcp_router_waiter_cancel(c->inst, c->via_node_id, &c->tx_waiter); + etcp_router_cancel_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter); tcp_conn_push_close(c->tc); return 1; } @@ -530,7 +530,7 @@ int socks_proxy_handle_etcp(struct socks_proxy_conn** head, int* count, if (subcmd == TCP_PROXY_SUBCMD_ERROR) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks_proxy: ERROR from exit sid=%08x", stream_id); c->rem_closed = 1; - etcp_router_waiter_cancel(c->inst, c->via_node_id, &c->tx_waiter); + etcp_router_cancel_send_ready(c->inst, c->via_node_id, ETCP_ID_TCP_PROXY, &c->tx_waiter); tcp_conn_push_close(c->tc); return 1; } diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 5a2f194c..aedb0e40 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -33,7 +33,6 @@ 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); @@ -218,14 +217,14 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) { int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 0); 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); return; } } else { - if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry"); + etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, + &rc->pause_waiter, pause_resume_cb, rc); return; } } @@ -262,7 +261,8 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) { memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry", (int)rc->tc->sock, rc->stream_id, rc->tx_len); - if (!rc->retry_timer) rc->retry_timer = uasync_set_timeout(rc->ua, 5000, rc, retry_timer_cb, "tps_retry"); + etcp_router_on_send_ready(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, + &rc->pause_waiter, pause_resume_cb, rc); } } @@ -273,25 +273,6 @@ static void pause_resume_cb(struct ll_queue* q, void* arg) { 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); - if (ret == 0) { - u_free(rc->tx_buf); rc->tx_buf = NULL; rc->tx_len = 0; - if (rc->dst_fin_deferred && rc->tc && 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"); - } -} - // ==================================================================== // Управление жизненным циклом коннекта // ==================================================================== @@ -323,8 +304,8 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) { struct tcp_proxy_server_conn** prev = &rc->ctx->conns; while (*prev) { if (*prev == rc) { *prev = rc->next; rc->ctx->conn_count--; break; } prev = &(*prev)->next; } } + if (rc->ctx && rc->ctx->inst) etcp_router_cancel_send_ready(rc->ctx->inst, rc->peer_node_id, ETCP_ID_TCP_PROXY_CLIENT, &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 7a65a674..2193eeb5 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/src/proxy/tcp_proxy_server.h @@ -36,7 +36,6 @@ struct tcp_proxy_server_conn { void* close_timer; // таймер повтора CLOSE/ERROR int close_backoff; // backoff: 50..5000 tb (5ms..500ms) void* diag_timer; // 1-секундный таймер диагностики - void* retry_timer; // таймер повтора tx_buf при backpressure (500ms) uint8_t dst_fin_deferred; // FIN от destination отложен (ждём drain read_queue + tx_buf) uint8_t freed; // 1 = уже освобождён, защита от double-free diff --git a/tests/test_socks_http_proxy.c b/tests/test_socks_http_proxy.c index 39c659f0..aa731de4 100644 --- a/tests/test_socks_http_proxy.c +++ b/tests/test_socks_http_proxy.c @@ -304,13 +304,13 @@ static int worker_main(int wtype, const char* socks_proxy, const char* http_prox unlink(f1); unlink(out); snprintf(cmd, sizeof(cmd), - "curl -s -I --connect-timeout 10 --max-time 30 -x http://%s %s/f_1k.bin 2>/dev/null | head -1 | grep -q '200 OK'", + "curl -s -I --connect-timeout 15 --max-time 30 -x http://%s %s/f_1k.bin 2>/dev/null | head -1 | grep -q '200 OK'", http_proxy, url_http); if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s HEAD 200\n", tname); return 1; } snprintf(out, sizeof(out), "%s/opt_%d.txt", tmpdir, getpid()); snprintf(cmd, sizeof(cmd), - "curl -s -i --connect-timeout 10 --max-time 30 -x http://%s -X OPTIONS -o %s %s/f_1k.bin 2>/dev/null", + "curl -s -i --connect-timeout 15 --max-time 30 -x http://%s -X OPTIONS -o %s %s/f_1k.bin 2>/dev/null", http_proxy, out, url_http); if (system(cmd) != 0) { fprintf(stderr, "[FAIL] %s OPTIONS curl\n", tname); return 1; } snprintf(cmd, sizeof(cmd), "grep -q 'Allow:' %s", out);