Browse Source

etcp_route_send: add force parameter to bypass send_q limit for control messages

CLOSE/ERROR/CONNECT and routing messages use force=1 to never be
dropped by backpressure. Regular DATA uses force=0 (default).
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
e765d62e43
  1. 10
      src/etcp_router.c
  2. 2
      src/etcp_router.h
  3. 4
      src/nat_transport.c
  4. 6
      src/proxy/icmp_proxy.c
  5. 14
      src/proxy/tcp_proxy_client.c
  6. 16
      src/proxy/tcp_proxy_server.c
  7. 4
      src/proxy/udp_proxy.c
  8. 2
      src/routing.c
  9. 6
      tests/test_etcp_router.c
  10. 2
      tests/test_etcp_router_unit.c
  11. 4
      tests/test_nat_transport.c

10
src/etcp_router.c

@ -184,8 +184,8 @@ static void router_send_resume_cb(void* arg) {
router_drain_send_q(rconn);
}
static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len) {
if (queue_entry_count(rconn->send_q) >= ROUTER_MAX_SEND_Q_PACKETS) {
static int router_enqueue_send(struct ETCP_ROUTER_CONN* rconn, const uint8_t* payload, size_t pl_len, int force) {
if (!force && queue_entry_count(rconn->send_q) >= ROUTER_MAX_SEND_Q_PACKETS) {
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_send_q: FULL svc_id=%u count=%d — backpressure",
rconn->svc_id, queue_entry_count(rconn->send_q));
return -1;
@ -534,7 +534,7 @@ int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) {
return 0;
}
int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry) {
int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry, int force) {
if (!inst || !entry) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "etcp_route_send: NULL inst=%p entry=%p", (void*)inst, (void*)entry);
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return -1;
@ -565,7 +565,7 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_
if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= ROUTER_MAX_INFLIGHT) {
int eq_ret = router_enqueue_send(rconn, entry->dgram + 1, payload_len);
int eq_ret = router_enqueue_send(rconn, entry->dgram + 1, payload_len, force);
queue_dgram_free(entry); queue_entry_free(entry);
return eq_ret;
}
@ -708,7 +708,7 @@ int etcp_router_conn_send(struct ETCP_ROUTER_CONN* rconn,
return -1;
}
if ((int32_t)(rconn->tx_seq - rconn->tx_acked) >= ROUTER_MAX_INFLIGHT) {
return router_enqueue_send(rconn, data, len);
return router_enqueue_send(rconn, data, len, 0);
}
return router_send_one(rconn, data, len);
}

2
src/etcp_router.h

@ -85,7 +85,7 @@ int etcp_router_bind(struct UTUN_INSTANCE* inst, uint8_t svc_id, etcp_recv_fn c
int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id);
// Отправить сервисный пакет (авто-conn, seq, inflight-контроль через send_q)
int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry);
int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry, int force);
// Найти/создать состояние seq-подключения по (remote_node_id, svc_id)
struct ETCP_ROUTER_CONN* etcp_router_conn_get(struct UTUN_INSTANCE* inst,

4
src/nat_transport.c

@ -47,7 +47,7 @@ static void nat_transport_client_tun_out_cb(struct ll_queue* q, void* arg) {
queue_dgram_free(pkt); queue_entry_free(pkt);
queue_resume_callback(q);
int ret = etcp_route_send(inst, tr->nat_via_node_id, new_entry);
int ret = etcp_route_send(inst, tr->nat_via_node_id, new_entry, 0);
if (ret != 0) {
DEBUG_WARN(DEBUG_CATEGORY_NAT, "NAT client: etcp_route_send to provider %016llx failed",
(unsigned long long)tr->nat_via_node_id);
@ -90,7 +90,7 @@ static void nat_transport_provider_tun_out_cb(struct ll_queue* q, void* arg) {
queue_dgram_free(pkt); queue_entry_free(pkt);
queue_resume_callback(q);
int send_ret = etcp_route_send(inst, entry->src_node_id, new_entry);
int send_ret = etcp_route_send(inst, entry->src_node_id, new_entry, 0);
if (send_ret != 0) {
DEBUG_WARN(DEBUG_CATEGORY_NAT, "NAT provider: etcp_route_send back to node %016llx failed",
(unsigned long long)entry->src_node_id);

6
src/proxy/icmp_proxy.c

@ -120,7 +120,7 @@ static void raw_read_cb(socket_t sock, void* arg) {
memcpy(e->dgram + 20, &icmp_hdr->icmp_seq, 2);
if (payload_len > 0) memcpy(e->dgram + ICMP_PROXY_HDR_SIZE, payload, payload_len);
e->len = ICMP_PROXY_HDR_SIZE + payload_len;
int ret = etcp_route_send(g_icmp_ctx->inst, r->client_node_id, e);
int ret = etcp_route_send(g_icmp_ctx->inst, r->client_node_id, e, 0);
if (ret != 0) DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "icmp_proxy: etcp_route_send reply failed: %d", ret);
else DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "icmp_proxy: reply forwarded to client %016llx", (unsigned long long)r->client_node_id);
u_free(r);
@ -158,7 +158,7 @@ static void exit_handle_request(struct ETCP_CONN* conn, struct ll_entry* entry)
memcpy(e->dgram + 20, &echo_seq, 2);
if (payload_len > 0) memcpy(e->dgram + ICMP_PROXY_HDR_SIZE, payload, payload_len);
e->len = ICMP_PROXY_HDR_SIZE + payload_len;
etcp_route_send(inst, client_node_id, e);
etcp_route_send(inst, client_node_id, e, 0);
} else queue_entry_free(e);
}
} else {
@ -223,7 +223,7 @@ int icmp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id,
memcpy(e->dgram + 20, &echo_seq, 2);
if (payload_len > 0) memcpy(e->dgram + ICMP_PROXY_HDR_SIZE, payload, payload_len);
e->len = ICMP_PROXY_HDR_SIZE + payload_len;
return etcp_route_send(inst, exit_node_id, e);
return etcp_route_send(inst, exit_node_id, e, 0);
}
// ====================================================================

14
src/proxy/tcp_proxy_client.c

@ -41,7 +41,7 @@ static err_t tcp_proxy_client_output_cb(void *arg, struct pbuf *p, uint32_t src_
static void tcp_proxy_client_feed_from_transport(struct tcp_proxy_client_conn *pc);
static struct ll_entry* tcp_proxy_client_entry_from_data(struct memory_pool* pool, const uint8_t* data, uint16_t len);
static void tcp_proxy_client_conn_free(struct tcp_proxy_client_conn *pc);
static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len);
static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force);
static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len);
static void tcp_proxy_client_tx_waiter_cb(struct ll_queue* q, void* arg);
@ -64,7 +64,7 @@ static struct ll_entry* tcp_proxy_client_entry_from_data(struct memory_pool* poo
// Протокол: сборка и отправка сообщений прокси через ETCP
// ====================================================================
static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
uint32_t sid, const uint8_t* data, size_t len) {
uint32_t sid, const uint8_t* data, size_t len, int force) {
struct ll_entry* e = queue_entry_new(0);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_client_send_msg: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
@ -74,7 +74,7 @@ static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, u
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len;
return etcp_route_send(inst, dst, e);
return etcp_route_send(inst, dst, e, force);
}
static int tcp_proxy_client_send_connect(struct tcp_proxy_client_conn* pc) {
@ -84,18 +84,18 @@ static int tcp_proxy_client_send_connect(struct tcp_proxy_client_conn* pc) {
pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3],
ntohs(pc->dest_port), (unsigned long long)pc->proxy->via_node_id);
return tcp_proxy_client_send_msg(pc->proxy->inst, pc->proxy->via_node_id,
TCP_PROXY_SUBCMD_CONNECT, pc->stream_id, buf, 6);
TCP_PROXY_SUBCMD_CONNECT, pc->stream_id, buf, 6, 1);
}
static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len) {
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY SEND sid=%08x len=%u", pc->stream_id, len);
return tcp_proxy_client_send_msg(pc->proxy->inst, pc->proxy->via_node_id,
TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len);
TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len, 0);
}
static int tcp_proxy_client_send_close(struct tcp_proxy_client_conn* pc) {
int ret = tcp_proxy_client_send_msg(pc->proxy->inst, pc->proxy->via_node_id,
TCP_PROXY_SUBCMD_CLOSE, pc->stream_id, NULL, 0);
TCP_PROXY_SUBCMD_CLOSE, pc->stream_id, NULL, 0, 1);
if (ret < 0) { pc->close_pending = 1; return -1; }
pc->close_pending = 0;
pc->close_sent = 1;
@ -104,7 +104,7 @@ static int tcp_proxy_client_send_close(struct tcp_proxy_client_conn* pc) {
static int tcp_proxy_client_send_error(struct tcp_proxy_client_conn* pc) {
int ret = tcp_proxy_client_send_msg(pc->proxy->inst, pc->proxy->via_node_id,
TCP_PROXY_SUBCMD_ERROR, pc->stream_id, NULL, 0);
TCP_PROXY_SUBCMD_ERROR, pc->stream_id, NULL, 0, 1);
if (ret < 0) { pc->close_pending = 1; return -1; }
pc->close_pending = 0;
pc->close_sent = 1;

16
src/proxy/tcp_proxy_server.c

@ -34,7 +34,7 @@ static void pause_resume_cb(struct ll_queue* q, void* arg);
static void close_retry_cb(void* arg);
static void diag_timer_cb(void* arg);
static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
uint32_t sid, const uint8_t* data, size_t len);
uint32_t sid, const uint8_t* data, size_t len, int force);
static void send_close(struct tcp_proxy_server_conn* rc);
static void send_error(struct tcp_proxy_server_conn* rc);
void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc);
@ -48,7 +48,7 @@ static inline int write_pending(struct tcp_conn* tc) {
// ====================================================================
static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
uint32_t sid, const uint8_t* data, size_t len) {
uint32_t sid, const uint8_t* data, size_t len, int force) {
struct ll_entry* e = queue_entry_new(0);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: queue_entry_new failed subcmd=%02x sid=%08x", subcmd, sid); return -1; }
e->dgram = u_malloc(TCP_PROXY_HDR_SIZE + len);
@ -58,7 +58,7 @@ static int send_msg(struct UTUN_INSTANCE* inst, uint64_t dst, uint8_t subcmd,
memcpy(e->dgram + 2, &sid, 4);
if (len > 0) memcpy(e->dgram + TCP_PROXY_HDR_SIZE, data, len);
e->len = TCP_PROXY_HDR_SIZE + len;
return etcp_route_send(inst, dst, e);
return etcp_route_send(inst, dst, e, force);
}
static void close_retry_cb(void* arg) {
@ -69,7 +69,7 @@ static void close_retry_cb(void* arg) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
uint8_t subcmd = (rc->tc && rc->tc->error) ? TCP_PROXY_SUBCMD_ERROR : TCP_PROXY_SUBCMD_CLOSE;
if (send_msg(inst, rc->peer_node_id, subcmd, rc->stream_id, NULL, 0) < 0) {
if (send_msg(inst, rc->peer_node_id, subcmd, rc->stream_id, NULL, 0, 1) < 0) {
rc->close_backoff = rc->close_backoff < 5000 ? rc->close_backoff * 2 : 5000;
rc->close_timer = uasync_set_timeout(rc->ua, rc->close_backoff, rc, close_retry_cb, "tps_close_retry");
return;
@ -80,7 +80,7 @@ static void close_retry_cb(void* arg) {
static void send_close(struct tcp_proxy_server_conn* rc) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
if (send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0) < 0) {
if (send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, NULL, 0, 1) < 0) {
rc->close_pending = 1;
rc->close_backoff = 50;
if (!rc->close_timer)
@ -91,7 +91,7 @@ static void send_close(struct tcp_proxy_server_conn* rc) {
static void send_error(struct tcp_proxy_server_conn* rc) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (!inst) return;
if (send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0) < 0) {
if (send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_ERROR, rc->stream_id, NULL, 0, 1) < 0) {
rc->close_pending = 1;
rc->close_backoff = 50;
if (!rc->close_timer)
@ -161,7 +161,7 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
struct ll_entry* e = queue_data_get(q);
if (!e) { queue_resume_callback(q); return; }
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len);
int ret = send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, e->dgram, e->len, 0);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%u total=%d",
(int)rc->tc->sock, rc->stream_id, e->len, rc->ctx ? rc->ctx->conn_count : 0);
memory_pool_free(rc->tc->data_pool, e->dgram);
@ -258,7 +258,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
if (ret < 0 && errno != EINPROGRESS) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: connect() to %d.%d.%d.%d:%d failed: %s",
dest_ip[0], dest_ip[1], dest_ip[2], dest_ip[3], ntohs(dest_port), strerror(errno));
send_msg(inst, src_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
send_msg(inst, src_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0, 1);
ctx->conn_count--; tcp_conn_destroy(rc->tc); u_free(rc);
queue_dgram_free(entry); queue_entry_free(entry); return -1;
}

4
src/proxy/udp_proxy.c

@ -64,7 +64,7 @@ static void flow_read_cb(socket_t sock, void* arg) {
memcpy(e->dgram + UDP_PROXY_HDR_SIZE, buf, n);
e->len = UDP_PROXY_HDR_SIZE + n;
etcp_route_send(g_udp_ctx->inst, f->client_node_id, e);
etcp_route_send(g_udp_ctx->inst, f->client_node_id, e, 0);
}
// ====================================================================
@ -169,7 +169,7 @@ int udp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id,
memcpy(e->dgram + 20, &dst_port, 2);
if (payload_len > 0) memcpy(e->dgram + UDP_PROXY_HDR_SIZE, payload, payload_len);
e->len = UDP_PROXY_HDR_SIZE + payload_len;
return etcp_route_send(inst, exit_node_id, e);
return etcp_route_send(inst, exit_node_id, e, 0);
}
// ====================================================================

2
src/routing.c

@ -147,7 +147,7 @@ void route_pkt(struct UTUN_INSTANCE* instance, struct ll_entry* entry, uint64_t
} else {
DEBUG_TRACE(DEBUG_CATEGORY_ROUTING, "route_pkt: sending %zu bytes to node %016llx dst=%s",
ip_len, (unsigned long long)nq->node.node_id, ip_to_str(&addr, AF_INET).str);
int send_err = etcp_route_send(instance, nq->node.node_id, entry);
int send_err = etcp_route_send(instance, nq->node.node_id, entry, 1);
if (send_err != 0) {
DEBUG_WARN(DEBUG_CATEGORY_ROUTING, "route_pkt: etcp_route_send failed: dst=%s err=%d",
ip_to_str(&addr, AF_INET).str, send_err);

6
tests/test_etcp_router.c

@ -136,7 +136,7 @@ static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) {
memcpy(buf + 10, expected_data, data_len);
struct ll_entry* re = queue_entry_new(0);
if (re) { re->dgram = u_malloc(10 + data_len); memcpy(re->dgram, buf, 10 + data_len); re->len = 10 + data_len;
etcp_route_send(srv, client_node_id, re); reply_sent++; }
etcp_route_send(srv, client_node_id, re, 0); reply_sent++; }
}
} else {
queue_entry_free(entry); queue_dgram_free(entry);
@ -187,7 +187,7 @@ static int test_loopback(void) {
uint8_t data[6] = { 0xF0, 0xAA, 0xBB, 0xCC, 0x00, 0x00 };
struct ll_entry* e = queue_entry_new(0);
e->dgram = u_malloc(6); memcpy(e->dgram, data, 6); e->len = 6;
etcp_route_send(srv, srv->node_id, e); // loopback
etcp_route_send(srv, srv->node_id, e, 0); // loopback
etcp_router_unbind(srv, 0xF0);
if (loop_rcvd != 1 || !loop_ok) {
printf("[FAIL] loopback: rcvd=%d ok=%d\n", loop_rcvd, loop_ok);
@ -249,7 +249,7 @@ static void monitor(void* arg) {
memcpy(buf + 10, expected_data, data_len);
struct ll_entry* e = queue_entry_new(0);
if (e) { e->dgram = u_malloc(10 + data_len); memcpy(e->dgram, buf, 10 + data_len); e->len = 10 + data_len;
if (etcp_route_send(cli, server_node_id, e) == 0) fwd_sent++;
if (etcp_route_send(cli, server_node_id, e, 0) == 0) fwd_sent++;
else { queue_entry_free(e); queue_dgram_free(e); }
}
}

2
tests/test_etcp_router_unit.c

@ -502,7 +502,7 @@ static int test_loopback(void) {
memcpy(e->dgram, buf, 4);
e->len = 4;
if (etcp_route_send(&inst, inst.node_id, e) != 0) FAIL("loopback send failed");
if (etcp_route_send(&inst, inst.node_id, e, 0) != 0) FAIL("loopback send failed");
if (rx.delivered != 1) FAIL("loopback not delivered");
if (rx.errors != 0) FAIL("loopback data mismatch");

4
tests/test_nat_transport.c

@ -299,7 +299,7 @@ static int test_provider_egress(void) {
entry->dgram = dgram;
entry->len = total;
int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry);
int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry, 0);
if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NAT, "etcp_route_send failed");
queue_dgram_free(entry);
@ -468,7 +468,7 @@ static int test_full_roundtrip(void) {
struct ll_entry* entry = queue_entry_new(0);
entry->dgram = dgram;
entry->len = total;
int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry);
int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry, 0);
if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_NAT, "etcp_route_send failed"); queue_dgram_free(entry); queue_entry_free(entry); return 0; }
// Poll for provider to process egress

Loading…
Cancel
Save