|
|
|
|
@ -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"); |
|
|
|
|
|