diff --git a/src/etcp_router.c b/src/etcp_router.c index 9014a812..4baf3b26 100644 --- a/src/etcp_router.c +++ b/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) { diff --git a/src/icmp_proxy.c b/src/icmp_proxy.c index 9cab6b1a..70f5cea0 100644 --- a/src/icmp_proxy.c +++ b/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); } // ==================================================================== diff --git a/src/nat_transport.c b/src/nat_transport.c index 80a39020..270661a3 100644 --- a/src/nat_transport.c +++ b/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 ==================== diff --git a/src/remote_proxy.c b/src/remote_proxy.c index 2ff23fe1..fb3e3c05 100644 --- a/src/remote_proxy.c +++ b/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; } diff --git a/src/tcp_proxy.c b/src/tcp_proxy.c index 44c448c7..4d0bda8b 100644 --- a/src/tcp_proxy.c +++ b/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=%uactive); 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); } // ==================================================================== diff --git a/src/udp_proxy.c b/src/udp_proxy.c index 686d5afc..119580bd 100644 --- a/src/udp_proxy.c +++ b/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); } // ==================================================================== diff --git a/tests/tcp_proxy_full/.gitignore b/tests/tcp_proxy_full/.gitignore new file mode 100644 index 00000000..d0b4e514 --- /dev/null +++ b/tests/tcp_proxy_full/.gitignore @@ -0,0 +1,2 @@ +/log/ +*.log diff --git a/tests/tcp_proxy_full/client.conf b/tests/tcp_proxy_full/client.conf new file mode 100644 index 00000000..f6e4c1d7 --- /dev/null +++ b/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 diff --git a/tests/tcp_proxy_full/echo_server.py b/tests/tcp_proxy_full/echo_server.py new file mode 100755 index 00000000..24f040e9 --- /dev/null +++ b/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]} [--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() diff --git a/tests/tcp_proxy_full/exit.conf b/tests/tcp_proxy_full/exit.conf new file mode 100644 index 00000000..bb7e896e --- /dev/null +++ b/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 diff --git a/tests/tcp_proxy_full/run_test.sh b/tests/tcp_proxy_full/run_test.sh new file mode 100755 index 00000000..0356e2e6 --- /dev/null +++ b/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 diff --git a/tests/tcp_proxy_full/stress_client.py b/tests/tcp_proxy_full/stress_client.py new file mode 100755 index 00000000..3bbcb5fc --- /dev/null +++ b/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() diff --git a/tests/tcp_proxy_full/tcp_client.py b/tests/tcp_proxy_full/tcp_client.py new file mode 100755 index 00000000..fb137644 --- /dev/null +++ b/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()