Browse Source

fix: initial consumer_ack in handle_connect + diag: send_q, BACKPRESSURE, ACK logs

handle_connect: send first consumer_ack after setting consumer_ack=1
to unblock client router — fixes rq=32/send_q=256 deadlock.

diag_timer_cb: added send_q=N to DIAG line for instant root cause visibility.
read_queue_drain_cb: SOCK:BACKPRESSURE log with send_q count.
write_on_get_cb: SOCK:ACK debug log on each consumer_ack.
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
4716641dff
  1. 30
      src/proxy/tcp_proxy_server.c

30
src/proxy/tcp_proxy_server.c

@ -104,10 +104,13 @@ static void send_error(struct tcp_proxy_server_conn* rc) {
// ==================================================================== // ====================================================================
static void write_on_get_cb(struct ll_queue* q, void* arg) { static void write_on_get_cb(struct ll_queue* q, void* arg) {
(void)q;
struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg;
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (inst) etcp_router_consumer_ack(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY); 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);
} }
static void on_fin_cb(struct tcp_conn* tc, void* arg) { static void on_fin_cb(struct tcp_conn* tc, void* arg) {
@ -140,6 +143,15 @@ static void diag_timer_cb(void* arg) {
int ent_f = tc->entry_pool->free_count; int ent_f = tc->entry_pool->free_count;
int dat_f = tc->data_pool->free_count; int dat_f = tc->data_pool->free_count;
int sq = 0;
{
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (inst) {
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY);
if (rconn && rconn->send_q) sq = queue_entry_count(rconn->send_q);
}
}
int rcv_buf = 0, snd_buf = 0; int rcv_buf = 0, snd_buf = 0;
socklen_t optlen = sizeof(int); socklen_t optlen = sizeof(int);
getsockopt(tc->sock, SOL_SOCKET, SO_RCVBUF, &rcv_buf, &optlen); getsockopt(tc->sock, SOL_SOCKET, SO_RCVBUF, &rcv_buf, &optlen);
@ -147,11 +159,12 @@ static void diag_timer_cb(void* arg) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, DEBUG_INFO(DEBUG_CATEGORY_SOCKET,
"SOCK:DIAG fd=%d sid=%08x conn=%d fin=%d " "SOCK:DIAG fd=%d sid=%08x conn=%d fin=%d "
"rq=%d(%zub) wq=%d(%zub) wbuf=%s " "rq=%d(%zub) wq=%d(%zub) send_q=%d wbuf=%s "
"ent_f=%d dat_f=%d tcp_rb=%zu tcp_sb=%zu", "ent_f=%d dat_f=%d tcp_rb=%zu tcp_sb=%zu",
(int)tc->sock, rc->stream_id, tc->connected, tc->fin, (int)tc->sock, rc->stream_id, tc->connected, tc->fin,
tc->read_queue->count, queue_total_bytes(tc->read_queue), tc->read_queue->count, queue_total_bytes(tc->read_queue),
tc->write_queue->count, queue_total_bytes(tc->write_queue), tc->write_queue->count, queue_total_bytes(tc->write_queue),
sq,
tc->write_buf ? "y" : "n", tc->write_buf ? "y" : "n",
ent_f, dat_f, ent_f, dat_f,
(size_t)rcv_buf, (size_t)snd_buf); (size_t)rcv_buf, (size_t)snd_buf);
@ -177,9 +190,11 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
queue_resume_callback(q); queue_resume_callback(q);
} else { } else {
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY); struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, rc->peer_node_id, ETCP_ID_TCP_PROXY);
if (rconn && rconn->send_q) if (rconn && rconn->send_q) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:BACKPRESSURE fd=%d sid=%08x send_q=%d — pausing drain",
(int)rc->tc->sock, rc->stream_id, queue_entry_count(rconn->send_q));
queue_waiter_wait(rconn->send_q, &rc->pause_waiter, pause_resume_cb, rc); queue_waiter_wait(rconn->send_q, &rc->pause_waiter, pause_resume_cb, rc);
else } else
etcp_router_waiter_register(inst, rc->peer_node_id, &rc->pause_waiter, pause_resume_cb, rc); etcp_router_waiter_register(inst, rc->peer_node_id, &rc->pause_waiter, pause_resume_cb, rc);
} }
} }
@ -274,7 +289,10 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
{ {
struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, src_node_id, ETCP_ID_TCP_PROXY); struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, src_node_id, ETCP_ID_TCP_PROXY);
if (rconn) rconn->consumer_ack = 1; if (rconn) {
rconn->consumer_ack = 1;
etcp_router_consumer_ack(inst, src_node_id, ETCP_ID_TCP_PROXY);
}
} }
queue_set_on_get(rc->tc->write_queue, write_on_get_cb, rc); queue_set_on_get(rc->tc->write_queue, write_on_get_cb, rc);

Loading…
Cancel
Save