Browse Source

fix use-after-free in queue_entry_free/dgram_free ordering + add tcp_proxy stress test + proxy diagnostics logging

- Fix queue_entry_free before queue_dgram_free in all proxy files (UAF bug causing segfault)
- Add PROXY/RP prefixed diagnostic logs for transparent proxy flow tracking
- Add tcp_proxy integration test: real utun instances, ETCP, remote_proxy, DNAT routing
- 7 test scenarios: basic_1mb, concurrent_5, concurrent_20, half_close, idle, conn_refused, stress (10 threads x 20 iterations)
- stress_client.py with --verbose per-thread phase timing (connect/send/recv)
congestion
Evgeny 4 months ago
parent
commit
a02c474b94
  1. 26
      src/etcp_router.c
  2. 6
      src/icmp_proxy.c
  3. 2
      src/nat_transport.c
  4. 32
      src/remote_proxy.c
  5. 44
      src/tcp_proxy.c
  6. 6
      src/udp_proxy.c
  7. 2
      tests/tcp_proxy_full/.gitignore
  8. 29
      tests/tcp_proxy_full/client.conf
  9. 76
      tests/tcp_proxy_full/echo_server.py
  10. 24
      tests/tcp_proxy_full/exit.conf
  11. 263
      tests/tcp_proxy_full/run_test.sh
  12. 122
      tests/tcp_proxy_full/stress_client.py
  13. 97
      tests/tcp_proxy_full/tcp_client.py

26
src/etcp_router.c

@ -12,14 +12,14 @@
// Обработчик ETCP_ID_SVC_ROUTE — вызывается etcp_int_recv на каждом узле
static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!conn || !entry) {
if (entry) { queue_entry_free(entry); queue_dgram_free(entry); }
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
return;
}
struct UTUN_INSTANCE* inst = conn->instance;
if (!inst || entry->len < SVC_ROUTE_HDR_SIZE) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: invalid packet inst=%p len=%zu min=%d",
(void*)inst, entry ? entry->len : 0, SVC_ROUTE_HDR_SIZE);
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
return;
}
uint8_t svc_id = entry->dgram[1];
@ -28,21 +28,21 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!inst->router_bindings.callbacks[svc_id]) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no handler for svc_id=%u dst=%016llx self=%016llx",
svc_id, (unsigned long long)dst_node_id, (unsigned long long)inst->node_id);
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
return;
}
size_t payload_len = entry->len - SVC_ROUTE_HDR_SIZE;
struct ll_entry* svc_entry = queue_entry_new(0);
if (!svc_entry) { queue_entry_free(entry); queue_dgram_free(entry); return; }
if (!svc_entry) { queue_dgram_free(entry); queue_entry_free(entry); return; }
svc_entry->len = 1 + payload_len;
svc_entry->dgram = u_malloc(svc_entry->len);
if (!svc_entry->dgram) { queue_entry_free(svc_entry); queue_entry_free(entry); queue_dgram_free(entry); return; }
if (!svc_entry->dgram) { queue_entry_free(svc_entry); queue_dgram_free(entry); queue_entry_free(entry); return; }
svc_entry->dgram[0] = svc_id;
if (payload_len > 0) memcpy(svc_entry->dgram + 1, entry->dgram + SVC_ROUTE_HDR_SIZE, payload_len);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: delivering svc_id=%u len=%zu to handler %p src=%016llx",
svc_id, payload_len, (void*)inst->router_bindings.callbacks[svc_id],
(unsigned long long)(*(uint64_t*)(entry->dgram + 10)));
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
inst->router_bindings.callbacks[svc_id](conn, svc_entry);
if (svc_id == ETCP_ID_ICMP_PROXY)
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_router: ICMP_PROXY reply delivered to handler self=%016llx",
@ -52,7 +52,7 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!next) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_router: no route to %016llx svc_id=%u, dropping",
(unsigned long long)dst_node_id, svc_id);
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
return;
}
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_router: forwarding svc_id=%u → %016llx via %s",
@ -101,12 +101,12 @@ int etcp_router_unbind(struct UTUN_INSTANCE* inst, uint8_t svc_id) {
int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_entry* entry) {
if (!inst || !entry) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: NULL inst=%p entry=%p", (void*)inst, (void*)entry);
if (entry) { queue_entry_free(entry); queue_dgram_free(entry); }
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
return -1;
}
if (!entry->dgram || entry->len < 1) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_route_send: empty entry");
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
return -1;
}
uint8_t svc_id = entry->dgram[0];
@ -117,14 +117,14 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_route_send: loopback svc_id=%u len=%zu", svc_id, payload_len);
if (inst->router_bindings.callbacks[svc_id])
inst->router_bindings.callbacks[svc_id](NULL, entry);
else { queue_entry_free(entry); queue_dgram_free(entry); }
else { queue_dgram_free(entry); queue_entry_free(entry); }
return 0;
}
// Упаковываем в routing header
size_t total_len = SVC_ROUTE_HDR_SIZE + payload_len;
uint8_t* new_dgram = u_malloc(total_len);
if (!new_dgram) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!new_dgram) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
new_dgram[0] = ETCP_ID_SVC_ROUTE;
new_dgram[1] = svc_id;
memcpy(new_dgram + 2, &dst_node_id, 8);
@ -132,11 +132,11 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, struct ll_
if (payload_len > 0) memcpy(new_dgram + SVC_ROUTE_HDR_SIZE, entry->dgram + 1, payload_len);
struct ll_entry* new_entry = queue_entry_new(0);
if (!new_entry) { u_free(new_dgram); queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!new_entry) { u_free(new_dgram); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
new_entry->dgram = new_dgram;
new_entry->len = total_len;
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, dst_node_id);
if (!conn) {

6
src/icmp_proxy.c

@ -166,7 +166,7 @@ static void exit_handle_request(struct ETCP_CONN* conn, struct ll_entry* entry)
}
drop:
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -186,7 +186,7 @@ static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry)
struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_icmp_ctx ? g_icmp_ctx->inst : NULL);
icmp_proxy_deliver_reply(inst, orig_src_ip, src_ip, echo_id, echo_seq,
entry->dgram + ICMP_PROXY_HDR_SIZE, entry->len - ICMP_PROXY_HDR_SIZE);
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -200,7 +200,7 @@ void icmp_proxy_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
uint8_t subcmd = entry->dgram[1];
if (subcmd == ICMP_PROXY_SUBCMD_REQUEST) { exit_handle_request(conn, entry); return; }
if (subcmd == ICMP_PROXY_SUBCMD_REPLY) { client_handle_reply(conn, entry); return; }
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================

2
src/nat_transport.c

@ -130,7 +130,7 @@ static void nat_transport_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry*
}
}
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ==================== Init / Destroy ====================

32
src/remote_proxy.c

@ -62,6 +62,7 @@ static void rp_sock_read_cb(socket_t sock, void* arg) {
if (!rc || rc->sock == SOCKET_INVALID) return;
uint8_t buf[8192]; ssize_t n = recv(rc->sock, buf, sizeof(buf), 0);
if (n > 0) {
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "RP RECV ← stream=%016llx len=%zd seq=%u", (unsigned long long)rc->stream_id, n, rc->send_seq);
struct UTUN_INSTANCE* inst = rc->ctx ? rc->ctx->inst : NULL;
if (inst) {
rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_DATA, rc->stream_id, rc->send_seq, buf, (size_t)n);
@ -132,7 +133,7 @@ static void rp_conn_free(struct remote_proxy_conn* rc) {
static void rp_standalone_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* i = conn ? conn->instance : (g_rp_ctx ? g_rp_ctx->inst : NULL);
if (!i || !entry || entry->dgram == NULL || entry->len < TCP_PROXY_HDR_SIZE) {
if (entry) { queue_entry_free(entry); queue_dgram_free(entry); }
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
return;
}
uint8_t subcmd = entry->dgram[1];
@ -140,8 +141,8 @@ static void rp_standalone_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry
uint64_t src = conn ? conn->peer_node_id : i->node_id;
if (subcmd == TCP_PROXY_SUBCMD_CONNECT) { remote_proxy_handle_connect(i, entry, sid, src); return; }
if (subcmd == TCP_PROXY_SUBCMD_DATA) { remote_proxy_handle_data(i, entry, sid); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { remote_proxy_handle_close(i, sid); queue_entry_free(entry); queue_dgram_free(entry); return; }
queue_entry_free(entry); queue_dgram_free(entry);
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { remote_proxy_handle_close(i, sid); queue_dgram_free(entry); queue_entry_free(entry); return; }
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -155,24 +156,24 @@ struct remote_proxy_conn* remote_proxy_find_conn(struct remote_proxy_ctx* ctx, u
int remote_proxy_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry,
uint64_t stream_id, uint64_t src_node_id) {
if (!inst || !inst->remote_proxy.enabled) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!inst || !inst->remote_proxy.enabled) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
struct remote_proxy_ctx* ctx = &inst->remote_proxy;
if (entry->len < TCP_PROXY_CONNECT_HDR_SIZE) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (entry->len < TCP_PROXY_CONNECT_HDR_SIZE) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
uint8_t* dest_ip = entry->dgram + TCP_PROXY_HDR_SIZE;
uint16_t dest_port = 0; memcpy(&dest_port, dest_ip + 4, 2);
struct remote_proxy_conn* rc = u_calloc(1, sizeof(struct remote_proxy_conn));
if (!rc) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!rc) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
rc->ctx = ctx; rc->stream_id = stream_id; rc->peer_node_id = src_node_id;
memcpy(rc->dest_ip, dest_ip, 4); rc->dest_port = dest_port;
rc->ua = inst->ua; rc->sock = SOCKET_INVALID;
rc->send_seq = 0; rc->recv_seq_init = 0;
rc->sock = socket(AF_INET, SOCK_STREAM, 0);
if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket() failed"); rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: socket() failed"); rp_conn_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
socket_set_nonblocking(rc->sock);
rc->read_id = uasync_add_socket_t(rc->ua, rc->sock, rp_sock_read_cb, rp_sock_write_cb, rp_sock_error_cb, rc);
if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: uasync_add_socket_t failed"); rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!rc->read_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: uasync_add_socket_t failed"); rp_conn_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }
struct sockaddr_in addr; memset(&addr, 0, sizeof(addr));
addr.sin_family = AF_INET; memcpy(&addr.sin_addr.s_addr, dest_ip, 4); addr.sin_port = dest_port;
@ -183,19 +184,19 @@ int remote_proxy_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* ent
if (ret < 0 && errno != EINPROGRESS) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: connect() failed: %s", strerror(errno));
rp_send_connected(inst, src_node_id, stream_id, 0, TCP_PROXY_CONNECTED_REFUSED);
rp_conn_free(rc); queue_entry_free(entry); queue_dgram_free(entry); return -1;
rp_conn_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1;
}
rc->next = ctx->conns; ctx->conns = rc;
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
return 0;
}
int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ll_entry* entry, uint64_t stream_id) {
if (!inst) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
struct remote_proxy_ctx* ctx = &inst->remote_proxy;
struct remote_proxy_conn* rc = remote_proxy_find_conn(ctx, stream_id);
if (!rc || rc->sock == SOCKET_INVALID || rc->connected != 1) { queue_entry_free(entry); queue_dgram_free(entry); return -1; }
if (!rc || rc->sock == SOCKET_INVALID || rc->connected != 1) { queue_dgram_free(entry); queue_entry_free(entry); return -1; }
uint16_t seq; memcpy(&seq, entry->dgram + 10, 2);
if (!rc->recv_seq_init) {
rc->recv_last_seq = seq; rc->recv_seq_init = 1;
@ -203,23 +204,24 @@ int remote_proxy_handle_data(struct UTUN_INSTANCE* inst, struct ll_entry* entry,
int16_t delta = (int16_t)(seq - rc->recv_last_seq);
if (delta <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "remote_proxy: duplicate DATA seq=%u stream=%016llx, discard", seq, (unsigned long long)stream_id);
queue_entry_free(entry); queue_dgram_free(entry); return 0;
queue_dgram_free(entry); queue_entry_free(entry); return 0;
}
if (delta > 1) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: seq gap seq=%u last=%u stream=%016llx, tearing down", seq, rc->recv_last_seq, (unsigned long long)stream_id);
rp_send_msg(inst, rc->peer_node_id, TCP_PROXY_SUBCMD_CLOSE, rc->stream_id, 0, NULL, 0);
rp_conn_free(rc);
queue_entry_free(entry); queue_dgram_free(entry); return -1;
queue_dgram_free(entry); queue_entry_free(entry); return -1;
}
rc->recv_last_seq = seq;
}
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "RP SEND → stream=%016llx seq=%u len=%zu", (unsigned long long)stream_id, seq, data_len);
uint8_t* data = entry->dgram + TCP_PROXY_HDR_SIZE;
ssize_t n = send(rc->sock, data, data_len, MSG_NOSIGNAL);
if (n < 0 && errno != EAGAIN && errno != EWOULDBLOCK) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "remote_proxy: send error %s", strerror(errno));
}
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
return 0;
}

44
src/tcp_proxy.c

@ -137,7 +137,7 @@ static err_t tcp_output_cb(void *arg, struct pbuf *p, uint32_t src_ip, uint32_t
pbuf_copy_partial(p, buf + 2, len, 0);
ssize_t n = write(proxy->ip_fd, buf, 2 + len); (void)n;
}
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "TCP OUT %u.%u.%u.%u→%u.%u.%u.%u len=%u",
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY TCP_OUT %u.%u.%u.%u:%u len=%u",
(uint8_t)(src_ip), (uint8_t)(src_ip>>8), (uint8_t)(src_ip>>16), (uint8_t)(src_ip>>24),
(uint8_t)(dst_ip), (uint8_t)(dst_ip>>8), (uint8_t)(dst_ip>>16), (uint8_t)(dst_ip>>24), len);
return LERR_OK;
@ -192,13 +192,13 @@ static void proxy_feed_from_transport(struct proxy_conn *pc) {
int entries_fed = 0;
while (1) {
uint16_t space = tcp_sndbuf(pc->pcb);
if (space < TCP_MSS / 2) { DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "FEED space=%u<ssthresh break active=%d", space, pc->active); break; }
if (space < TCP_MSS / 2) { DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY FEED → stream=%016llx space=%u break", (unsigned long long)pc->remote_stream_id, space); break; }
struct ll_entry *e = queue_data_get(pc->transport_to_uip);
if (!e) break;
uint16_t len = e->len;
if (len > space) len = space;
err_t ret = tcp_write(pc->pcb, e->dgram, len, TCP_WRITE_FLAG_COPY);
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "FEED tcp_write len=%u ret=%d active=%d cwnd=%u snd_wnd=%u unsent=%p", len, ret, pc->active, pc->pcb->cwnd, pc->pcb->snd_wnd, (void*)pc->pcb->unsent);
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY FEED → stream=%016llx len=%u ret=%d sndbuf=%u", (unsigned long long)pc->remote_stream_id, len, ret, pc->pcb->snd_buf);
if (ret == LERR_OK) {
sent_any = 1;
if (len >= e->len) { queue_dgram_free(e); queue_entry_free(e); }
@ -311,6 +311,7 @@ static err_t proxy_poll_cb(void *arg, struct tcp_pcb *pcb) {
struct proxy_conn *pc = (struct proxy_conn *)arg;
if (!pc) return LERR_OK;
if (pc->closing_tun && pc->closing_rem && pc->pcb) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE stream=%016llx tun=1 rem=1", (unsigned long long)pc->remote_stream_id);
struct tcp_pcb *save = pc->pcb;
tcp_arg(save, NULL); pc->pcb = NULL; tcp_close(save);
struct tcp_proxy *proxy = pc->proxy;
@ -396,9 +397,8 @@ static void tcp_proxy_raw_read(int fd, void* arg) {
uint32_t src_ip, dst_ip;
memcpy(&src_ip, pkt + 12, 4); memcpy(&dst_ip, pkt + 16, 4);
uint16_t tcp_len = ip_total - ip_hdr_len;
DEBUG_ERROR(DEBUG_CATEGORY_TRAFFIC, "TCP IN %u.%u.%u.%u→%u.%u.%u.%u iplen=%u",
(uint8_t)(src_ip), (uint8_t)(src_ip>>8), (uint8_t)(src_ip>>16), (uint8_t)(src_ip>>24),
(uint8_t)(dst_ip), (uint8_t)(dst_ip>>8), (uint8_t)(dst_ip>>16), (uint8_t)(dst_ip>>24), ip_total);
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY TCP_IN %u.%u.%u.%u:%u len=%u",
(uint8_t)(src_ip), (uint8_t)(src_ip>>8), (uint8_t)(src_ip>>16), (uint8_t)(src_ip>>24), ntohs(dport_net), tcp_len);
struct pbuf *pb = pbuf_alloc(PBUF_RAW, tcp_len);
if (pb) { pbuf_take(pb, pkt + ip_hdr_len, tcp_len); lwip_tcp_input(p->lwip, pb, src_ip, dst_ip); }
}
@ -645,15 +645,16 @@ static struct proxy_conn* find_pc_by_stream(struct tcp_proxy* p, uint64_t stream
static void handle_connected(struct tcp_proxy* p, uint64_t stream_id, struct ll_entry* entry) {
struct proxy_conn* pc = find_pc_by_stream(p, stream_id);
if (!pc || !pc->transport) {
queue_entry_free(entry); queue_dgram_free(entry);
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CONNECTED stream=%016llx pc=%p transport=%p — drop", (unsigned long long)stream_id, (void*)pc, pc ? pc->transport : NULL);
queue_dgram_free(entry); queue_entry_free(entry);
return;
}
struct etcp_transport* et = (struct etcp_transport*)pc->transport;
if (et->base.ops != &etcp_transport_ops) { queue_entry_free(entry); queue_dgram_free(entry); return; }
if (et->base.ops != &etcp_transport_ops) { queue_dgram_free(entry); queue_entry_free(entry); return; }
if (entry->len >= TCP_PROXY_CONNECTED_HDR_SIZE) {
uint8_t status = entry->dgram[TCP_PROXY_HDR_SIZE + 2];
if (status == TCP_PROXY_CONNECTED_OK) { et->connected = 1;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote connected stream=%016llx", (unsigned long long)stream_id);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY CONNECTED stream=%016llx status=OK", (unsigned long long)stream_id);
struct ll_entry* e;
while ((e = queue_data_get(pc->uip_to_transport)) != NULL) {
if (e->len > 0) {
@ -665,7 +666,7 @@ static void handle_connected(struct tcp_proxy* p, uint64_t stream_id, struct ll_
else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: remote refused stream=%016llx", (unsigned long long)stream_id);
pc->closing_rem = 1; }
}
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -673,7 +674,7 @@ static void handle_connected(struct tcp_proxy* p, uint64_t stream_id, struct ll_
// ====================================================================
void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) {
if (entry) { queue_entry_free(entry); queue_dgram_free(entry); }
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
return;
}
uint8_t subcmd = entry->dgram[1];
@ -690,7 +691,7 @@ void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct remote_proxy_conn* rc = remote_proxy_find_conn(&inst->remote_proxy, stream_id);
if (rc) {
if (subcmd == TCP_PROXY_SUBCMD_DATA) { remote_proxy_handle_data(inst, entry, stream_id); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { remote_proxy_handle_close(inst, stream_id); queue_entry_free(entry); queue_dgram_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { remote_proxy_handle_close(inst, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
}
}
if (proxy) {
@ -707,22 +708,23 @@ void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
int16_t delta = (int16_t)(seq - et->recv_last_seq);
if (delta <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy: duplicate DATA seq=%u stream=%016llx, discard", seq, (unsigned long long)stream_id);
queue_entry_free(entry); queue_dgram_free(entry); return;
queue_dgram_free(entry); queue_entry_free(entry); return;
}
if (delta > 1) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy: seq gap seq=%u last=%u stream=%016llx, tearing down", seq, et->recv_last_seq, (unsigned long long)stream_id);
pc->closing_rem = 1;
if (pc->transport) { pc->transport->ops->destroy(pc->transport); pc->transport = NULL; }
queue_entry_free(entry); queue_dgram_free(entry); return;
queue_dgram_free(entry); queue_entry_free(entry); return;
}
et->recv_last_seq = seq;
}
}
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
if (data_len > 0) {
struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len);
if (e) queue_data_put(pc->transport_to_uip, e);
proxy_feed_from_transport(pc);
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
DEBUG_INFO(DEBUG_CATEGORY_TRAFFIC, "PROXY DATA ← stream=%016llx seq=%u len=%zu", (unsigned long long)stream_id, seq, data_len);
if (data_len > 0) {
struct ll_entry* e = entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len);
if (e) queue_data_put(pc->transport_to_uip, e);
proxy_feed_from_transport(pc);
}
}
}
}
@ -734,7 +736,7 @@ void tcp_proxy_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
}
}
}
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================

6
src/udp_proxy.c

@ -107,7 +107,7 @@ static void exit_handle_data(struct ETCP_CONN* conn, struct ll_entry* entry) {
sendto(f->sock, payload, payload_len, 0, (struct sockaddr*)&addr, sizeof(addr));
drop:
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -128,7 +128,7 @@ static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry)
struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_udp_ctx ? g_udp_ctx->inst : NULL);
udp_proxy_deliver_reply(inst, src_ip, src_port, dst_ip, dst_port, payload, payload_len);
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
@ -145,7 +145,7 @@ void udp_proxy_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
else client_handle_reply(conn, entry);
return;
}
queue_entry_free(entry); queue_dgram_free(entry);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================

2
tests/tcp_proxy_full/.gitignore vendored

@ -0,0 +1,2 @@
/log/
*.log

29
tests/tcp_proxy_full/client.conf

@ -0,0 +1,29 @@
[global]
my_node_id=0xAAAA000000000002
my_private_key=4813d31d28b7e9829247f488c6be7672f2bdf61b2508333128e386d1759afed2
my_public_key=c594f33c91f3a2222795c2c110c527bf214ad1009197ce14556cb13df3c461b3c373bed8f205a8dd1fc0c364f90bf471d7c6f5db49564c33e4235d268569ac71
tun_ip=10.200.20.1/24
tun_ifname=tun_test_cli
debug_level=error
[server: s1]
addr=127.0.0.1:15002
type=public
[client: c1]
keepalive=1
link=s1:127.0.0.1:15001
peer_public_key=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9
[tcp_proxy]
enabled=yes
tun_name=tun_test_proxy
tun_ip=10.200.30.1
via_node=0xAAAA000000000001
[debug]
connection=info
socket=info
general=info
traffic=info
tun=info

76
tests/tcp_proxy_full/echo_server.py

@ -0,0 +1,76 @@
#!/usr/bin/env python3
"""TCP echo server for tcp_proxy integration test.
Accepts connections, reads all data, sends back.
Supports --read-timeout to close after idle (for half-close tests).
"""
import socket
import sys
import threading
import signal
running = True
def handle(conn: socket.socket, read_timeout: float):
try:
while True:
if read_timeout > 0:
conn.settimeout(read_timeout)
try:
data = conn.recv(65536)
except socket.timeout:
break
if not data:
break
conn.sendall(data)
except OSError:
pass
finally:
conn.close()
def main():
if len(sys.argv) < 3:
print(f"Usage: {sys.argv[0]} <host> <port> [--read-timeout SEC]", file=sys.stderr)
sys.exit(1)
host = sys.argv[1]
port = int(sys.argv[2])
read_timeout = 0.0
for i in range(3, len(sys.argv)):
if sys.argv[i] == "--read-timeout" and i + 1 < len(sys.argv):
read_timeout = float(sys.argv[i + 1])
global running
def _sigterm(sig, frame):
global running
running = False
signal.signal(signal.SIGTERM, _sigterm)
signal.signal(signal.SIGINT, _sigterm)
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
s.bind((host, port))
s.listen(128)
print(f"ECHO: listening on {host}:{port}" + (f" timeout={read_timeout}s" if read_timeout > 0 else ""), flush=True)
while running:
try:
s.settimeout(1.0)
conn, addr = s.accept()
except socket.timeout:
continue
except OSError:
break
t = threading.Thread(target=handle, args=(conn, read_timeout), daemon=True)
t.start()
s.close()
if __name__ == "__main__":
main()

24
tests/tcp_proxy_full/exit.conf

@ -0,0 +1,24 @@
[global]
my_node_id=0xAAAA000000000001
my_private_key=67b705a92b41bcaae105af2d6a17743faa7b26ccebba8b3b9b0af05e9cd1d5fb
my_public_key=1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9
tun_ip=10.200.10.1/24
tun_ifname=tun_test_exit
debug_level=error
[server: s1]
addr=127.0.0.1:15001
type=public
[allowed_keys]
allow_all=1
[remote_proxy]
enabled=yes
[debug]
connection=info
socket=info
general=info
traffic=info
tun=info

263
tests/tcp_proxy_full/run_test.sh

@ -0,0 +1,263 @@
#!/bin/bash
# run_test.sh — интеграционный тест tcp_proxy (transparent proxy)
#
# Трафик:
# tcp_client.py (SO_MARK=1, dest=10.200.100.N) → table 100 → tun_test_proxy
# → tcp_proxy(lwIP) → ETCP → exit(remote_proxy)
# → connect(10.200.100.N) → iptables DNAT (!mark=1) → 127.0.0.1 → echo → обратно
#
# Запуск: sudo ./run_test.sh
set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)"
UTUN_BIN="$SCRIPT_DIR/../../utun"
LOG_DIR="$SCRIPT_DIR/log"
CLIENT_PY="$SCRIPT_DIR/tcp_client.py"
ECHO_PY="$SCRIPT_DIR/echo_server.py"
STRESS_PY="$SCRIPT_DIR/stress_client.py"
ECHO_PORT=19090
HC_PORT=19091
REFUSED_PORT=19099
N_IP=20
PASS=0; FAIL=0
# -------- setup / cleanup --------
setup_net() {
echo "=== Setting up networking ==="
sysctl -w net.ipv4.conf.all.rp_filter=0
# policy routing: marked packets → TUN (priority 199)
ip rule add fwmark 1 priority 199 table 100 2>/dev/null || true
local gw
gw=$(ip route show default | awk '/via/ {print $3; exit}')
[ -n "$gw" ] && ip route replace default via "$gw" table 100 2>/dev/null || true
# DNAT: unmarked packets to fake IPs → real localhost
for i in $(seq 1 $N_IP); do
iptables -t nat -C OUTPUT -d "10.200.100.${i}" -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$ECHO_PORT" 2>/dev/null \
|| iptables -t nat -A OUTPUT -d "10.200.100.${i}" -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$ECHO_PORT"
done
echo " Done"
}
setup_tun_route() {
for i in $(seq 1 60); do
if [ -e "/proc/sys/net/ipv4/conf/tun_test_proxy/rp_filter" ]; then
sysctl -w net.ipv4.conf.tun_test_proxy.rp_filter=0
ip route replace 10.200.100.0/24 dev tun_test_proxy table 100
ip route replace 10.200.100.0/24 dev tun_test_proxy # main table anti-martian
echo " tun route ready (${i}x0.3s)"
return 0
fi
sleep 0.3
done
echo " WARN: tun_test_proxy not found after 18s"
}
cleanup_net() {
echo "=== Cleaning up networking ==="
for i in $(seq 1 $N_IP); do
iptables -t nat -D OUTPUT -d "10.200.100.${i}" -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$ECHO_PORT" 2>/dev/null || true
done
ip rule del fwmark 1 priority 199 table 100 2>/dev/null || true
ip route del 10.200.100.0/24 dev tun_test_proxy table 100 2>/dev/null || true
ip route del 10.200.100.0/24 dev tun_test_proxy 2>/dev/null || true
ip route flush table 100 2>/dev/null || true
sysctl -w net.ipv4.conf.all.rp_filter=1
sysctl -w net.ipv4.conf.tun_test_proxy.rp_filter=1 2>/dev/null || true
echo " Done"
}
cleanup() {
echo ""; echo "=== Cleanup ==="
kill $EXIT_PID 2>/dev/null || true
kill $CLIENT_PID 2>/dev/null || true
kill $ECHO_PID 2>/dev/null || true
wait $EXIT_PID 2>/dev/null || true
wait $CLIENT_PID 2>/dev/null || true
wait $ECHO_PID 2>/dev/null || true
sleep 0.5; cleanup_net
}
wait_for_etcp() {
for i in $(seq 1 30); do
grep -q "Connection established" "$LOG_DIR/exit_utun.log" 2>/dev/null && return 0
sleep 0.5
done
return 1
}
# -------- test runners --------
run_test() {
local name=$1; shift
local t_start log rc elapsed
t_start=$(date +%s%3N)
echo -n " $name ... "
log="$LOG_DIR/tcp_client_${name}.log"
if python3 "$CLIENT_PY" "$@" --name "$name" >"$log" 2>&1; then
rc=0
else
rc=$?
fi
elapsed=$(($(date +%s%3N) - t_start))
if [ $rc -eq 0 ]; then
echo "PASS (${elapsed}ms)"
tail -1 "$log"
return 0
else
echo "FAIL (rc=$rc, ${elapsed}ms)"
tail -3 "$log"
return 1
fi
}
run_test_neg() {
local name=$1; shift
local t_start log rc elapsed
t_start=$(date +%s%3N)
echo -n " $name ... "
log="$LOG_DIR/tcp_client_${name}.log"
if ! python3 "$CLIENT_PY" "$@" --name "$name" >"$log" 2>&1; then
rc=0
else
rc=1
fi
elapsed=$(($(date +%s%3N) - t_start))
if [ $rc -eq 0 ]; then
echo "PASS (expected error, ${elapsed}ms)"; return 0
else
echo "FAIL (expected error, got success, ${elapsed}ms)"; return 1
fi
}
run_concurrent() {
local name=$1 count=$2 size=$3; shift 3
local t_start pids ok elapsed
t_start=$(date +%s%3N)
pids=()
for i in $(seq 1 $count); do
python3 "$CLIENT_PY" \
--host "10.200.100.${i}" --port "$ECHO_PORT" \
--size "$size" --verify --timeout 15 \
--name "${name}_${i}" \
>"$LOG_DIR/tcp_client_${name}_${i}.log" 2>&1 &
pids+=($!)
done
ok=0
for pid in "${pids[@]}"; do wait "$pid"; [ $? -eq 0 ] && ((++ok)); done
elapsed=$(($(date +%s%3N) - t_start))
echo -n " $name ($count) ... "
if [ "$ok" -eq "$count" ]; then
echo "PASS ($ok/$count, ${elapsed}ms)"; return 0
else
echo "FAIL ($ok/$count, ${elapsed}ms)"
local fail=''
for i in $(seq 1 $count); do
grep -q "\[PASS\]" "$LOG_DIR/tcp_client_${name}_${i}.log" 2>/dev/null || fail+=" ${i}"
done
echo " failed:$fail"; return 1
fi
}
run_stress() {
local t_start elapsed
t_start=$(date +%s%3N)
echo -n " stress_sessions ... "
local log="$LOG_DIR/stress.log"
if python3 "$STRESS_PY" \
--host 10.200.100.1 --port-base "$ECHO_PORT" --ports 20 \
--threads 10 --iters 20 --min-size 1024 --max-size 65536 --timeout 60 --verbose \
>"$log" 2>&1; then
elapsed=$(($(date +%s%3N) - t_start))
echo "PASS (${elapsed}ms)"
head -1 "$log"; return 0
else
elapsed=$(($(date +%s%3N) - t_start))
echo "FAIL (${elapsed}ms)"
tail -5 "$log"; return 1
fi
}
# -------- main --------
if [ "$(id -u)" -ne 0 ]; then echo "ERROR: must be root"; exit 1; fi
[ -x "$UTUN_BIN" ] || { echo "ERROR: utun not found. Build first."; exit 1; }
echo "=== tcp_proxy full integration test ==="
mkdir -p "$LOG_DIR"; rm -f "$LOG_DIR"/*.log
setup_net
echo "Starting echo server on 127.0.0.1:$ECHO_PORT ..."
python3 "$ECHO_PY" 127.0.0.1 "$ECHO_PORT" >"$LOG_DIR/echo.log" 2>&1 &
ECHO_PID=$!; sleep 0.3
echo "Starting utun exit ..."
"$UTUN_BIN" -c "$SCRIPT_DIR/exit.conf" -f -l "$LOG_DIR/exit_utun.log" >"$LOG_DIR/exit_stdout.log" 2>&1 &
EXIT_PID=$!
echo "Starting utun client ..."
"$UTUN_BIN" -c "$SCRIPT_DIR/client.conf" -f -l "$LOG_DIR/client_utun.log" >"$LOG_DIR/client_stdout.log" 2>&1 &
CLIENT_PID=$!
trap cleanup EXIT
setup_tun_route
echo "Waiting for ETCP connection ..."
wait_for_etcp || { echo "ERROR: ETCP timeout"; tail -20 "$LOG_DIR/client_utun.log"; exit 1; }
echo "ETCP ready"; sleep 1
echo ""; echo "=== Running tests ==="
# 1: stress first (catch resource leaks)
run_stress && ((++PASS)) || ((++FAIL))
sleep 2
# 2: basic_1mb
run_test basic_1mb --host 10.200.100.1 --port "$ECHO_PORT" --size 1048576 --verify --timeout 10 \
&& ((++PASS)) || ((++FAIL))
sleep 1
# 3: concurrent_5
run_concurrent concurrent_5 5 204800 && ((++PASS)) || ((++FAIL))
sleep 1
# 4: concurrent_20
run_concurrent concurrent_20 20 51200 && ((++PASS)) || ((++FAIL))
sleep 1
# 5: half_close
echo " half_close: one-shot echo on $HC_PORT ..."
python3 "$ECHO_PY" 127.0.0.1 "$HC_PORT" --read-timeout 2 >"$LOG_DIR/echo_hc.log" 2>&1 &
HC_PID=$!; sleep 0.3
run_test half_close --host 10.200.100.1 --port "$HC_PORT" --size 65536 --half-close --timeout 8 \
&& ((++PASS)) || ((++FAIL))
kill $HC_PID 2>/dev/null || true; wait $HC_PID 2>/dev/null || true
sleep 1
# 6: idle
run_test idle --host 10.200.100.1 --port "$ECHO_PORT" --size 32768 --verify --sleep-before 1 --timeout 5 \
&& ((++PASS)) || ((++FAIL))
sleep 1
# 7: conn_refused
run_test_neg conn_refused --host 10.200.100.99 --port "$REFUSED_PORT" --size 1 --verify --timeout 5 \
&& ((++PASS)) || ((++FAIL))
sleep 1
echo ""; echo "=========================================="
echo "Results: $PASS passed, $FAIL failed ($((PASS + FAIL)) total)"
echo "Logs: $LOG_DIR/"
echo "=========================================="
trap - EXIT; cleanup
[ "$FAIL" -gt 0 ] && exit 1
exit 0

122
tests/tcp_proxy_full/stress_client.py

@ -0,0 +1,122 @@
#!/usr/bin/env python3
"""Stress test: parallel threads, sequential echo requests, data verification."""
import os
import socket
import sys
import threading
import random
import time
SO_MARK = 36
def worker(thread_id, host, ports, iters, min_size, max_size, timeout, verbose, results_raw):
ok = 0
times = []
for i in range(iters):
size = random.randint(min_size, max_size)
port = ports[i % len(ports)]
payload = os.urandom(size)
t0 = time.monotonic()
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.setsockopt(socket.SOL_SOCKET, SO_MARK, 1)
s.settimeout(timeout)
try:
t1 = time.monotonic()
s.connect((host, port))
t2 = time.monotonic()
s.sendall(payload)
t3 = time.monotonic()
received = bytearray()
while len(received) < size:
chunk = s.recv(min(65536, size - len(received)))
if not chunk:
break
received += chunk
t4 = time.monotonic()
if len(received) != size or received != payload:
raise RuntimeError(f"data mismatch")
ok += 1
times.append((size, int((t2-t1)*1000), int((t3-t2)*1000), int((t4-t3)*1000)))
except Exception as e:
results_raw[thread_id] = (False, f"FAIL iter={i} size={size}: {e} "
f"connect={int((t2-t1)*1000)}ms send={int((t3-t2)*1000)}ms recv={int((t4-t3)*1000)}ms")
try:
s.close()
except Exception:
pass
return
finally:
if s:
try:
s.close()
except Exception:
pass
agg_connect = sum(t[1] for t in times)
agg_send = sum(t[2] for t in times)
agg_recv = sum(t[3] for t in times)
total_bytes = sum(t[0] for t in times)
results_raw[thread_id] = (True, f"PASS {ok}/{iters} c={agg_connect}ms s={agg_send}ms r={agg_recv}ms total={total_bytes}B")
if verbose:
for i, (sz, cm, sm, rm) in enumerate(times):
print(f" TH{thread_id:02d} iter={i:02d} size={sz:6d} c={cm:4d}ms s={sm:4d}ms r={rm:4d}ms", flush=True)
def main():
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--host", required=True)
parser.add_argument("--port-base", type=int, required=True)
parser.add_argument("--ports", type=int, default=20)
parser.add_argument("--threads", type=int, default=10)
parser.add_argument("--iters", type=int, default=20)
parser.add_argument("--min-size", type=int, default=1024)
parser.add_argument("--max-size", type=int, default=1048576)
parser.add_argument("--timeout", type=float, default=60)
parser.add_argument("--verbose", action="store_true")
args = parser.parse_args()
ports = [args.port_base + i for i in range(args.ports)]
random.seed(os.urandom(8))
threads = []
results_raw = [None] * args.threads
t0 = time.monotonic()
for t in range(args.threads):
th = threading.Thread(
target=worker,
args=(t, args.host, ports, args.iters, args.min_size, args.max_size, args.timeout, args.verbose, results_raw),
daemon=True,
)
th.start()
threads.append(th)
for th in threads:
th.join()
elapsed = time.monotonic() - t0
passed = sum(1 for r in results_raw if r and r[0])
total_iters = args.threads * args.iters
print(f"[{'PASS' if passed == args.threads else 'FAIL'}] stress: "
f"{passed}/{args.threads} threads ({total_iters} sessions) in {elapsed:.1f}s", flush=True)
for i, r in enumerate(results_raw):
if r:
print(f" thread {i}: {r[1]}", flush=True)
if passed < args.threads:
sys.exit(1)
if __name__ == "__main__":
main()

97
tests/tcp_proxy_full/tcp_client.py

@ -0,0 +1,97 @@
#!/usr/bin/env python3
"""TCP test client for tcp_proxy integration test.
Connects to host:port with SO_MARK=1, sends N bytes, receives echo, verifies.
"""
import os
import socket
import sys
import time
SO_MARK = 36
def main():
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--host", required=True)
parser.add_argument("--port", type=int, required=True)
parser.add_argument("--size", type=int, required=True)
parser.add_argument("--verify", action="store_true")
parser.add_argument("--half-close", action="store_true")
parser.add_argument("--sleep-before", type=float, default=0)
parser.add_argument("--timeout", type=float, default=30.0)
parser.add_argument("--mark", type=int, default=1, help="SO_MARK (0=off)")
parser.add_argument("--name", default="test")
args = parser.parse_args()
payload = os.urandom(args.size)
t0 = time.monotonic()
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
if args.mark != 0:
s.setsockopt(socket.SOL_SOCKET, SO_MARK, args.mark)
s.settimeout(args.timeout)
try:
s.connect((args.host, args.port))
except OSError as e:
print(f"[FAIL] {args.name}: connect failed: {e}", flush=True)
s.close()
sys.exit(1)
if args.sleep_before > 0:
time.sleep(args.sleep_before)
try:
s.sendall(payload)
except OSError as e:
print(f"[FAIL] {args.name}: sendall failed: {e}", flush=True)
s.close()
sys.exit(1)
if args.half_close:
s.shutdown(socket.SHUT_WR)
if not args.verify and not args.half_close:
elapsed = time.monotonic() - t0
print(f"[PASS] {args.name}: {args.size} bytes sent in {elapsed:.3f}s "
f"({args.size / 1e6 / elapsed:.2f} MB/s)", flush=True)
s.close()
sys.exit(0)
received = bytearray()
expected = args.size
try:
while len(received) < expected:
chunk = s.recv(min(65536, expected - len(received)))
if not chunk:
break
received += chunk
except OSError as e:
print(f"[FAIL] {args.name}: recv error after {len(received)} bytes: {e}", flush=True)
s.close()
sys.exit(1)
s.close()
elapsed = time.monotonic() - t0
if len(received) != expected:
print(f"[FAIL] {args.name}: {len(received)}/{expected} bytes in {elapsed:.3f}s", flush=True)
sys.exit(1)
if payload != received:
for i in range(min(len(payload), len(received))):
if payload[i] != received[i]:
print(f"[FAIL] {args.name}: data mismatch at byte {i}: "
f"sent={payload[i]:02x} recv={received[i]:02x}", flush=True)
sys.exit(1)
throughput = (expected * 2) / 1e6 / elapsed
print(f"[PASS] {args.name}: {expected} bytes echoed in {elapsed:.3f}s ({throughput:.2f} MB/s)", flush=True)
if __name__ == "__main__":
main()
Loading…
Cancel
Save