Browse Source

fix: proxy ERROR handling — single CLOSE on error, no ERROR response cascade

- tcp_proxy_client_handle_error: set error flag + send CLOSE once if flag was clear
- tcp_proxy_client_handle_data/handle_close: drop silently, don't send ERROR back
- tcp_proxy_server_handle_data/handle_close: drop silently, don't send ERROR back
- Add recv→send 1MB integration test (3 clients × 3 requests, shell + python)
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
e7eb24537d
  1. 8
      src/proxy/tcp_proxy_client.c
  2. 6
      src/proxy/tcp_proxy_server.c
  3. 133
      tests/tcp_proxy_full/recv_send_client.py
  4. 99
      tests/tcp_proxy_full/recv_send_server.py
  5. 162
      tests/tcp_proxy_full/run_recv_send_test.sh

8
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);
}

6
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) {

133
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()

99
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]} <host> <port>", 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()

162
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
Loading…
Cancel
Save