diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 90a9ab2b..36596283 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/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) { - (void)q; struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; 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) { @@ -140,6 +143,15 @@ static void diag_timer_cb(void* arg) { int ent_f = tc->entry_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; socklen_t optlen = sizeof(int); 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, "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", (int)tc->sock, rc->stream_id, tc->connected, tc->fin, tc->read_queue->count, queue_total_bytes(tc->read_queue), tc->write_queue->count, queue_total_bytes(tc->write_queue), + sq, tc->write_buf ? "y" : "n", ent_f, dat_f, (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); } else { 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); - else + } else 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); - 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);