diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index f0600388..e0037e13 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -214,10 +214,12 @@ static void diag_timer_cb(void* arg) { static void read_queue_drain_cb(struct ll_queue* q, void* arg) { struct tcp_proxy_server_conn* rc = (struct tcp_proxy_server_conn*)arg; struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL; + struct tcp_conn* tc = rc->tc; struct ll_entry* e; if (rc->tx_buf) { int ret = send_msg(inst, TOPO_GROUP_UTUN, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->tx_buf, rc->tx_len, 0); + if (!tc->read_queue) { u_free(rc->tx_buf); return; } DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "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; @@ -246,7 +248,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) { } if (rc->cli_closed) { do { - memory_pool_free(rc->tc->data_pool, e->dgram); + memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); e = queue_data_get(q); } while (e); @@ -254,19 +256,21 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) { return; } int ret = send_msg(inst, TOPO_GROUP_UTUN, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len, 0); + if (!tc->read_queue) { memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); return; } rc->bytes_relayed += e->len; rc->drain_count++; DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "TPS DRAIN #%u sid=%08x len=%u → ret=%d rq=%d", rc->drain_count, rc->stream_id, e->len, ret, q->count); if (ret == 0) { - memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); + memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(q); } else { rc->tx_buf = u_malloc(e->len); if (rc->tx_buf) { memcpy(rc->tx_buf, e->dgram, e->len); rc->tx_len = e->len; } else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TPS tx_buf malloc=%u failed sid=%08x — drop", e->len, rc->stream_id); } - memory_pool_free(rc->tc->data_pool, e->dgram); queue_entry_free(e); + memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); + if (!tc->read_queue) return; DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry", - (int)rc->tc->sock, rc->stream_id, rc->tx_len); + (int)tc->sock, rc->stream_id, rc->tx_len); etcp_router_on_send_ready(inst, TOPO_GROUP_UTUN, rc->peer_node_id, ETCP_RT_ID_TCP_PROXY_CLIENT, &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");