Browse Source

fix: remove retry_timer, use waiter-based round-robin on send_q for fairness

- tcp_proxy_server: drain_cb uses etcp_router_on_send_ready instead of retry_timer
- socks_proxy: on_read_cb + tx_waiter_cb use etcp_router_on_send_ready instead of waiter_register
- Result: FIFO round-robin via ll_queue waiters on send_q, all streams get fair share
- No timers — wakeup is event-driven via queue_data_get → check_waiters
- TCP backpressure: read_queue callback_suspended prevents over-reading
chatgui
Evgeny 3 months ago
parent
commit
ba14e31e09
  1. 10
      src/proxy/socks_proxy.c
  2. 29
      src/proxy/tcp_proxy_server.c
  3. 1
      src/proxy/tcp_proxy_server.h
  4. 4
      tests/test_socks_http_proxy.c

10
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;
}

29
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; }

1
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

4
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);

Loading…
Cancel
Save