From 265421eb968029d57e555ec02e946502e87570a1 Mon Sep 17 00:00:00 2001 From: evgeny Date: Fri, 25 Sep 2026 23:36:45 +0300 Subject: [PATCH] =?UTF-8?q?tun=20proxy:=20=D1=84=D0=B8=D0=BA=D1=81=20TCP-?= =?UTF-8?q?=D1=81=D1=83=D0=BC=D0=BC=D1=8B=20(=D0=BD=D0=B5=D1=87=D1=91?= =?UTF-8?q?=D1=82=D0=BD=D1=8B=D0=B9=20pbuf),=20FIN=20=D0=B4=D0=BE=20=D0=B4?= =?UTF-8?q?=D1=80=D0=B5=D0=BD=D0=B0=D0=B6=D0=B0=20to=5Flwip,=20tcp=5Fbind/?= =?UTF-8?q?rexmit=5Ffast/active=5Fpcbs=5Fchanged=20+=20=D0=BD=D0=B0=D0=B3?= =?UTF-8?q?=D1=80=D1=83=D0=B7=D0=BE=D1=87=D0=BD=D1=8B=D0=B9=20burst-=D1=82?= =?UTF-8?q?=D0=B5=D1=81=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- check.sh | 6 + src/lwip_tcp/lwip_tcp.c | 13 +- src/lwip_tcp/lwip_tcp.h | 1 + src/lwip_tcp/lwip_tcp_in.c | 3 - src/lwip_tcp/lwip_tcp_out.c | 29 ++- src/lwip_tcp/lwip_tcp_priv.h | 4 +- src/proxy/tcp_proxy_client.c | 19 +- src/proxy/tcp_proxy_client.h | 1 + tests/Makefile.am | 10 + tests/tcp_proxy_full/burst_peer.py | 339 +++++++++++++++++++++++++ tests/tcp_proxy_full/client.conf | 2 +- tests/tcp_proxy_full/run_burst_test.sh | 168 ++++++++++++ tests/tcp_proxy_full/run_load_test.sh | 21 ++ 13 files changed, 601 insertions(+), 15 deletions(-) create mode 100755 tests/tcp_proxy_full/burst_peer.py create mode 100755 tests/tcp_proxy_full/run_burst_test.sh create mode 100755 tests/tcp_proxy_full/run_load_test.sh diff --git a/check.sh b/check.sh index 19127923..e9992b14 100755 --- a/check.sh +++ b/check.sh @@ -20,9 +20,15 @@ rm -f "$TMP" echo "" if [ "$(id -u)" -eq 0 ]; then make -C tests check-proxy + make -C tests check-burst + make -C tests check-load elif command -v sudo >/dev/null 2>&1 && sudo -n true 2>/dev/null; then echo "=== tcp_proxy_full integration test (sudo) ===" sudo -n make -C tests check-proxy + echo "=== tcp_proxy_full burst test (sudo) ===" + sudo -n make -C tests check-burst + echo "=== tcp_proxy_full load test (sudo) ===" + sudo -n make -C tests check-load else echo "[SKIP] tcp_proxy_full integration test (requires passwordless sudo)" fi diff --git a/src/lwip_tcp/lwip_tcp.c b/src/lwip_tcp/lwip_tcp.c index daa4c051..9c64a42c 100644 --- a/src/lwip_tcp/lwip_tcp.c +++ b/src/lwip_tcp/lwip_tcp.c @@ -410,7 +410,8 @@ err_t tcp_bind(struct tcp_pcb *pcb, uint32_t ipaddr, uint16_t port) if (pcb->state != CLOSED) return LERR_VAL; ctx = pcb->ctx; - uint16_t new_port = port; + // Порт принимаем в network-order (конвенция вызывающих), дальше работаем в host-order. + uint16_t new_port = ntohs(port); if (new_port == 0) { new_port = tcp_new_port(pcb); if (new_port == 0) return LERR_BUF; @@ -429,7 +430,7 @@ err_t tcp_bind(struct tcp_pcb *pcb, uint32_t ipaddr, uint16_t port) } if (ipaddr != 0) pcb->local_ip = ipaddr; - pcb->local_port = ntohs(new_port); + pcb->local_port = new_port; TCP_REG(&ctx->bound_pcbs, pcb); return LERR_OK; } @@ -645,6 +646,7 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx) ctx->slowtmr_ctr++; // ---- process active PCBs ---- +tcp_slowtmr_start: prev = NULL; pcb = ctx->active_pcbs; @@ -781,7 +783,9 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx) (uint8_t)pcb2->local_port); tcp_free(pcb2); + ctx->active_pcbs_changed = 0; TCP_EVENT_ERR(last_state, err_fn, err_arg, LERR_ABRT); + if (ctx->active_pcbs_changed) goto tcp_slowtmr_start; } else { prev = pcb; pcb = pcb->next; @@ -790,7 +794,9 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx) if (prev->polltmr >= prev->pollinterval) { prev->polltmr = 0; err = LERR_OK; + ctx->active_pcbs_changed = 0; TCP_EVENT_POLL(prev, err); + if (ctx->active_pcbs_changed) goto tcp_slowtmr_start; if (err == LERR_OK) { tcp_output(prev); } @@ -854,6 +860,7 @@ void tcp_fasttmr(struct lwip_tcp_ctx *ctx) ctx->slowtmr_ctr++; +tcp_fasttmr_start: pcb = ctx->active_pcbs; while (pcb != NULL) { if (pcb->last_timer != ctx->slowtmr_ctr) { @@ -873,7 +880,9 @@ void tcp_fasttmr(struct lwip_tcp_ctx *ctx) next = pcb->next; if (pcb->refused_data != NULL) { + ctx->active_pcbs_changed = 0; tcp_process_refused_data(pcb); + if (ctx->active_pcbs_changed) goto tcp_fasttmr_start; } pcb = next; } else { diff --git a/src/lwip_tcp/lwip_tcp.h b/src/lwip_tcp/lwip_tcp.h index 525ee954..8bb1bf94 100644 --- a/src/lwip_tcp/lwip_tcp.h +++ b/src/lwip_tcp/lwip_tcp.h @@ -212,6 +212,7 @@ struct lwip_tcp_ctx { uint32_t ticks; uint8_t slowtmr_ctr; // for last_timer check (incremented by fasttmr and slowtmr) uint8_t tmr_phase; // for even/odd slowtmr decision (incremented in tcp_tmr_cb) + uint8_t active_pcbs_changed; // коллбэк изменил список active_pcbs — перезапустить обход uint16_t iss_seed; uint16_t ip_id; uint16_t port_seed; diff --git a/src/lwip_tcp/lwip_tcp_in.c b/src/lwip_tcp/lwip_tcp_in.c index f6ae604f..476bd2eb 100644 --- a/src/lwip_tcp/lwip_tcp_in.c +++ b/src/lwip_tcp/lwip_tcp_in.c @@ -725,7 +725,6 @@ static void tcp_receive(struct tcp_pcb *pcb) { int16_t m; uint32_t right_wnd_edge; - int found_dupack = 0; if (flags & TCP_ACK) { right_wnd_edge = pcb->snd_wnd + pcb->snd_wl2; @@ -745,7 +744,6 @@ static void tcp_receive(struct tcp_pcb *pcb) if (pcb->snd_wl2 + pcb->snd_wnd == right_wnd_edge) { if (pcb->rtime >= 0) { if (pcb->lastack == ackno) { - found_dupack = 1; if ((uint8_t)(pcb->dupacks + 1) > pcb->dupacks) ++pcb->dupacks; if (pcb->dupacks > 3) @@ -756,7 +754,6 @@ static void tcp_receive(struct tcp_pcb *pcb) } } } - if (!found_dupack) pcb->dupacks = 0; } else if (TCP_SEQ_BETWEEN(ackno, pcb->lastack + 1, pcb->snd_nxt)) { tcpwnd_size_t acked; diff --git a/src/lwip_tcp/lwip_tcp_out.c b/src/lwip_tcp/lwip_tcp_out.c index c40520de..f745e347 100644 --- a/src/lwip_tcp/lwip_tcp_out.c +++ b/src/lwip_tcp/lwip_tcp_out.c @@ -623,6 +623,9 @@ static uint16_t tcp_pseudo_checksum(uint32_t src_ip, uint32_t dst_ip, uint8_t proto, uint16_t tcp_len, const struct pbuf *p) { uint32_t sum = 0; + uint8_t leftover = 0; + int have_leftover = 0; + const struct pbuf *q; sum += (uint32_t)ntohs((uint16_t)((src_ip >> 16) & 0xFFFF)); sum += (uint32_t)ntohs((uint16_t)(src_ip & 0xFFFF)); @@ -631,15 +634,27 @@ tcp_pseudo_checksum(uint32_t src_ip, uint32_t dst_ip, uint8_t proto, uint16_t tc sum += (uint32_t)(uint16_t)proto; sum += (uint32_t)tcp_len; - const struct pbuf *q; + // Суммируем цепочку pbuf как непрерывный поток байт: нечётный хвост одного + // pbuf должен пароваться с первым байтом следующего, а не дополняться нулём + // по месту (иначе чексумма неверна для конкатенированных сегментов с нечётной + // длиной не-последнего pbuf). for (q = p; q != NULL; q = q->next) { const uint8_t *b = (const uint8_t *)q->payload; uint16_t remaining = q->len; - uint16_t i; - for (i = 0; i + 1 < remaining; i += 2) + uint16_t i = 0; + if (have_leftover && remaining > 0) { + sum += (uint32_t)(((uint16_t)leftover << 8) | b[0]); + i = 1; + have_leftover = 0; + } + for (; i + 1 < remaining; i += 2) sum += (uint32_t)((uint16_t)(b[i] << 8) | b[i + 1]); - if (remaining & 1) sum += (uint32_t)((uint16_t)b[remaining - 1] << 8); + if (i < remaining) { + leftover = b[i]; + have_leftover = 1; + } } + if (have_leftover) sum += (uint32_t)((uint16_t)leftover << 8); while (sum >> 16) sum = (sum & 0xFFFF) + (sum >> 16); return (uint16_t)~sum; @@ -796,9 +811,11 @@ tcp_rexmit_fast(struct tcp_pcb *pcb) { if (pcb == NULL) return; if (pcb->unacked != NULL && !(pcb->flags & TF_INFR)) { + uint32_t rseq = ntohl(pcb->unacked->tcphdr->seqno); + uint16_t rlen = pcb->unacked->len; if (tcp_rexmit(pcb) == LERR_OK) { - lwip_tcp_trace_record(pcb->ctx, 'F', ntohl(pcb->unacked->tcphdr->seqno), - pcb->unacked->len, pcb->ssthresh, (uint16_t)pcb->rto, pcb->rtime, (uint8_t)pcb->state); + lwip_tcp_trace_record(pcb->ctx, 'F', rseq, + rlen, pcb->ssthresh, (uint16_t)pcb->rto, pcb->rtime, (uint8_t)pcb->state); pcb->ssthresh = pcb->cwnd; if ((uint32_t)pcb->snd_wnd < pcb->cwnd) pcb->ssthresh = pcb->snd_wnd; pcb->ssthresh = (tcpwnd_size_t)(pcb->ssthresh / 2); diff --git a/src/lwip_tcp/lwip_tcp_priv.h b/src/lwip_tcp/lwip_tcp_priv.h index f901922b..a254477a 100644 --- a/src/lwip_tcp/lwip_tcp_priv.h +++ b/src/lwip_tcp/lwip_tcp_priv.h @@ -251,8 +251,8 @@ void lwip_tcp_stats_clear(struct lwip_tcp_ctx *ctx); (npcb)->next = NULL; \ } while(0) -#define TCP_REG_ACTIVE(ctx, npcb) TCP_REG(&(ctx)->active_pcbs, npcb) -#define TCP_RMV_ACTIVE(ctx, npcb) TCP_RMV(&(ctx)->active_pcbs, npcb) +#define TCP_REG_ACTIVE(ctx, npcb) do { TCP_REG(&(ctx)->active_pcbs, npcb); (ctx)->active_pcbs_changed = 1; } while(0) +#define TCP_RMV_ACTIVE(ctx, npcb) do { TCP_RMV(&(ctx)->active_pcbs, npcb); (ctx)->active_pcbs_changed = 1; } while(0) #ifdef __cplusplus diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index 0ba71f62..1fdfda3d 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -225,6 +225,15 @@ static void tcp_proxy_client_feed_from_transport(struct tcp_proxy_client_conn *p pc->stream_id, sent_any, q_pre, queue_entry_count(pc->to_lwip), pc->pcb->snd_wnd, pc->pcb->cwnd, unsent); tcp_output(pc->pcb); } + // FIN от exit был отложен до дренажа to_lwip — теперь очередь пуста, шлём FIN локальной стороне + if (pc->fin_deferred && pc->pcb && queue_entry_count(pc->to_lwip) == 0 && + pc->pcb->state != TIME_WAIT && pc->pcb->state != CLOSED) { + pc->fin_deferred = 0; + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN RELAY (deferred) sid=%08x — to_lwip drained, shutdown write", pc->stream_id); + tcp_shutdown(pc->pcb, 0, 1); + tcp_sent(pc->pcb, NULL); + tcp_poll(pc->pcb, NULL, 0); + } } // ==================================================================== @@ -633,8 +642,16 @@ static void tcp_proxy_client_handle_error(struct tcp_proxy_client* p, uint32_t s static void tcp_proxy_client_handle_fin(struct tcp_proxy_client* p, uint32_t stream_id) { struct tcp_proxy_client_conn* pc = tcp_proxy_client_find_conn(p, stream_id); if (!pc || !pc->pcb) return; - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN FROM exit sid=%08x pcb_state=%u — shutdown write (send FIN to local)", stream_id, pc->pcb->state); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN FROM exit sid=%08x pcb_state=%u to_lwip=%d — shutdown write (send FIN to local)", + stream_id, pc->pcb->state, pc->to_lwip ? queue_entry_count(pc->to_lwip) : -1); pc->fin_remote = 1; + // Если в to_lwip ещё есть данные от exit — откладываем FIN до их дренажа, + // иначе остаток потеряется (после tcp_shutdown писать в lwIP нельзя). + if (pc->to_lwip && queue_entry_count(pc->to_lwip) > 0 && pc->pcb->state != TIME_WAIT && pc->pcb->state != CLOSED) { + pc->fin_deferred = 1; + tcp_proxy_client_feed_from_transport(pc); + return; + } if (pc->pcb->state != TIME_WAIT && pc->pcb->state != CLOSED) { tcp_shutdown(pc->pcb, 0, 1); tcp_sent(pc->pcb, NULL); diff --git a/src/proxy/tcp_proxy_client.h b/src/proxy/tcp_proxy_client.h index e30ea755..ca385156 100644 --- a/src/proxy/tcp_proxy_client.h +++ b/src/proxy/tcp_proxy_client.h @@ -38,6 +38,7 @@ struct tcp_proxy_client_conn { uint8_t fin_local; // локальная сторона отправила FIN (lwIP) uint8_t fin_remote; // exit отправил FIN (назначение закрылось) + uint8_t fin_deferred; // FIN от exit пришёл, но отложен до дренажа to_lwip uint8_t rem_closed; // exit отправил CLOSE uint8_t error; // ошибка, немедленная очистка uint8_t close_sent; // отправили CLOSE/ERROR в exit diff --git a/tests/Makefile.am b/tests/Makefile.am index 768aa187..29dbaff9 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -560,3 +560,13 @@ check-local: $(check_PROGRAMS) .PHONY: check-proxy check-proxy: @cd $(srcdir)/tcp_proxy_full && bash run_test.sh + +# Burst-тест tcp_proxy: двусторонний byte-exact трафик со случайными idle (требует sudo). +.PHONY: check-burst +check-burst: + @cd $(srcdir)/tcp_proxy_full && bash run_burst_test.sh + +# Нагрузочный тест tcp_proxy: N параллельных соединений, byte-exact, замер пропускной (требует sudo). +.PHONY: check-load +check-load: + @cd $(srcdir)/tcp_proxy_full && bash run_load_test.sh diff --git a/tests/tcp_proxy_full/burst_peer.py b/tests/tcp_proxy_full/burst_peer.py new file mode 100755 index 00000000..1cf717bf --- /dev/null +++ b/tests/tcp_proxy_full/burst_peer.py @@ -0,0 +1,339 @@ +#!/usr/bin/env python3 +"""Burst/load peer for tcp_proxy integration test (full-duplex, byte-exact). + +Both sides stream a known pseudo-random byte sequence in random chunks +(10B..1MB) with optional random idle (100..500ms) between chunks, over full-duplex +TCP connections, and verify the received stream byte-for-byte. + +Because the sequence is deterministic (self-contained xorshift64* generator), +every received byte is checked against a pre-known value — any corruption, +cross-connection mixing or stall is detected exactly. + +Multiple concurrent connections are supported (--connections N); each connection +carries its index in a 4-byte handshake so both peers derive distinct per-connection +seeds and any cross-connection byte mixing is caught. + +Roles: + server: burst_peer.py --server --host 127.0.0.1 --port 19193 --connections 8 + client: burst_peer.py --host 10.200.100.1 --port 19193 --mark 1 --connections 8 + +Exit code: 0 = PASS (all connections byte-exact), 1 = FAIL. +""" + +import argparse +import math +import random +import select +import signal +import socket +import struct +import sys +import threading +import time + +SO_MARK = 36 +MASK64 = (1 << 64) - 1 + + +class ByteStream: + """Deterministic xorshift64* byte stream (version-independent).""" + + def __init__(self, seed): + self.s = (seed ^ 0x9E3779B97F4A7C15) & MASK64 + if self.s == 0: + self.s = 0x123456789ABCDEF + + def next_u64(self): + x = self.s + x ^= (x >> 12) & MASK64 + x ^= (x << 25) & MASK64 + x ^= (x >> 27) & MASK64 + self.s = x + return (x * 0x2545F4914F6CDD1D) & MASK64 + + def next_bytes(self, n): + out = bytearray() + while len(out) < n: + out += self.next_u64().to_bytes(8, "little") + return bytes(out[:n]) + + +def rand_chunk(rnd, lo, hi): + """Log-uniform chunk size in [lo, hi] — covers small (keep-alive) and large.""" + if hi <= lo: + return lo + v = math.exp(rnd.uniform(math.log(lo), math.log(hi))) + return max(lo, min(hi, int(round(v)))) + + +class StuckError(Exception): + pass + + +def send_all(sock, data): + view = memoryview(data) + while len(view): + n = sock.send(view) + if n <= 0: + raise RuntimeError("send failed (connection closed)") + view = view[n:] + + +def recv_exact(sock, n): + out = bytearray() + while len(out) < n: + c = sock.recv(n - len(out)) + if not c: + raise RuntimeError("connection closed") + out += c + return bytes(out) + + +def receive_session(sock, recv_seed, total, stall_ms, name): + """Read frames until END, verify byte-for-byte. Returns (bytes, frames).""" + rng = ByteStream(recv_seed) + received = 0 + nframes = 0 + last_progress = time.monotonic() + last_log = time.monotonic() + + def read_n(n): + nonlocal last_progress + out = bytearray() + while len(out) < n: + r, _, _ = select.select([sock], [], [], 0.1) + if r: + c = sock.recv(n - len(out)) + if not c: + raise RuntimeError("connection closed at %d/%d bytes" % (received, total)) + out += c + last_progress = time.monotonic() + else: + if time.monotonic() - last_progress >= stall_ms: + raise StuckError("no progress for %dms (received %d/%d)" % + (stall_ms, received, total)) + return bytes(out) + + while True: + hdr = read_n(4) + (length,) = struct.unpack(">I", hdr) + if length == 0: + if received != total: + raise RuntimeError("peer ended early: %d/%d" % (received, total)) + return received, nframes + if received + length > total: + raise RuntimeError("frame overflow: %d+%d > %d" % (received, length, total)) + payload = read_n(length) + expected = rng.next_bytes(length) + if payload != expected: + off = 0 + for i in range(length): + if payload[i] != expected[i]: + off = i + break + raise RuntimeError("byte mismatch at offset %d: got %02x want %02x" % + (received + off, payload[off], expected[off])) + received += length + nframes += 1 + if time.monotonic() - last_log >= 1.0: + print("[%s] recv %d/%d bytes" % (name, received, total), flush=True) + last_log = time.monotonic() + + +def send_session(sock, send_seed, total, min_chunk, max_chunk, min_idle, max_idle): + """Send `total` deterministic bytes in random chunks with random idle, then END.""" + rng = ByteStream(send_seed) + rnd = random.Random(send_seed ^ 0xABCDEF) + remaining = total + nframes = 0 + while remaining > 0: + length = rand_chunk(rnd, min_chunk, max_chunk) + if length > remaining: + length = remaining + if length < 1: + length = 1 + payload = rng.next_bytes(length) + send_all(sock, struct.pack(">I", length) + payload) + remaining -= length + nframes += 1 + if remaining > 0: + time.sleep(rnd.uniform(min_idle, max_idle) / 1000.0) + send_all(sock, struct.pack(">I", 0)) + return nframes, total + + +def run_peer(sock, args, role, name, idx): + """Run one full-duplex session. Returns (ok, recv_bytes, send_bytes, err).""" + if role == "client": + seed_out, seed_in = args.seed_c2s ^ idx, args.seed_s2c ^ idx + else: + seed_out, seed_in = args.seed_s2c ^ idx, args.seed_c2s ^ idx + + sender_res = {} + + def sender(): + try: + sender_res["frames"], sender_res["bytes"] = send_session( + sock, seed_out, args.total, + args.min_chunk, args.max_chunk, + args.min_idle_ms, args.max_idle_ms) + sender_res["ok"] = True + except Exception as e: + sender_res["ok"] = False + sender_res["err"] = str(e) + + t = threading.Thread(target=sender, daemon=True) + t.start() + + try: + recv, _ = receive_session(sock, seed_in, args.total, args.stall_ms, name) + except StuckError as e: + return False, 0, 0, "STUCK: %s" % e + except Exception as e: + return False, 0, 0, str(e) + + t.join() + if not sender_res.get("ok"): + return False, recv, 0, "send error: %s" % sender_res.get("err") + + return True, recv, sender_res["bytes"], None + + +def main(): + p = argparse.ArgumentParser() + p.add_argument("--server", action="store_true") + p.add_argument("--host", required=True) + p.add_argument("--port", type=int, required=True) + p.add_argument("--name", default="burst") + p.add_argument("--connections", type=int, default=1) + p.add_argument("--total", type=int, default=4 * 1024 * 1024) + p.add_argument("--min-chunk", type=int, default=10) + p.add_argument("--max-chunk", type=int, default=1048576) + p.add_argument("--min-idle-ms", type=int, default=100) + p.add_argument("--max-idle-ms", type=int, default=500) + p.add_argument("--stall-ms", type=int, default=2000) + p.add_argument("--mark", type=int, default=1) + p.add_argument("--seed-c2s", type=lambda x: int(x, 0), default=0x11443322778855A) + p.add_argument("--seed-s2c", type=lambda x: int(x, 0), default=0x2233445511223344) + args = p.parse_args() + + if args.server: + run_server(args) + else: + run_client(args) + + +def run_client(args): + results = [] + lock = threading.Lock() + t0 = time.monotonic() + + def worker(i): + name = "%s[%d]" % (args.name, i) + s = None + try: + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + if args.mark: + s.setsockopt(socket.SOL_SOCKET, SO_MARK, args.mark) + s.connect((args.host, args.port)) + s.sendall(struct.pack(">I", i)) + ok, recv, send, err = run_peer(s, args, "client", name, i) + if ok: + print("[PASS] %s: recv %d sent %d" % (name, recv, send), flush=True) + else: + print("[FAIL] %s: %s" % (name, err), flush=True) + with lock: + results.append((ok, recv, send)) + except OSError as e: + print("[FAIL] %s: connect: %s" % (name, e), flush=True) + with lock: + results.append((False, 0, 0)) + finally: + if s: + s.close() + + threads = [threading.Thread(target=worker, args=(i,)) for i in range(args.connections)] + for t in threads: + t.start() + for t in threads: + t.join() + + dt = time.monotonic() - t0 + ok = len(results) == args.connections and all(r[0] for r in results) + total_bytes = sum(r[1] + r[2] for r in results) + if ok: + print("[PASS] %s: %d/%d connections, %d bytes in %.2fs (%.2f MB/s)" % + (args.name, len(results), args.connections, total_bytes, dt, + total_bytes / 1e6 / dt if dt > 0 else 0), flush=True) + sys.exit(0) + failed = args.connections - sum(1 for r in results if r[0]) + print("[FAIL] %s: %d/%d connections failed" % (args.name, failed, args.connections), flush=True) + sys.exit(1) + + +def run_server(args): + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + s.bind((args.host, args.port)) + s.listen(256) + print("BURST: listening on %s:%d connections=%d total=%d" % + (args.host, args.port, args.connections, args.total), flush=True) + + results = [] + lock = threading.Lock() + running = [True] + + def _stop(sig, frame): + running[0] = False + + signal.signal(signal.SIGTERM, _stop) + signal.signal(signal.SIGINT, _stop) + + def handle(conn): + idx = 0 + try: + idxb = recv_exact(conn, 4) + (idx,) = struct.unpack(">I", idxb) + name = "%s[%d]" % (args.name, idx) + ok, recv, send, err = run_peer(conn, args, "server", name, idx) + if ok: + print("[PASS] %s: recv %d sent %d" % (name, recv, send), flush=True) + else: + print("[FAIL] %s: %s" % (name, err), flush=True) + with lock: + results.append(ok) + except OSError as e: + print("[FAIL] %s: %s" % (args.name, e), flush=True) + with lock: + results.append(False) + finally: + conn.close() + + s.settimeout(1.0) + threads = [] + while running[0] and len(threads) < args.connections: + try: + conn, addr = s.accept() + except socket.timeout: + continue + except OSError: + break + t = threading.Thread(target=handle, args=(conn,), daemon=True) + t.start() + threads.append(t) + + for t in threads: + t.join() + s.close() + + ok = len(results) == args.connections and all(results) + if ok: + print("[PASS] %s: %d/%d connections" % (args.name, len(results), args.connections), flush=True) + sys.exit(0) + print("[FAIL] %s: %d/%d connections ok" % + (args.name, sum(1 for r in results if r), args.connections), flush=True) + sys.exit(1) + + +if __name__ == "__main__": + main() diff --git a/tests/tcp_proxy_full/client.conf b/tests/tcp_proxy_full/client.conf index 5f6338b5..ab2a71a8 100644 --- a/tests/tcp_proxy_full/client.conf +++ b/tests/tcp_proxy_full/client.conf @@ -29,4 +29,4 @@ traffic=info tun=info proxy=debug etcp_route=debug -debug=debug +debug=info diff --git a/tests/tcp_proxy_full/run_burst_test.sh b/tests/tcp_proxy_full/run_burst_test.sh new file mode 100755 index 00000000..1362aa7c --- /dev/null +++ b/tests/tcp_proxy_full/run_burst_test.sh @@ -0,0 +1,168 @@ +#!/bin/bash +# run_burst_test.sh — burst-тест tcp_proxy: двусторонний byte-exact трафик со случайными idle +# +# Топология (3 узла, как в run_test.sh): +# burst_peer.py --client (SO_MARK=1) → 10.200.100.1 → tun_test_proxy +# → tcp_proxy_client (lwIP) → ETCP → intermediate (транзит) → ETCP → exit +# → tcp_proxy_server → connect(10.200.100.1) → iptables DNAT(!mark) → 127.0.0.1:KA_PORT +# → burst_peer.py --server +# +# Проверяется: оба направления шлют известную псевдослучайную последовательность +# случайными пачками (10Б..1МБ) со случайными паузами (100..500мс), принимающая сторона +# сверяет каждый байт. Если за stall_ms нет прогресса при неполном приёме — «застряло». +# +# Запуск: sudo ./run_burst_test.sh + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" +UTUN_BIN="$SCRIPT_DIR/../../utun" +LOG_DIR="$SCRIPT_DIR/log" +BURST_PY="$SCRIPT_DIR/burst_peer.py" + +KA_PORT=19193 +BURST_IP="10.200.100.1" + +NAME="${BURST_NAME:-burst}" +CONNECTIONS="${BURST_CONNECTIONS:-1}" +TOTAL="${BURST_TOTAL:-4194304}" # байт на направление на соединение +STALL_MS="${BURST_STALL_MS:-2000}" +MIN_CHUNK="${BURST_MIN_CHUNK:-10}" +MAX_CHUNK="${BURST_MAX_CHUNK:-1048576}" +MIN_IDLE_MS="${BURST_MIN_IDLE_MS:-100}" +MAX_IDLE_MS="${BURST_MAX_IDLE_MS:-500}" +CLIENT_TIMEOUT="${BURST_CLIENT_TIMEOUT:-180}" + +# -------- setup / cleanup -------- + +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 + + iptables -t nat -C OUTPUT -d "$BURST_IP" -p tcp --dport "$KA_PORT" -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$KA_PORT" 2>/dev/null \ + || iptables -t nat -A OUTPUT -d "$BURST_IP" -p tcp --dport "$KA_PORT" -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$KA_PORT" + 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 ===" + iptables -t nat -D OUTPUT -d "$BURST_IP" -p tcp --dport "$KA_PORT" -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$KA_PORT" 2>/dev/null || true + 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 $SERVER_PID 2>/dev/null || true + kill $EXIT_PID 2>/dev/null || true + kill $INTER_PID 2>/dev/null || true + kill $CLIENT_PID 2>/dev/null || true + wait $SERVER_PID 2>/dev/null || true + wait $EXIT_PID 2>/dev/null || true + wait $INTER_PID 2>/dev/null || true + wait $CLIENT_PID 2>/dev/null || true + sleep 0.5; cleanup_net +} + +wait_for_etcp() { + for i in $(seq 1 40); do + grep -q "Connection UP" "$LOG_DIR/exit_utun.log" 2>/dev/null \ + && grep -q "Connection UP" "$LOG_DIR/intermediate_utun.log" 2>/dev/null \ + && grep -q "Connection UP" "$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 burst test (full-duplex, byte-exact) ===" +echo "name=$NAME connections=$CONNECTIONS total=$TOTAL/direction chunk=$MIN_CHUNK..$MAX_CHUNK idle=${MIN_IDLE_MS}..${MAX_IDLE_MS}ms stall=${STALL_MS}ms" + +mkdir -p "$LOG_DIR"; rm -f "$LOG_DIR"/burst_*.log +setup_net + +echo "Starting burst server on 127.0.0.1:$KA_PORT ..." +python3 "$BURST_PY" --server --host 127.0.0.1 --port "$KA_PORT" \ + --connections "$CONNECTIONS" --total "$TOTAL" \ + --min-chunk "$MIN_CHUNK" --max-chunk "$MAX_CHUNK" \ + --min-idle-ms "$MIN_IDLE_MS" --max-idle-ms "$MAX_IDLE_MS" \ + --stall-ms "$STALL_MS" --name "$NAME" \ + >"$LOG_DIR/burst_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 intermediate ..." +"$UTUN_BIN" -c "$SCRIPT_DIR/intermediate.conf" -f -l "$LOG_DIR/intermediate_utun.log" >"$LOG_DIR/intermediate_stdout.log" 2>&1 & +INTER_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 burst (client) ===" +set +e +timeout "$CLIENT_TIMEOUT" python3 "$BURST_PY" --host "$BURST_IP" --port "$KA_PORT" --mark 1 \ + --connections "$CONNECTIONS" --total "$TOTAL" \ + --min-chunk "$MIN_CHUNK" --max-chunk "$MAX_CHUNK" \ + --min-idle-ms "$MIN_IDLE_MS" --max-idle-ms "$MAX_IDLE_MS" \ + --stall-ms "$STALL_MS" --name "$NAME" \ + >"$LOG_DIR/burst_client.log" 2>&1 +CLIENT_RC=$? +wait "$SERVER_PID" 2>/dev/null +SERVER_RC=$? +set -e + +echo ""; echo "=== $NAME results ===" +echo " client: rc=$CLIENT_RC $(tail -1 "$LOG_DIR/burst_client.log" 2>/dev/null)" +echo " server: rc=$SERVER_RC $(tail -1 "$LOG_DIR/burst_server.log" 2>/dev/null)" + +trap - EXIT; cleanup + +if [ "$CLIENT_RC" -eq 0 ] && [ "$SERVER_RC" -eq 0 ]; then + echo "$NAME test: PASS" + exit 0 +else + echo "$NAME test: FAIL (client_rc=$CLIENT_RC server_rc=$SERVER_RC)" + echo " Logs: $LOG_DIR/burst_client.log, $LOG_DIR/burst_server.log, $LOG_DIR/client_utun.log" + exit 1 +fi diff --git a/tests/tcp_proxy_full/run_load_test.sh b/tests/tcp_proxy_full/run_load_test.sh new file mode 100755 index 00000000..60db946c --- /dev/null +++ b/tests/tcp_proxy_full/run_load_test.sh @@ -0,0 +1,21 @@ +#!/bin/bash +# run_load_test.sh — нагрузочный тест tcp_proxy: N параллельных соединений, +# двусторонний byte-exact трафик без idle, замер пропускной способности. +# +# Это пресет поверх run_burst_test.sh (тот же движок, другие параметры). +# +# Параметры (env): +# LOAD_CONNECTIONS — число параллельных соединений (по умолчанию 8) +# LOAD_TOTAL — байт на направление на соединение (по умолчанию 8MB) +# +# Запуск: sudo ./run_load_test.sh + +set -euo pipefail + +BURST_NAME=load \ +BURST_CONNECTIONS="${LOAD_CONNECTIONS:-8}" \ +BURST_TOTAL="${LOAD_TOTAL:-8388608}" \ +BURST_MIN_IDLE_MS=0 \ +BURST_MAX_IDLE_MS=0 \ +BURST_STALL_MS="${LOAD_STALL_MS:-30000}" \ +bash "$(dirname "$0")/run_burst_test.sh"