From 3db891bdcf9f313adf90d76b321c6d1031c4d994 Mon Sep 17 00:00:00 2001 From: W11 Date: Fri, 25 Sep 2026 02:14:27 +0800 Subject: [PATCH] =?UTF-8?q?test(e2e):=2050=E8=A7=82=E7=9C=8B=E7=AB=AF?= =?UTF-8?q?=E5=B9=B6=E5=8F=91=E7=A9=BF=E9=80=8F=E6=AE=8B=E7=95=99=E5=8E=8B?= =?UTF-8?q?=E6=B5=8B=E2=80=94=E2=80=94=E5=B9=B2=E5=87=80FIN/RST=E5=B4=A9?= =?UTF-8?q?=E6=BA=83/=E8=BF=9F=E8=B5=B0=E4=B8=89=E7=A7=8D=E7=A6=BB?= =?UTF-8?q?=E5=9C=BA=E6=B7=B7=E5=90=88,=20=E5=85=A8=E9=83=A8=E7=A6=BB?= =?UTF-8?q?=E5=9C=BA=E5=90=8E=E7=9B=AE=E6=A0=87=E7=AB=AF=E9=9B=B6=E8=BF=9E?= =?UTF-8?q?=E6=8E=A5=E6=AE=8B=E7=95=99=E3=80=81=E9=9A=A7=E9=81=93=E7=AB=8B?= =?UTF-8?q?=E5=8D=B3=E5=8F=AF=E5=A4=8D=E7=94=A8(=E5=B3=B0=E5=80=BC18?= =?UTF-8?q?=E5=AD=98=E6=B4=BB,=E6=94=B6=E5=8F=A30.12s)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- e2e/leak_test.go | 125 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 125 insertions(+) create mode 100644 e2e/leak_test.go diff --git a/e2e/leak_test.go b/e2e/leak_test.go new file mode 100644 index 0000000..2e0652a --- /dev/null +++ b/e2e/leak_test.go @@ -0,0 +1,125 @@ +package e2e + +// 并发残留压测:多观看端并发穿透后全部离场(混合断开方式), 目标端必须 +// 零连接残留, 且隧道立即恢复可用——断开同步的正确性验收。 + +import ( + "fmt" + "io" + "net" + "sync" + "sync/atomic" + "testing" + "time" +) + +// countedEcho 计数回显目标: 跟踪存活连接数与峰值, 连接保持到对端关闭。 +type countedEcho struct { + ln net.Listener + live atomic.Int64 + peak atomic.Int64 +} + +func startCountedEcho(t *testing.T) *countedEcho { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { ln.Close() }) + ce := &countedEcho{ln: ln} + go func() { + for { + c, err := ln.Accept() + if err != nil { + return + } + go func(c net.Conn) { + defer c.Close() + cur := ce.live.Add(1) + defer ce.live.Add(-1) + for { // 峰值 CAS 刷高 + p := ce.peak.Load() + if cur <= p || ce.peak.CompareAndSwap(p, cur) { + break + } + } + io.Copy(c, c) + }(c) + } + }() + return ce +} + +// TestConcurrentViewersNoTargetLeak 50 观看端并发穿透、三种离场方式 +// (干净 FIN / RST 崩溃 / 迟走), 全部离场后目标端零残留, 隧道立即可复用。 +func TestConcurrentViewersNoTargetLeak(t *testing.T) { + n, a := startEnv(t) + ce := startCountedEcho(t) + _, port := mustRegister(t, n, a, ce.ln.Addr().String()) + + const viewers = 50 + var wg sync.WaitGroup + var mu sync.Mutex + var late []net.Conn // 迟走者: 并发段结束后统一关闭 + + for i := 0; i < viewers; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + conn, err := net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", port), 3*time.Second) + if err != nil { + t.Errorf("观看端%d 拨号失败: %v", i, err) + return + } + _ = conn.SetDeadline(time.Now().Add(5 * time.Second)) + payload := []byte(fmt.Sprintf("viewer-%02d", i)) + if _, err := conn.Write(payload); err != nil { + t.Errorf("观看端%d 写失败: %v", i, err) + conn.Close() + return + } + got := make([]byte, len(payload)) + if _, err := io.ReadFull(conn, got); err != nil { + t.Errorf("观看端%d 读回失败: %v", i, err) + conn.Close() + return + } + switch i % 3 { + case 0: // 干净 FIN(页面正常关闭) + conn.Close() + case 1: // RST 硬断(观看端进程崩溃) + if tc, ok := conn.(*net.TCPConn); ok { + tc.SetLinger(0) + } + conn.Close() + default: // 迟走: 压测段结束后统一关 + mu.Lock() + late = append(late, conn) + mu.Unlock() + } + }(i) + } + wg.Wait() + t.Logf("并发 %d 观看端完成往返: 目标端存活=%d 峰值=%d", viewers, ce.live.Load(), ce.peak.Load()) + + for _, c := range late { + c.Close() + } + + deadline := time.Now().Add(10 * time.Second) + for time.Now().Before(deadline) { + if ce.live.Load() == 0 { + break + } + time.Sleep(100 * time.Millisecond) + } + if left := ce.live.Load(); left != 0 { + t.Fatalf("全部观看端离场后目标端残留 %d 条连接(应为 0)", left) + } + + // 隧道压测后立即可复用: 新观看端一轮完整往返 + if got := roundTrip(t, port, []byte("after-storm")); string(got) != "after-storm" { + t.Fatalf("压测后隧道不可用, 回读=%q", got) + } +}