diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index 4d1a2cb2..e36316f5 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -427,8 +427,7 @@ static struct tcp_proxy_client_conn* tcp_proxy_client_find_conn(struct tcp_proxy static void tcp_proxy_client_handle_data(struct tcp_proxy_client* p, struct ETCP_CONN* conn, uint32_t stream_id, struct ll_entry* entry) { struct tcp_proxy_client_conn* pc = tcp_proxy_client_find_conn(p, stream_id); if (!pc || pc->rem_closed || pc->error || pc->tun_closed) { - DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — no conn/closed, sending ERROR", stream_id); - tcp_proxy_client_send_msg(p->inst, p->via_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0); + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — no conn/closed, dropping", stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; } size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; @@ -448,8 +447,7 @@ static void tcp_proxy_client_handle_close(struct tcp_proxy_client* p, uint32_t s stream_id, pc->tun_closed); pc->rem_closed = 1; } else { - DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE sid=%08x — no conn, sending ERROR", stream_id); - tcp_proxy_client_send_msg(p->inst, p->via_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0); + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE sid=%08x — no conn, dropping", stream_id); } } @@ -458,7 +456,7 @@ static void tcp_proxy_client_handle_error(struct tcp_proxy_client* p, uint32_t s if (pc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY ERROR from exit sid=%08x tun_closed=%d rem_closed=%d", stream_id, pc->tun_closed, pc->rem_closed); - pc->error = 1; + if (!pc->error) { pc->error = 1; tcp_proxy_client_send_close(pc); } } else { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — no conn, silent drop", stream_id); } diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 4247c969..251dcace 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -330,8 +330,7 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id); if (!rc || rc->sock == SOCKET_INVALID || rc->error) { - DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TPS handle_data: no/error conn for sid=%08x, sending ERROR", stream_id); - if (conn) tcp_proxy_server_send_msg(inst, conn->peer_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0); + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TPS handle_data: no/error conn for sid=%08x, dropping", stream_id); queue_dgram_free(entry); queue_entry_free(entry); return -1; } @@ -370,8 +369,7 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_i } prev = &rc->next; } - DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn, sending ERROR", stream_id); - tcp_proxy_server_send_msg(inst, inst->node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0); + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn, dropping", stream_id); } int tcp_proxy_server_init(struct UTUN_INSTANCE* inst) { diff --git a/tests/tcp_proxy_full/recv_send_client.py b/tests/tcp_proxy_full/recv_send_client.py new file mode 100755 index 00000000..410ea6a2 --- /dev/null +++ b/tests/tcp_proxy_full/recv_send_client.py @@ -0,0 +1,133 @@ +#!/usr/bin/env python3 +"""TCP test client for tcp_proxy recv→send integration test. + +Connects to host:port with SO_MARK=1, sends 1MB random data, +receives 4-byte sentinel + 1MB response, verifies. +Supports --count N for multiple sequential requests. +""" + +import os +import socket +import struct +import sys +import time + +SO_MARK = 36 +ONE_MB = 1048576 +SENTINEL = 0xBEEF0102 + + +def recv_exact(sock: socket.socket, n: int) -> bytes: + data = bytearray() + while len(data) < n: + chunk = sock.recv(min(65536, n - len(data))) + if not chunk: + break + data += chunk + return bytes(data) + + +def send_exact(sock: socket.socket, data: bytes) -> None: + sent = 0 + while sent < len(data): + n = sock.send(data[sent:]) + if n <= 0: + raise OSError("send failed") + sent += n + + +def run_one(args, payload: bytes, idx: int) -> int: + 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}[{idx}]: connect: {e}", flush=True) + s.close() + return 1 + + try: + send_exact(s, payload) + except OSError as e: + print(f"[FAIL] {args.name}[{idx}]: send: {e}", flush=True) + s.close() + return 1 + + try: + got = recv_exact(s, 4 + args.size) + except OSError as e: + print(f"[FAIL] {args.name}[{idx}]: recv: {e}", flush=True) + s.close() + return 1 + + s.close() + elapsed = time.monotonic() - t0 + + if len(got) != 4 + args.size: + print(f"[FAIL] {args.name}[{idx}]: len={len(got)}/{4 + args.size} in {elapsed:.3f}s", flush=True) + return 1 + + sentinel, = struct.unpack(">I", got[:4]) + if sentinel != SENTINEL: + print(f"[FAIL] {args.name}[{idx}]: sentinel={sentinel:#010x} expected={SENTINEL:#010x}", flush=True) + return 1 + + received = got[4:] + if not args.verify: + throughput = (args.size * 2) / 1e6 / elapsed + print(f"[PASS] {args.name}[{idx}]: {args.size} bytes in {elapsed:.3f}s ({throughput:.2f} MB/s)", flush=True) + return 0 + + if payload != received: + for i in range(min(len(payload), len(received))): + if payload[i] != received[i]: + print(f"[FAIL] {args.name}[{idx}]: mismatch at {i}: " + f"sent={payload[i]:02x} recv={received[i]:02x}", flush=True) + return 1 + print(f"[FAIL] {args.name}[{idx}]: size mismatch", flush=True) + return 1 + + throughput = (args.size * 2) / 1e6 / elapsed + print(f"[PASS] {args.name}[{idx}]: {args.size} bytes in {elapsed:.3f}s ({throughput:.2f} MB/s)", flush=True) + return 0 + + +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, default=ONE_MB) + parser.add_argument("--verify", action="store_true") + parser.add_argument("--timeout", type=float, default=60.0) + parser.add_argument("--mark", type=int, default=1, help="SO_MARK (0=off)") + parser.add_argument("--count", type=int, default=1, help="number of sequential requests") + parser.add_argument("--name", default="test") + parser.add_argument("--seed", type=int, default=0, help="random seed") + args = parser.parse_args() + + import random + if args.seed: + random.seed(args.seed) + + payloads = [os.urandom(args.size) for _ in range(args.count)] + fails = 0 + + for i in range(args.count): + if run_one(args, payloads[i], i + 1) != 0: + fails += 1 + time.sleep(0.1) + + if fails > 0: + print(f"[FAIL] {args.name}: {fails}/{args.count} failed", flush=True) + sys.exit(1) + print(f"[PASS] {args.name}: all {args.count} ok", flush=True) + + +if __name__ == "__main__": + main() diff --git a/tests/tcp_proxy_full/recv_send_server.py b/tests/tcp_proxy_full/recv_send_server.py new file mode 100755 index 00000000..a1d10d1f --- /dev/null +++ b/tests/tcp_proxy_full/recv_send_server.py @@ -0,0 +1,99 @@ +#!/usr/bin/env python3 +"""TCP server for tcp_proxy integration test: recv 1MB → send 1MB. + +Each connection receives exactly 1MB fully, then sends response: +[sentinel(4) || received_data(1MB)] = 1048580 bytes total. +""" + +import socket +import struct +import sys +import threading +import signal + +running = True +conn_seq = 0 +lock = threading.Lock() + +ONE_MB = 1048576 +SENTINEL = 0xBEEF0102 + + +def recv_exact(conn: socket.socket, n: int) -> bytes: + data = bytearray() + while len(data) < n: + chunk = conn.recv(min(65536, n - len(data))) + if not chunk: + break + data += chunk + return bytes(data) + + +def send_exact(conn: socket.socket, data: bytes) -> None: + sent = 0 + while sent < len(data): + n = conn.send(data[sent:]) + if n <= 0: + raise OSError("send failed") + sent += n + + +def handle(conn: socket.socket): + global conn_seq + with lock: + conn_seq += 1 + seq = conn_seq + try: + received = recv_exact(conn, ONE_MB) + print(f"SERVER conn={seq} recv={len(received)}", flush=True) + if len(received) == ONE_MB: + resp = struct.pack(">I", SENTINEL) + received + send_exact(conn, resp) + print(f"SERVER conn={seq} sent={len(resp)}", flush=True) + else: + print(f"SERVER conn={seq} short_recv={len(received)}/{ONE_MB}", flush=True) + except OSError as e: + print(f"SERVER conn={seq} error: {e}", flush=True) + finally: + conn.close() + + +def main(): + if len(sys.argv) < 3: + print(f"Usage: {sys.argv[0]} ", file=sys.stderr) + sys.exit(1) + + host = sys.argv[1] + port = int(sys.argv[2]) + + global running + + def _stop(sig, frame): + global running + running = False + + signal.signal(signal.SIGTERM, _stop) + signal.signal(signal.SIGINT, _stop) + + 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"SERVER: listening on {host}:{port}", 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,), daemon=True) + t.start() + + s.close() + + +if __name__ == "__main__": + main() diff --git a/tests/tcp_proxy_full/run_recv_send_test.sh b/tests/tcp_proxy_full/run_recv_send_test.sh new file mode 100755 index 00000000..1809c8ef --- /dev/null +++ b/tests/tcp_proxy_full/run_recv_send_test.sh @@ -0,0 +1,162 @@ +#!/bin/bash +# run_recv_send_test.sh — интеграционный тест tcp_proxy: recv→send 1MB +# +# 3 клиента параллельно, каждый по 3 запроса (recv 1MB → send 1MB). +# Всего 9 запросов, макс 3 одновременных TCP-соединения через проксю. +# +# Топология: +# recv_send_client.py (SO_MARK=1) → tun_test_proxy +# → tcp_proxy(lwIP) → ETCP → exit(remote_proxy) +# → connect(10.200.100.N) → iptables DNAT → 127.0.0.1:19092 → recv_send_server.py +# +# Запуск: sudo ./run_recv_send_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/recv_send_client.py" +SERVER_PY="$SCRIPT_DIR/recv_send_server.py" + +SERVER_PORT=19092 +CLIENT_COUNT=3 +REQ_PER_CLIENT=3 +REQ_SIZE=1048576 +N_IP=20 + +# -------- networking -------- + +setup_net() { + echo "=== Setting up networking ===" + sysctl -w net.ipv4.conf.all.rp_filter=0 + + 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 + + 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:"$SERVER_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:"$SERVER_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 + 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:"$SERVER_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 $SERVER_PID 2>/dev/null || true + wait $EXIT_PID 2>/dev/null || true + wait $CLIENT_PID 2>/dev/null || true + wait $SERVER_PID 2>/dev/null || true + sleep 0.5; cleanup_net +} + +wait_for_etcp() { + for i in $(seq 1 30); do + grep -q "Connection established\|initialized and marked as UP (client)" "$LOG_DIR/exit_utun.log" 2>/dev/null && return 0 + grep -q "initialized and marked as UP (client)" "$LOG_DIR/client_utun.log" 2>/dev/null && return 0 + sleep 0.5 + done + return 1 +} + +# -------- 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 recv→send 1MB integration test ===" +echo "Clients: $CLIENT_COUNT parallel × $REQ_PER_CLIENT requests × $(($REQ_SIZE / 1048576))MB = $(($CLIENT_COUNT * $REQ_PER_CLIENT)) total" + +mkdir -p "$LOG_DIR"; rm -f "$LOG_DIR"/*.log +setup_net + +echo "Starting recv_send_server on 127.0.0.1:$SERVER_PORT ..." +python3 "$SERVER_PY" 127.0.0.1 "$SERVER_PORT" >"$LOG_DIR/server.log" 2>&1 & +SERVER_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 test ===" + +T0=$(date +%s%3N) +PASS=0; FAIL=0; pids=() + +declare -A client_results + +for c in $(seq 1 $CLIENT_COUNT); do + ip="10.200.100.${c}" + name="cl${c}" + python3 "$CLIENT_PY" \ + --host "$ip" --port "$SERVER_PORT" \ + --size "$REQ_SIZE" --verify --timeout 120 \ + --count "$REQ_PER_CLIENT" --name "$name" \ + >"$LOG_DIR/client_${c}.log" 2>&1 & + pids+=($!) +done + +for pid in "${pids[@]}"; do + wait "$pid" && ((++PASS)) || ((++FAIL)) +done + +ELAPSED=$(($(date +%s%3N) - T0)) +echo ""; echo "==========================================" +echo "Results: $PASS passed, $FAIL failed ($((PASS + FAIL)) clients)" +echo "Total time: ${ELAPSED}ms (${ELAPSED}ms for ${REQ_PER_CLIENT}×${CLIENT_COUNT}=$(($CLIENT_COUNT * $REQ_PER_CLIENT)) requests × $(($REQ_SIZE / 1048576))MB)" +echo "Logs: $LOG_DIR/" +echo "==========================================" + +for c in $(seq 1 $CLIENT_COUNT); do + log="$LOG_DIR/client_${c}.log" + if [ -f "$log" ]; then + echo " client_${c}: $(tail -1 "$log")" + fi +done + +trap - EXIT; cleanup +[ "$FAIL" -gt 0 ] && exit 1 +exit 0