Browse Source

add backpressure support: uip_retry_timer, feed drain at closing_tun/sent, send buffering with retry

- proxy_uip_retry_cb: retry send from uip_to_transport every 5ms
- proxy_recv_cb: buffer to uip_to_transport when send fails, set retry timer
- proxy_feed: drain transport_to_uip when closing_tun || closing_sent (prevents RST)
- proxy_recv_cb(p==NULL): only tcp_recved + closing_tun, no drain
- proxy_poll_cb/tcp_proxy_destroy: cancel retry timer on cleanup
- uip_retry_timer field in proxy_conn
congestion
Evgeny 4 months ago
parent
commit
c3eccd7b35
  1. 37
      src/tcp_proxy.c
  2. 1
      src/tcp_proxy.h

37
src/tcp_proxy.c

@ -53,6 +53,7 @@ static void etcp_transport_close(struct tcp_proxy_transport* t);
static void etcp_transport_destroy(struct tcp_proxy_transport* t);
static void sock_transport_connect(struct sock_transport* st);
static void proxy_try_close(struct proxy_conn *pc);
static void proxy_uip_retry_cb(void* arg);
static struct sock_transport* sock_transport_create(struct proxy_conn* pc, struct UASYNC* ua);
static struct etcp_transport* etcp_transport_create(struct proxy_conn* pc, struct UTUN_INSTANCE* inst, uint64_t remote_node_id);
@ -200,6 +201,14 @@ static void transport_to_uip_cb(struct ll_queue* q, void* arg) {
// ====================================================================
static void proxy_feed_from_transport(struct proxy_conn *pc) {
if (!pc->transport_to_uip || !pc->pcb) return;
if (pc->closing_tun || pc->closing_sent) {
struct ll_entry *e;
while ((e = queue_data_get(pc->transport_to_uip))) {
queue_dgram_free(e); queue_entry_free(e);
queue_resume_callback(pc->transport_to_uip);
}
return;
}
int sent_any = 0;
int entries_fed = 0;
while (1) {
@ -277,9 +286,6 @@ static err_t proxy_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbuf *p, err_t
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY FIN stream=%016llx tun=1 active=%d", (unsigned long long)(pc->remote_stream_id ? pc->remote_stream_id : 0), pc->active);
{ uint16_t wnd_gap = TCP_WND_MAX(pcb) - pcb->rcv_wnd;
if (wnd_gap > 0) tcp_recved(pcb, wnd_gap); }
if (pc->transport_to_uip) {
struct ll_entry* e; while ((e = queue_data_get(pc->transport_to_uip))) { queue_dgram_free(e); queue_entry_free(e); }
}
pc->closing_tun = 1;
return LERR_OK;
}
@ -297,9 +303,11 @@ static err_t proxy_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbuf *p, err_t
pbuf_copy_partial(p, data, len, 0);
int ret = pc->transport->ops->send(pc->transport, data, len);
if (ret < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY send failed stream=%016llx ret=%d len=%u", (unsigned long long)pc->remote_stream_id, ret, len);
struct ll_entry *e = entry_from_data(pc->proxy->entry_pool, data, len);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY send queued stream=%016llx ret=%d len=%u", (unsigned long long)pc->remote_stream_id, ret, len);
struct ll_entry *e = entry_from_data(pc->proxy->entry_pool, p->payload, p->tot_len);
if (e) queue_data_put(pc->uip_to_transport, e);
if (!pc->uip_retry_timer)
pc->uip_retry_timer = uasync_set_timeout(pc->proxy->ua, 50, pc, proxy_uip_retry_cb, "uip_retry");
}
u_free(data);
}
@ -317,6 +325,21 @@ static err_t proxy_sent_cb(void *arg, struct tcp_pcb *pcb, uint16_t len) {
return LERR_OK;
}
static void proxy_uip_retry_cb(void* arg) {
struct proxy_conn* pc = (struct proxy_conn*)arg;
pc->uip_retry_timer = NULL;
struct ll_entry* e = queue_data_get(pc->uip_to_transport);
if (!e) { queue_resume_callback(pc->uip_to_transport); return; }
int ret = pc->transport ? pc->transport->ops->send(pc->transport, e->dgram, e->len) : -1;
if (ret < 0) {
queue_data_put_first(pc->uip_to_transport, e);
pc->uip_retry_timer = uasync_set_timeout(pc->proxy->ua, 50, pc, proxy_uip_retry_cb, "uip_retry");
} else {
queue_dgram_free(e); queue_entry_free(e);
}
queue_resume_callback(pc->uip_to_transport);
}
static void proxy_err_cb(void *arg, err_t err) {
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY proxy_err_cb: pc=NULL err=%d", err); return; }
@ -358,6 +381,7 @@ static err_t proxy_poll_cb(void *arg, struct tcp_pcb *pcb) {
struct proxy_conn **prev = &proxy->conns;
while (*prev) { if (*prev == pc) { *prev = pc->next; proxy->conn_count--; break; } prev = &(*prev)->next; }
if (pc->transport) { pc->transport->ops->destroy(pc->transport); pc->transport = NULL; }
if (pc->uip_retry_timer) { uasync_cancel_timeout(pc->proxy->ua, pc->uip_retry_timer); pc->uip_retry_timer = NULL; }
if (pc->eim_mapping) { pc->eim_mapping->delete_at_tb = get_time_tb() + proxy->eim_timeout_tb; pc->eim_mapping = NULL; }
if (pc->uip_to_transport) { struct ll_entry *e2; while ((e2 = queue_data_get(pc->uip_to_transport))) { queue_dgram_free(e2); queue_entry_free(e2); } queue_free(pc->uip_to_transport); }
if (pc->transport_to_uip) { struct ll_entry *e2; while ((e2 = queue_data_get(pc->transport_to_uip))) { queue_dgram_free(e2); queue_entry_free(e2); } queue_free(pc->transport_to_uip); }
@ -613,7 +637,7 @@ static void sock_transport_error_callback(socket_t sock, void* arg) {
// ====================================================================
static int etcp_transport_send(struct tcp_proxy_transport* t, const uint8_t* data, size_t len) {
struct etcp_transport* et = (struct etcp_transport*)t;
if (!et->connected) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY send not connected stream=%016llx len=%zu", (unsigned long long)et->stream_id, len); return -1; }
if (!et->connected) return -1;
struct ll_entry* e = queue_entry_new(0);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY send alloc entry fail stream=%016llx", (unsigned long long)et->stream_id); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
@ -878,6 +902,7 @@ void tcp_proxy_destroy(struct tcp_proxy* p) {
struct proxy_conn* pc = p->conns;
while (pc) { struct proxy_conn* next = pc->next;
if (pc->transport) pc->transport->ops->destroy(pc->transport);
if (pc->uip_retry_timer) { uasync_cancel_timeout(p->ua, pc->uip_retry_timer); pc->uip_retry_timer = NULL; }
if (pc->uip_to_transport) { struct ll_entry* e; while ((e = queue_data_get(pc->uip_to_transport))) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->uip_to_transport); }
if (pc->transport_to_uip) { struct ll_entry* e; while ((e = queue_data_get(pc->transport_to_uip))) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->transport_to_uip); }
u_free(pc); pc = next; }

1
src/tcp_proxy.h

@ -63,6 +63,7 @@ struct proxy_conn {
int active; // 1 = outgoing connection
struct tcp_proxy_mapping* eim_mapping;
uint64_t remote_stream_id; // stream ID for remote proxy (0=local)
void* uip_retry_timer; // timer for backpressure retry
};
// Main TCP proxy module

Loading…
Cancel
Save