diff --git a/.github/workflows/soak.yml b/.github/workflows/soak.yml index 8c7d350..145a6b1 100644 --- a/.github/workflows/soak.yml +++ b/.github/workflows/soak.yml @@ -14,6 +14,10 @@ on: description: '浸泡时长(秒)' required: false default: '600' + seed: + description: '可选随机 seed(重放动作序列)' + required: false + default: '' permissions: contents: read @@ -54,6 +58,7 @@ jobs: - name: Soak tests (${{ matrix.tfm }}) env: STREAMFRAME_SOAK_SECONDS: ${{ inputs.seconds || '600' }} + STREAMFRAME_SOAK_SEED: ${{ inputs.seed || '' }} run: >- dotnet test test/StreamFrame.Tests/StreamFrame.Tests.csproj -c Release --no-build -f ${{ matrix.tfm }} diff --git a/CHANGELOG.md b/CHANGELOG.md index 0632dee..8274541 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,12 +6,16 @@ ## [Unreleased] ### 修复 +- **Soak 队列续发与收尾([#67](https://github.com/CSJ608/StreamFrame/issues/67))**:发送 worker 在每次出队前检查会话取消,避免上一帧写完后拆除的旧 worker 取走并丢弃待续发普通条目;新增确定性门控回归。长流测试精确排空后关闭读 socket 并观察任务,修复 net48 仅取消令牌时读任务不收敛。 - **TCP KeepAlive 单位([#62](https://github.com/CSJ608/StreamFrame/issues/62))**:现代 .NET 将正毫秒值安全向上取整为秒,默认 30000/5000ms 正确设置为 30/5s;保留 netstandard2.0 IOControl 毫秒语义及默认关闭。增加真实 Socket 配置读回与边界回归,双语文档明确粒度差异。 - **终态连接等待([#65](https://github.com/CSJ608/StreamFrame/issues/65))**:取得连接等待器后复查停机状态,关闭首次注册与 Shutdown 的竞态;Dispose 或生命周期取消后的新等待也以任务取消结束,无需调用方令牌。补充终态调用、1000 次首次注册竞速及调用方取消隔离回归。 - **坏头后完整帧交付([#64](https://github.com/CSJ608/StreamFrame/issues/64))**:定界器返回 false 但已消费无效前缀时继续解析现有缓冲,无进展才等待输入;保留非法长度头四字节丢弃策略,避免完整帧挂起及误触发半帧超时。补充推进契约及连续坏头、粘包、半帧、EOF、自定义定界器和错误元数据回归。 - **监听绑定失败恢复([#63](https://github.com/CSJ608/StreamFrame/issues/63))**:监听 Socket 完成配置、Bind/Listen 后才发布,失败释放并保持字段为空,端口释放后可自动重试接入;发布与停机释放互斥,保留接受循环代次及单客户端监听关闭语义。新增同步信号门控的端口占用恢复、双向消息及取消/Dispose 竞态回归。 - **Soak 目标矩阵([#66](https://github.com/CSJ608/StreamFrame/issues/66))**:显式选择 Ubuntu net8.0/net10.0、Windows net8.0/net10.0/net48,避免 Ubuntu 启动 .NET Framework 测试;job 名、详细日志与始终尝试上传的 TRX 产物按 OS/TFM 区分,测试断言失败继续使 job 失败。普通 CI 必需检查名不变。 +### 变更 +- **Soak 可观测协议([#67](https://github.com/CSJ608/StreamFrame/issues/67))**:保持 SendAsync 仅入队契约,记录可重放 seed、动作、尝试 ID、会话及发送/读取阶段;业务测试使用 ACK、重试、去重并精确校验,重试前独立检查队列续发与会话边界。所有阶段有界等待,未启用长测报告 Skip;完整 600 秒五目标连续三次验收另行跟踪,见 docs/SOAK.md。 + ## [2.6.0] - 2026-08-29 ### 新增 diff --git a/docs/SOAK.md b/docs/SOAK.md new file mode 100644 index 0000000..6ead718 --- /dev/null +++ b/docs/SOAK.md @@ -0,0 +1,27 @@ +# Soak 验收与重放 + +`.github/workflows/soak.yml` 每夜北京时间 03:00 执行三个长测,每项 600 秒,覆盖 Ubuntu net8/net10 和 Windows net8/net10/net48。未设置有效正数 `STREAMFRAME_SOAK_SECONDS` 时报告 Skip。 + +本地先以单目标短时诊断,例如 PowerShell: + +```powershell +$env:STREAMFRAME_SOAK_SECONDS = '10' +$env:STREAMFRAME_SOAK_SEED = '67002' +dotnet test test/StreamFrame.Tests/StreamFrame.Tests.csproj -c Release -f net48 --filter FullyQualifiedName~Soak --logger 'console;verbosity=detailed' +``` + +seed 固定随机动作序列;TCP、调度及按时间截止的迭代数不保证逐字节确定重放。日志同时记录动作序号、随机值、消息/尝试 ID、SessionId、入队任务结果、编码认领、实际本机写出字节、对端采集、ACK 和停止阶段。手动 workflow 输入 seed 可重用失败 seed;未指定则生成并输出 seed。 + +## 三个独立断言面 + +- 无故障长流:发送结束后使用同一 FIFO 队列的哨兵,等待对端 ACK,精确核对全部消息及顺序,不允许尾部少收。 +- 故障下库保证:普通 `SendAsync` 仅保证入队。编码入口观测已认领;原始字节回调观测本机写出,不代表对端确认。终局在重试前排空队列,检查未认领条目送达。确定性 `SendQueueContinuationTests` 独立门控取消发生在首帧写完与下次出队之间,禁止旧 worker 消耗待续发条目。绑定消息仍检查合法任务结局、无重复、接收会话等于绑定会话;每个接收会话内所有尝试保持发送顺序。 +- 业务交付:`p业务ID/尝试号`,对端按业务 ID 去重应用,每次完整接收后发 ACK;只重试未 ACK 的业务 ID,并精确校验所有业务 ID 的应用与 ACK 集合。半帧注入会暂停该连接的 ACK,终局在干净会话重试。`DeliveryProtocolTests` 强制 ACK 丢失,验证入队/应用不会伪造 ACK,重试不会重复应用。 + +未确认尝试按观测分类:未认领、已认领但写出不完整/无法归属、本机整帧已写出但远端未知、对端已采集但 ACK 缺失。`RawBytesSent` 公共事件没有不可变会话编号,日志中的 sid 是当时快照;拆除恰好切断分片时保留归属不确定性,不把快照当作新的库保证。所有历史读任务在故障关闭时被观察,终局不提前取消采集。 + +## 有界停止与验收 + +.NET Framework 的测试端 `NetworkStream.ReadAsync` 在本例中不因令牌取消而结束。排空后关闭所属 socket,再有界观察读任务;连接、发送、排空、ACK、读任务和服务端释放均有命名阶段期限。失败保留原始非零退出和 TRX,不能靠增加 job 上限消除悬挂。 + +短测不算完整验收。Issue #67 要求相同实现的完整 600 秒五目标矩阵连续至少 3 次通过,记录运行链接、提交、seed 和每项实际时长;优先复用夜间运行,串行计数,不重复 dispatch 同一验收。实现变更或完整运行失败重置连续通过计数。验收未齐时 Issue 保持开放,即使实现 PR 已合并。 diff --git a/src/StreamFrame/Connection/StreamConnection.cs b/src/StreamFrame/Connection/StreamConnection.cs index 53aeede..0216d22 100644 --- a/src/StreamFrame/Connection/StreamConnection.cs +++ b/src/StreamFrame/Connection/StreamConnection.cs @@ -782,8 +782,14 @@ private async Task SendWorkerAsync(CancellationToken ct, long sessionId) var reader = _sendQueue.Reader; while (await reader.WaitToReadAsync(ct).ConfigureAwait(false)) { - while (reader.TryRead(out var entry)) + while (true) { + // 上一帧写完后也可能已拆除会话;取消的 worker 不再认领下一条, + // 否则它会出队后在写入处抛取消,丢掉本应由新会话续发的普通条目。 + ct.ThrowIfCancellationRequested(); + if (!reader.TryRead(out var entry)) + break; + if (!entry.IsSessionBound) { await SendFramedAsync(entry.Message, ct).ConfigureAwait(false); diff --git a/test/StreamFrame.Tests/DeliveryProtocolTests.cs b/test/StreamFrame.Tests/DeliveryProtocolTests.cs new file mode 100644 index 0000000..8638c2b --- /dev/null +++ b/test/StreamFrame.Tests/DeliveryProtocolTests.cs @@ -0,0 +1,29 @@ +using StreamFrame; + +namespace StreamFrame.Tests; + +public class DeliveryProtocolTests +{ + private readonly Xunit.Abstractions.ITestOutputHelper _output; + public DeliveryProtocolTests(Xunit.Abstractions.ITestOutputHelper output) => _output = output; + + [Fact] + public async Task LostAck_RetryUsesSameBusinessId_ReceiverAppliesOnce() + { + await using var run = new SoakRun(new SoakTrace(_output)); + var first = await run.ConnectAsync(); + await first.InjectPartialAsync(); // Peer receives, but deliberately withholds ACK on this broken session. + await run.SendBusinessAsync(42); + await run.WaitAsync(() => run.Delivered.ContainsKey(42), "first application"); + await run.ClosePeerAsync(first); + Assert.Empty(run.Acked); // Enqueue and peer application do not fabricate acknowledgement. + await run.WaitAsync(() => run.Server.State != ConnectionState.Connected, "old session ended"); + var second = await run.ConnectAsync(); + await run.SendBusinessAsync(42); + await run.BarrierAsync(second); + Assert.Single(run.Delivered); + Assert.Single(run.Acked); + Assert.Equal(new[] { "p42/0", "p42/1" }, run.Received.Where(x => x.Wire.StartsWith("p", StringComparison.Ordinal)).Select(x => x.Wire)); + await run.ClosePeerAsync(second); // Includes net48 pending ReadAsync shutdown and task observation. + } +} diff --git a/test/StreamFrame.Tests/SendQueueContinuationTests.cs b/test/StreamFrame.Tests/SendQueueContinuationTests.cs new file mode 100644 index 0000000..3c9dc00 --- /dev/null +++ b/test/StreamFrame.Tests/SendQueueContinuationTests.cs @@ -0,0 +1,95 @@ +using System.Buffers.Binary; +using System.Net; +using System.Net.Sockets; +using System.Text; +using StreamFrame; + +namespace StreamFrame.Tests; + +public class SendQueueContinuationTests +{ + [Fact] + public async Task CancelledWorker_AfterSuccessfulWrite_LeavesUnclaimedPlainEntryForNextSession() + { + var listener = new TcpListener(IPAddress.Loopback, 0); + listener.Start(); + var port = ((IPEndPoint)listener.LocalEndpoint).Port; + listener.Stop(); + await using var server = new StreamConnection(new LengthPrefixFramer(), StringCodec.Instance, + IPAddress.Loopback, port, isActive: false); + var written = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var release = new ManualResetEventSlim(); + var first = 0; + var gateTimedOut = 0; + server.RawBytesSent += _ => + { + if (Interlocked.Increment(ref first) != 1) return; + written.TrySetResult(true); + if (!release.Wait(TimeSpan.FromSeconds(10))) Interlocked.Exchange(ref gateTimedOut, 1); + }; + server.Start(CancellationToken.None); + using var peer1 = new TcpClient(); + await SoakRun.BoundedAsync(peer1.ConnectAsync(IPAddress.Loopback, port), "first regression peer connect"); + await UntilAsync(() => server.CurrentSessionId != 0); + var oldSession = server.CurrentSessionId; + var workerField = typeof(StreamConnection).GetField("_sendWorkerTask", + System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic)!; + // Connected is published before StartSession finishes; wait for the worker to be installed. + await UntilAsync(() => workerField.GetValue(server) is Task); + var oldWorker = (Task)workerField.GetValue(server)!; + try + { + await server.SendAsync("already-written"); + Assert.Same(written.Task, await Task.WhenAny(written.Task, Task.Delay(5000))); + await server.SendAsync("must-continue"); + // Callback gates the old worker after the complete write, before its next dequeue. + await SoakRun.BoundedAsync(Task.Run(() => server.Reconnect()), "regression reconnect"); + Assert.NotEqual(oldSession, server.CurrentSessionId); + } + finally { release.Set(); } + await SoakRun.BoundedAsync(oldWorker, "old worker exit", 5000); + Assert.Equal(0, Volatile.Read(ref gateTimedOut)); + peer1.Dispose(); + using var peer2 = new TcpClient(); + await SoakRun.BoundedAsync(peer2.ConnectAsync(IPAddress.Loopback, port), "second regression peer connect"); + await UntilAsync(() => server.CurrentSessionId > oldSession); + await server.SendAsync("sentinel"); + Assert.Equal("must-continue", await ReadAsync(peer2)); + Assert.Equal("sentinel", await ReadAsync(peer2)); + } + + private static async Task UntilAsync(Func condition) + { + var end = TestClock.TickCount64 + 5000; + while (!condition() && TestClock.TickCount64 < end) await Task.Delay(10); + Assert.True(condition()); + } + + private static async Task ReadAsync(TcpClient client) + { + async Task ReadBytes(int count) + { + var bytes = new byte[count]; + var offset = 0; + while (offset < count) + { + var task = client.GetStream().ReadAsync(bytes, offset, count - offset); + if (await Task.WhenAny(task, Task.Delay(5000)) != task) + { + client.Dispose(); + try { await SoakRun.BoundedAsync(task, "failed read cleanup"); } + catch (Exception ex) when (ex is IOException or ObjectDisposedException) { } + Assert.Fail("frame read timed out"); + } + var read = await task; + Assert.True(read > 0); + offset += read; + } + return bytes; + } + var header = await ReadBytes(4); + var length = BinaryPrimitives.ReadInt32BigEndian(header); + Assert.InRange(length, 1, 1024); + return Encoding.UTF8.GetString(await ReadBytes(length)); + } +} diff --git a/test/StreamFrame.Tests/SoakProtocol.cs b/test/StreamFrame.Tests/SoakProtocol.cs new file mode 100644 index 0000000..a743d97 --- /dev/null +++ b/test/StreamFrame.Tests/SoakProtocol.cs @@ -0,0 +1,415 @@ +using System.Buffers; +using System.Buffers.Binary; +using System.Collections.Concurrent; +using System.Net; +using System.Net.Sockets; +using System.Text; +using StreamFrame; + +namespace StreamFrame.Tests; + +#pragma warning disable MA0048 // Cohesive test-only protocol helpers, kept out of the library API. + +internal sealed class SoakFactAttribute : FactAttribute +{ + public SoakFactAttribute() + { + var raw = Environment.GetEnvironmentVariable("STREAMFRAME_SOAK_SECONDS"); + if (!double.TryParse(raw, System.Globalization.NumberStyles.Float, + System.Globalization.CultureInfo.InvariantCulture, out var seconds) || seconds <= 0 || double.IsInfinity(seconds) || double.IsNaN(seconds)) + Skip = "Set STREAMFRAME_SOAK_SECONDS to execute the long-running TCP suite."; + } +} + +internal sealed class SoakTrace +{ + private readonly Xunit.Abstractions.ITestOutputHelper _output; + private readonly long _start = TestClock.TickCount64; + private readonly bool _verbose; + public SoakTrace(Xunit.Abstractions.ITestOutputHelper output, bool verbose = true) + { + _output = output; + _verbose = verbose; + } + public void Log(string text) + { + if (_verbose || text.StartsWith("phase=", StringComparison.Ordinal) || text.StartsWith("state=", StringComparison.Ordinal) || text.StartsWith("long-stream", StringComparison.Ordinal)) + _output.WriteLine($"[{TestClock.TickCount64 - _start}ms] {text}"); + } +} + +internal sealed class SoakAttempt +{ + public string Wire { get; } + public long ClaimedSession; + public int Claimed; + public int RawComplete; + public int Sequence { get; } + public SoakAttempt(string wire, int sequence) { Wire = wire; Sequence = sequence; } +} + +/// Test-only protocol: wire attempt IDs, application ACK, retry and receiver deduplication. +/// Codec entry observes dequeue; RawBytesSent observes local bytes, never a remote acknowledgement. +internal sealed class SoakRun : IAsyncDisposable +{ + private readonly SoakTrace _trace; + private readonly List _peers = new(); + private readonly ConcurrentDictionary _retry = new(); + private readonly ConcurrentDictionary _acks = new(); + private readonly CancellationTokenSource _ackStop = new(); + private readonly Task _ackReader; + private readonly List _raw = new(); + private long _rawSession; + private int _barrier; + private int _sequence; + private readonly List _sends = new(); + private bool _stopped; + private readonly int _port; + public StreamConnection Server { get; } + public ConcurrentDictionary Attempts { get; } = new(); + public ConcurrentQueue<(string Wire, long SessionId)> Received { get; } = new(); + public ConcurrentDictionary Delivered { get; } = new(); + public ConcurrentDictionary Acked { get; } = new(); + + public SoakRun(SoakTrace trace) + { + _trace = trace; + var listener = new TcpListener(IPAddress.Loopback, 0); + listener.Start(); + _port = ((IPEndPoint)listener.LocalEndpoint).Port; + listener.Stop(); + Server = new StreamConnection(new LengthPrefixFramer(), new ObservedCodec(this), + IPAddress.Loopback, _port, isActive: false, + new StreamConnectionOptions { SendQueueCapacity = 256, AcceptRetryDelayMs = 200 }, + logger: new FaultLogger(trace)); + Server.ConnectionChanged += (_, state) => trace.Log($"state={state} sid={Server.CurrentSessionId}"); + Server.RawBytesSent += ObserveRaw; + _ackReader = ReadAcksAsync(); + Server.Start(CancellationToken.None); + } + + private sealed class FaultLogger : Microsoft.Extensions.Logging.ILogger + { + private readonly SoakTrace _trace; + public FaultLogger(SoakTrace trace) => _trace = trace; + public IDisposable? BeginScope(TState state) where TState : notnull => null; + public bool IsEnabled(Microsoft.Extensions.Logging.LogLevel logLevel) => logLevel >= Microsoft.Extensions.Logging.LogLevel.Warning; + public void Log(Microsoft.Extensions.Logging.LogLevel logLevel, Microsoft.Extensions.Logging.EventId eventId, + TState state, Exception? exception, Func formatter) + => _trace.Log($"fault level={logLevel} {formatter(state, exception)} {exception}"); + } + + private sealed class ObservedCodec : ICodec + { + private readonly SoakRun _run; + public ObservedCodec(SoakRun run) => _run = run; + public string Decode(in ReadOnlySequence frame, CancellationToken ct = default) => StringCodec.Instance.Decode(frame, ct); + public void Encode(string message, IBufferWriter writer, CancellationToken ct = default) + { + var attempt = _run.Attempts[message]; + attempt.ClaimedSession = _run.Server.CurrentSessionId; + Interlocked.Exchange(ref attempt.Claimed, 1); + _run._trace.Log($"claim id={message} sid-snapshot={attempt.ClaimedSession} cancelled={ct.IsCancellationRequested}"); + StringCodec.Instance.Encode(message, writer, ct); + } + } + + private void ObserveRaw(ReadOnlyMemory bytes) + { + lock (_raw) + { + var sid = Server.CurrentSessionId; + if (_rawSession != sid) + { + if (_raw.Count > 0) _trace.Log($"raw-partial sid-snapshot={_rawSession} bytes={_raw.Count}"); + _raw.Clear(); + _rawSession = sid; + } + _trace.Log($"raw bytes={bytes.Length} sid-snapshot={sid}"); + _raw.AddRange(bytes.ToArray()); + while (_raw.Count >= 4) + { + var length = BinaryPrimitives.ReadInt32BigEndian(_raw.Take(4).ToArray()); + // The public raw callback has a current-state snapshot, not immutable session metadata. + // A teardown in a partial write may change that snapshot: retain uncertainty in the log. + if (length < 1 || length > 4096) + { + _trace.Log("raw snapshot changed within a frame; attribution unavailable"); + _raw.Clear(); + return; + } + if (_raw.Count < length + 4) return; + var wire = Encoding.UTF8.GetString(_raw.Skip(4).Take(length).ToArray()); + if (Attempts.TryGetValue(wire, out var attempt)) Interlocked.Exchange(ref attempt.RawComplete, 1); + _trace.Log($"local-frame id={wire} sid-snapshot={sid}"); + _raw.RemoveRange(0, length + 4); + } + } + } + + public void Register(string wire) => Assert.True(Attempts.TryAdd(wire, new SoakAttempt(wire, _sequence++)), $"duplicate attempt {wire}"); + + public void SendBound(long sid, string wire) + { + Register(wire); + _sends.Add(ObserveAsync()); + async Task ObserveAsync() + { + try + { + await Server.SendInSessionAsync(sid, wire); + _trace.Log($"bound-result id={wire} full-local-write"); + } + catch (Exception ex) when (ex is SessionExpiredException or SocketException or OperationCanceledException) + { + _trace.Log($"bound-result id={wire} {ex.GetType().Name}"); + } + } + } + + public Task ObserveSendsAsync() => BoundedAsync(Task.WhenAll(_sends), "all bound sends"); + + public async Task EnqueueAsync(string wire) + { + Register(wire); + using var timeout = new CancellationTokenSource(15000); + var send = Server.SendAsync(wire, timeout.Token); + try { await BoundedAsync(send, $"enqueue {wire}"); } + finally { _trace.Log($"enqueue-result id={wire} status={send.Status} sid={Server.CurrentSessionId}"); } + } + + public Task SendBusinessAsync(int id) + { + var attempt = _retry.AddOrUpdate(id, 0, (_, previous) => previous + 1); + return EnqueueAsync($"p{id}/{attempt}"); + } + + public async Task ConnectAsync() + { + _trace.Log("phase=connect begin"); + var end = TestClock.TickCount64 + 15000; + while (true) + { + var client = new TcpClient(); + try + { + await BoundedAsync(client.ConnectAsync(IPAddress.Loopback, _port), "TCP connect", 5000); + await WaitAsync(() => Server.CurrentSessionId != 0, "session publication"); + var peer = new SoakPeer(this, client, Server.CurrentSessionId, _trace); + _peers.Add(peer); + _trace.Log($"phase=connect done sid={peer.SessionId}"); + return peer; + } + catch (SocketException) when (TestClock.TickCount64 < end) + { + client.Dispose(); + await Task.Delay(100); + } + catch { client.Dispose(); throw; } + } + } + + public async Task BarrierAsync(SoakPeer peer) + { + var wire = $"z{_barrier++}"; + _trace.Log($"phase=drain begin id={wire} sid={peer.SessionId}"); + await EnqueueAsync(wire); + await WaitAsync(() => _acks.ContainsKey(wire), $"drain ACK {wire}"); + _trace.Log($"phase=drain done id={wire}"); + } + + private async Task ReadAcksAsync() + { + try + { + await foreach (var message in Server.GetSessionMessages(_ackStop.Token)) + { + var ack = message.Message; + Assert.StartsWith("ack/", ack, StringComparison.Ordinal); + var wire = ack.Substring(4); + Assert.True(Attempts.ContainsKey(wire), $"unknown ACK {wire}"); + _acks.TryAdd(wire, 0); + if (wire.StartsWith("p", StringComparison.Ordinal)) Acked.TryAdd(BusinessId(wire), 0); + _trace.Log($"ACK id={wire} sid={message.SessionId}"); + } + } + catch (OperationCanceledException) when (_ackStop.IsCancellationRequested) { } + } + + public static int BusinessId(string wire) => int.Parse(wire.Substring(1).Split('/')[0], System.Globalization.CultureInfo.InvariantCulture); + + public void OnFrame(string wire, long sid) + { + Assert.True(Attempts.ContainsKey(wire), $"unknown peer frame {wire}"); + Received.Enqueue((wire, sid)); + if (wire.StartsWith("p", StringComparison.Ordinal)) + { + var first = Delivered.TryAdd(BusinessId(wire), 0); + _trace.Log($"business id={wire} sid={sid} applied={first}"); + } + _trace.Log($"peer-frame id={wire} sid={sid}"); + } + + public string[] UnclaimedPlain() => Attempts.Values + .Where(x => x.Wire.StartsWith("p", StringComparison.Ordinal) && Volatile.Read(ref x.Claimed) == 0) + .Select(x => x.Wire).ToArray(); + + public void AssertQueueContinuation(string[] unclaimed) + { + var seen = new HashSet(Received.Select(x => x.Wire)); + foreach (var wire in unclaimed) Assert.Contains(wire, seen); + foreach (var attempt in Attempts.Values.Where(x => x.Wire.StartsWith("p", StringComparison.Ordinal))) + Assert.Equal(1, Volatile.Read(ref attempt.Claimed)); + } + + public void ReportUnconfirmed() + { + var seen = new HashSet(Received.Select(x => x.Wire)); + foreach (var attempt in Attempts.Values.Where(x => x.Wire.StartsWith("p", StringComparison.Ordinal) && !_acks.ContainsKey(x.Wire))) + { + var stage = seen.Contains(attempt.Wire) ? "peer-collected/ACK-missing" : + Volatile.Read(ref attempt.RawComplete) != 0 ? "local-write-complete/remote-unconfirmed" : + Volatile.Read(ref attempt.Claimed) != 0 ? "claimed/write-incomplete-or-raw-attribution-unavailable" : "unclaimed"; + _trace.Log($"unconfirmed id={attempt.Wire} stage={stage} claim-sid-snapshot={attempt.ClaimedSession}"); + } + } + + public async Task WaitAsync(Func condition, string phase) + { + var end = TestClock.TickCount64 + 15000; + while (!condition() && TestClock.TickCount64 < end) + { + if (_ackReader.IsFaulted) await _ackReader; + foreach (var peer in _peers) if (peer.Reader.IsFaulted) await peer.Reader; + await Task.Delay(10); + } + Assert.True(condition(), $"phase={phase} timed out, state={Server.State} sid={Server.CurrentSessionId}"); + } + + public async Task ClosePeerAsync(SoakPeer peer) + { + _trace.Log($"phase=reader-stop sid={peer.SessionId}"); + peer.Close(); // NetworkStream.ReadAsync cancellation alone does not abort pending net48 I/O. + await BoundedAsync(peer.Reader, $"reader stop sid={peer.SessionId}"); + _trace.Log($"phase=reader-observed sid={peer.SessionId} status={peer.Reader.Status}"); + } + + public async Task StopAsync() + { + if (_stopped) return; + _stopped = true; + foreach (var peer in _peers) peer.Close(); + try { await BoundedAsync(Task.WhenAll(_peers.Select(x => x.Reader)), "all peer readers"); } + finally + { + _ackStop.Cancel(); + try { await BoundedAsync(_ackReader, "ACK reader stop"); } + finally + { + await BoundedAsync(Task.Run(async () => await Server.DisposeAsync()), "server dispose"); + await ObserveSendsAsync(); + _ackStop.Dispose(); + _trace.Log("phase=server-disposed"); + } + } + } + + public ValueTask DisposeAsync() => new(StopAsync()); + + public static async Task BoundedAsync(Task task, string phase, int milliseconds = 15000) + { + using var timer = new CancellationTokenSource(); + if (await Task.WhenAny(task, Task.Delay(milliseconds, timer.Token)) != task) + { + // Observe a late fault even when the bounded failure has already unwound the owner. + _ = task.ContinueWith(t => { _ = t.Exception; }, CancellationToken.None, + TaskContinuationOptions.OnlyOnFaulted | TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default); + throw new TimeoutException($"phase={phase} exceeded {milliseconds}ms, status={task.Status}"); + } + timer.Cancel(); + await task; + } +} + +internal sealed class SoakPeer +{ + private readonly SoakRun _run; + private readonly SoakTrace _trace; + private readonly NetworkStream _stream; + private readonly SemaphoreSlim _write = new(1, 1); + private int _closing; + private bool _poisoned; + public TcpClient Client { get; } + public long SessionId { get; } + public Task Reader { get; } + public SoakPeer(SoakRun run, TcpClient client, long sessionId, SoakTrace trace) + { + _run = run; + _trace = trace; + Client = client; + SessionId = sessionId; + _stream = client.GetStream(); + Reader = ReadAsync(); + } + public void Close() + { + Interlocked.Exchange(ref _closing, 1); + Client.Dispose(); + } + public async Task InjectPartialAsync() + { + await SoakRun.BoundedAsync(_write.WaitAsync(), "partial write lock"); + try + { + _poisoned = true; // No ACK writes behind a deliberately incomplete incoming frame. + var junk = new byte[] { 0, 0, 15, 160 }; + await SoakRun.BoundedAsync(_stream.WriteAsync(junk, 0, junk.Length), "partial write"); + _trace.Log($"partial-frame sid={SessionId}"); + } + finally { _write.Release(); } + } + private async Task ReadAsync() + { + try + { + while (true) + { + var header = await ReadBytesAsync(4); + if (header is null) return; + var length = BinaryPrimitives.ReadInt32BigEndian(header); + Assert.InRange(length, 1, 4096); + var payload = await ReadBytesAsync(length); + if (payload is null) return; // Fault injection may truncate an in-flight frame. + var wire = Encoding.UTF8.GetString(payload); + _run.OnFrame(wire, SessionId); + await SoakRun.BoundedAsync(_write.WaitAsync(), "ACK write lock"); + try + { + if (_poisoned) continue; + var ack = Encoding.UTF8.GetBytes("ack/" + wire); + var frame = new byte[ack.Length + 4]; + BinaryPrimitives.WriteInt32BigEndian(frame.AsSpan(0, 4), ack.Length); + ack.CopyTo(frame, 4); + await SoakRun.BoundedAsync(_stream.WriteAsync(frame, 0, frame.Length), $"ACK write {wire}"); + } + finally { _write.Release(); } + } + } + catch (Exception ex) when (ex is IOException or SocketException || ex is ObjectDisposedException && Volatile.Read(ref _closing) != 0) + { + _trace.Log($"reader-end sid={SessionId} exception={ex.GetType().Name} closing={_closing}"); + } + } + private async Task ReadBytesAsync(int length) + { + var bytes = new byte[length]; + var offset = 0; + while (offset < length) + { + var count = await _stream.ReadAsync(bytes, offset, length - offset); + if (count == 0) return null; + offset += count; + } + return bytes; + } +} diff --git a/test/StreamFrame.Tests/SoakTests.cs b/test/StreamFrame.Tests/SoakTests.cs index db1ff33..f369511 100644 --- a/test/StreamFrame.Tests/SoakTests.cs +++ b/test/StreamFrame.Tests/SoakTests.cs @@ -1,8 +1,4 @@ -using System.Buffers.Binary; -using System.Collections.Concurrent; -using System.Net; using System.Net.Sockets; -using System.Text; using StreamFrame; namespace StreamFrame.Tests; @@ -13,7 +9,7 @@ namespace StreamFrame.Tests; /// /// 启用方式(本地手动触发,常规 CI 与默认 dotnet test 均不运行): /// STREAMFRAME_SOAK_SECONDS=120 dotnet test -f net8.0 --filter FullyQualifiedName~Soak -/// 未设置环境变量时本套件立即返回(视为通过),不影响提交前的全量验证。 +/// 未设置环境变量时报告 Skip,不把门控跳过计为长测通过。 /// 长期观察:.github/workflows/soak.yml 每夜以 600s 跑本套件(ubuntu + windows), /// 覆盖用户重连竞速场景(Soak_ReconnectRacing_LongRun,含状态机转移合法性与 /// 会话编号不变式的全程校验)。 @@ -24,20 +20,6 @@ public class SoakTests public SoakTests(Xunit.Abstractions.ITestOutputHelper output) => _output = output; - /// 把库内 Warning(accept 失败等)透传到测试输出——长期观察失败时的第一手诊断。 - private sealed class SoakDumpLogger : Microsoft.Extensions.Logging.ILogger - { - private readonly Xunit.Abstractions.ITestOutputHelper _sink; - public SoakDumpLogger(Xunit.Abstractions.ITestOutputHelper sink) => _sink = sink; - public IDisposable? BeginScope(TState state) where TState : notnull => null; - public bool IsEnabled(Microsoft.Extensions.Logging.LogLevel logLevel) => true; - public void Log(Microsoft.Extensions.Logging.LogLevel logLevel, Microsoft.Extensions.Logging.EventId eventId, TState state, Exception? exception, Func formatter) - { - var line = $"[LOG {logLevel}] {formatter(state, exception)}{(exception is null ? string.Empty : " (" + exception.GetType().Name + ")")}"; - _sink.WriteLine(line); - } - } - /// /// 状态机转移的全程记录与校验器(长期观察的核心断言面): /// 合法边检查、Retry 后必须收敛回 Connected、非 Connected 态会话编号必须为 0、 @@ -48,6 +30,7 @@ private sealed class StateTransitionRecorder { private readonly List<(ConnectionState State, long SessionId)> _transitions = new(); // 兼作监视锁(net48 需 object 锁) private long _lastConnectedId; + private readonly List _errors = new(); public int DirectConnectedMigrations; @@ -73,7 +56,7 @@ public void Record(ConnectionState state, long sessionId) _ => false, }; if (!legal) - throw new InvalidOperationException($"非法状态转移:{previous} → {state}"); + _errors.Add($"非法状态转移:{previous} → {state}"); if (previous == ConnectionState.Connected && state == ConnectionState.Connected) DirectConnectedMigrations++; @@ -88,6 +71,7 @@ public void Validate() { lock (_transitions) { + Assert.Empty(_errors); // 事件处理器异常被库隔离,必须在测试任务内断言。 foreach (var (state, sessionId) in _transitions) { if (state == ConnectionState.Connected) @@ -108,572 +92,134 @@ public void Validate() } } - private static int GetFreePort() - { - var listener = new TcpListener(IPAddress.Loopback, 0); - listener.Start(); - var port = ((IPEndPoint)listener.LocalEndpoint).Port; - listener.Stop(); - return port; - } - private static double SoakSeconds() { var raw = Environment.GetEnvironmentVariable("STREAMFRAME_SOAK_SECONDS"); return double.TryParse(raw, System.Globalization.NumberStyles.Float, System.Globalization.CultureInfo.InvariantCulture, out var value) && value > 0 ? value : 0; } - private static async Task ConnectWithRetryAsync(TcpClient client, int port) - { - // 混沌场景容忍慢收敛:服务端重连可能处于 2s 级的 accept 重试延迟(TIME_WAIT 竞争等), - // 连接预算给足 60s,避免把"暂时拒绝"误判为测试失败 - for (var attempt = 0; ; attempt++) - { - try - { - await client.ConnectAsync(IPAddress.Loopback, port); - return; - } - catch (SocketException) when (attempt < 120) - { - await Task.Delay(500); - } - } - } - - private static async Task WaitForStateAsync( - StreamConnection connection, - Func predicate, - int timeoutMs = 15_000) - { - var deadline = TestClock.TickCount64 + timeoutMs; - while (TestClock.TickCount64 < deadline) - { - if (predicate(connection.State)) - return; - await Task.Delay(20); - } - - Assert.True(predicate(connection.State), $"等待连接状态超时({timeoutMs}ms),当前 {connection.State}。"); - } - - [Fact] + [SoakFact] public async Task Soak_LongRun_MessageIntegrity() { var seconds = SoakSeconds(); - if (seconds <= 0) - return; // 未启用浸泡模式:跳过(保持默认套件快速) - - var port = GetFreePort(); - await using var server = new StreamConnection( - new LengthPrefixFramer(), StringCodec.Instance, IPAddress.Loopback, port, isActive: false); - server.Start(CancellationToken.None); - - using var client = new TcpClient(); - await ConnectWithRetryAsync(client, port); - await WaitForStateAsync(server, s => s == ConnectionState.Connected); - - // 服务端 → 客户端方向:序列消息流;对端读尽并核对顺序与完整性 - var received = new ConcurrentQueue(); - var readCts = new CancellationTokenSource(TimeSpan.FromSeconds(seconds + 30)); - var readerTask = Task.Run(async () => - { - var stream = client.GetStream(); - var header = new byte[4]; - var body = new byte[4096]; - try - { - while (true) - { - var read = 0; - while (read < 4) - { - var n = await stream.ReadAsync(header, read, 4 - read, readCts.Token); - if (n == 0) - return; - read += n; - } - - var length = BinaryPrimitives.ReadInt32BigEndian(header); - var payload = new byte[length]; - read = 0; - while (read < length) - { - var n = await stream.ReadAsync(payload, read, length - read, readCts.Token); - if (n == 0) - return; - read += n; - } - - received.Enqueue(Encoding.UTF8.GetString(payload)); - } - } - catch (OperationCanceledException) - { - } - }); - - var deadline = TestClock.TickCount64 + (int)(seconds * 1000); + var trace = new SoakTrace(_output, verbose: false); + await using var run = new SoakRun(trace); + var peer = await run.ConnectAsync(); + var end = TestClock.TickCount64 + (long)(seconds * 1000); var sent = 0; - while (TestClock.TickCount64 < deadline) + while (TestClock.TickCount64 < end) { - await server.SendAsync($"msg-{sent}", CancellationToken.None); - sent++; - await Task.Delay(5); // ~200 msg/s + await run.EnqueueAsync($"m{sent++}"); + await Task.Delay(5); } + await run.BarrierAsync(peer); + await run.ClosePeerAsync(peer); + var frames = run.Received.Select(x => x.Wire).Where(x => x.StartsWith("m", StringComparison.Ordinal)).ToArray(); + Assert.Equal(Enumerable.Range(0, sent).Select(i => $"m{i}"), frames); + trace.Log($"long-stream exact delivery={sent}, reader observed after socket close"); + } - readCts.Cancel(); - await readerTask; + [SoakFact] + public Task Soak_Chaos_RandomFaults_NoHangAndIntegrity() => RunFaultsAsync(racing: false); - // 完整性:顺序、无丢失、无重复 - Assert.True(received.Count >= sent - 5, $"发送 {sent},对端仅收 {received.Count}(允许尾部少数在途)"); - var index = 0; - foreach (var frame in received) - { - Assert.Equal($"msg-{index}", frame); - index++; - } - } + [SoakFact] + public Task Soak_ReconnectRacing_LongRun() => RunFaultsAsync(racing: true); - [Fact] - public async Task Soak_Chaos_RandomFaults_NoHangAndIntegrity() + private async Task RunFaultsAsync(bool racing) { - var seconds = SoakSeconds(); - if (seconds <= 0) - return; // 未启用浸泡模式:跳过 - - var seed = Environment.TickCount; + var rawSeed = Environment.GetEnvironmentVariable("STREAMFRAME_SOAK_SEED"); + var seed = string.IsNullOrEmpty(rawSeed) ? Environment.TickCount : int.Parse(rawSeed, System.Globalization.CultureInfo.InvariantCulture); var random = new Random(seed); - var port = GetFreePort(); - await using var server = new StreamConnection( - new LengthPrefixFramer(), StringCodec.Instance, IPAddress.Loopback, port, - isActive: false, new StreamConnectionOptions { SendQueueCapacity = 256 }); - server.Start(CancellationToken.None); - - var peerFrames = new ConcurrentQueue(); - var pendingPlain = new ConcurrentQueue(); - var allTasks = new List(); - TcpClient? client = null; - NetworkStream? stream = null; - var readCts = new CancellationTokenSource(TimeSpan.FromSeconds(seconds + 30)); - var readerTask = Task.CompletedTask; - - var deadline = TestClock.TickCount64 + (int)(seconds * 1000); - var boundSeq = 0; - var plainSeq = 0; - - try - { - while (TestClock.TickCount64 < deadline) - { - var action = random.Next(100); - if (client is null) - { - // (重)连接对端并开始读 - client = new TcpClient(); - await ConnectWithRetryAsync(client, port); - await WaitForStateAsync(server, s => s == ConnectionState.Connected); - stream = client.GetStream(); - var currentStream = stream; - readerTask = Task.Run(async () => - { - var header = new byte[4]; - try - { - while (true) - { - var read = 0; - while (read < 4) - { - var n = await currentStream.ReadAsync(header, read, 4 - read, readCts.Token); - if (n == 0) - return; - read += n; - } - - var length = BinaryPrimitives.ReadInt32BigEndian(header); - var payload = new byte[length]; - read = 0; - while (read < length) - { - var n = await currentStream.ReadAsync(payload, read, length - read, readCts.Token); - if (n == 0) - return; - read += n; - } - - peerFrames.Enqueue(Encoding.UTF8.GetString(payload)); - } - } - catch (Exception ex) when (ex is ObjectDisposedException or IOException or OperationCanceledException) - { - } - }); - } - else if (action < 50) - { - // 会话绑定突发 - var id = server.CurrentSessionId; - if (id != 0) - { - for (var i = 0; i < random.Next(1, 5); i++) - { - var label = $"b{boundSeq++}"; - allTasks.Add(server.SendInSessionAsync(id, label)); - await Task.Delay(10); - } - } - } - else if (action < 70) - { - // 普通发送(跨会话续发):记录标签,终局核对全部送达 - for (var i = 0; i < random.Next(1, 4); i++) - { - var label = $"p{plainSeq++}"; - pendingPlain.Enqueue(label); - await server.SendAsync(label); - } - } - else if (action < 80 && stream is not null) - { - // 半帧注入(未完成帧超时未启用:仅占住解码缓冲) - var junk = new byte[] { 0x00, 0x00, 0x0F, 0xA0 }; - await stream.WriteAsync(junk); - } - else if (action < 95) - { - // 杀对端(FIN 路径:正常释放) - if (client is not null) - { - client.Dispose(); - client = null; - stream = null; - await WaitForStateAsync(server, s => s != ConnectionState.Connected); - } - } - else if (action < 97) - { - // 杀对端(RST 路径:Linger(true,0) 硬复位) - if (client is not null) - { - client.LingerState = new LingerOption(true, 0); - client.Dispose(); - client = null; - stream = null; - await WaitForStateAsync(server, s => s != ConnectionState.Connected); - } - } - else - { - // 用户显式重连(#47 观察到的楔死触发器,与本会话并发故障竞速): - // 不等当前会话的自动重连,立即发起服务端主动拆除 - server.Reconnect(); - if (client is not null) - { - await WaitForStateAsync(server, s => s != ConnectionState.Connected); - client.Dispose(); - client = null; - stream = null; - } - } - - await Task.Delay(random.Next(20, 120)); - } - } - finally - { - readCts.Cancel(); - - // 循环结束时对端可能仍存活:单客户端模式下监听已关闭(AcceptFirstClientOnly), - // 必须先释放存量对端并等服务端离开 Connected,终局对端才能接入 - client?.Dispose(); - if (client is not null) - { - // 非断言式等待(尽力推进;真正的状态断言留给终局阶段) - var exitDeadline = TestClock.TickCount64 + 15_000; - while (TestClock.TickCount64 < exitDeadline && server.State == ConnectionState.Connected) - await Task.Delay(20); - } - } - - // 1) 所有会话绑定发送必须在时限内终结,失败类型合法;成功帧不得重复投递 - var boundSuccesses = new HashSet(); - foreach (var task in allTasks) + var trace = new SoakTrace(_output); + trace.Log($"seed={seed} racing={racing} seconds={SoakSeconds()}"); + var recorder = new StateTransitionRecorder(); + await using var run = new SoakRun(trace); + run.Server.ConnectionChanged += (_, state) => recorder.Record(state, run.Server.CurrentSessionId); + SoakPeer? peer = null; + var plain = 0; + var bound = 0; + var actionId = 0; + var end = TestClock.TickCount64 + (long)(SoakSeconds() * 1000); + while (TestClock.TickCount64 < end) { - var done = await Task.WhenAny(task, Task.Delay(20_000)); - Assert.True(ReferenceEquals(task, done), "混沌后存在悬挂的会话绑定发送。"); - if (task.Status != TaskStatus.RanToCompletion) + var action = random.Next(100); + trace.Log($"action={actionId++} draw={action} sid={run.Server.CurrentSessionId} peer={peer?.SessionId}"); + if (peer is null) + peer = await run.ConnectAsync(); + else if (racing ? action < 70 : action >= 80) { - var ex = task.Exception?.InnerExceptions[0]; - Assert.True(ex is SessionExpiredException or SocketException or OperationCanceledException, - $"混沌中失败类型意外:{ex?.GetType().Name}"); + if (racing && action < 50) + { + var victim = peer; + var kill = Task.Run(() => victim.Close()); + await SoakRun.BoundedAsync(Task.Run(() => run.Server.Reconnect()), "racing reconnect"); + await SoakRun.BoundedAsync(kill, "peer kill"); + } + else if (racing || action >= 97) + await SoakRun.BoundedAsync(Task.Run(() => run.Server.Reconnect()), "explicit reconnect"); + else if (action >= 95) + peer.Client.LingerState = new LingerOption(true, 0); + await run.ClosePeerAsync(peer); + await run.WaitAsync(() => run.Server.State != ConnectionState.Connected, "leave Connected"); + peer = null; } - } - - // 2) 重连一个干净对端,等普通消息全部续发送达(FIFO 完整性) - using (var finalClient = new TcpClient()) - { - await ConnectWithRetryAsync(finalClient, port); - await WaitForStateAsync(server, s => s == ConnectionState.Connected); - var finalStream = finalClient.GetStream(); - var finalCts = new CancellationTokenSource(TimeSpan.FromSeconds(20)); - var finalReader = Task.Run(async () => + else if (!racing && action < 50) { - var header = new byte[4]; - try - { - while (true) - { - var read = 0; - while (read < 4) - { - var n = await finalStream.ReadAsync(header, read, 4 - read, finalCts.Token); - if (n == 0) - return; - read += n; - } - - var length = BinaryPrimitives.ReadInt32BigEndian(header); - var payload = new byte[length]; - read = 0; - while (read < length) - { - var n = await finalStream.ReadAsync(payload, read, length - read, finalCts.Token); - if (n == 0) - return; - read += n; - } - - peerFrames.Enqueue(Encoding.UTF8.GetString(payload)); - } - } - catch (Exception ex) when (ex is ObjectDisposedException or IOException or OperationCanceledException) + var sid = run.Server.CurrentSessionId; + var count = random.Next(1, 5); + for (var i = 0; i < count; i++) { + var wire = $"b{bound++}/{sid}"; + run.SendBound(sid, wire); + trace.Log($"bound-submit id={wire} sid={sid}"); + await Task.Delay(10); } - }); - - var finalDeadline = TestClock.TickCount64 + 20_000; - while (TestClock.TickCount64 < finalDeadline) - { - var plainSeen = 0; - foreach (var f in peerFrames) - if (f.StartsWith("p", StringComparison.Ordinal)) - plainSeen++; - if (plainSeen >= pendingPlain.Count) - break; - await Task.Delay(100); } - - finalCts.Cancel(); - await finalReader; - } - - // 3) 终局核对:普通消息无丢失无重复(FIFO 序列完整) - var plainFrames = peerFrames.Where(f => f.StartsWith("p", StringComparison.Ordinal)).ToArray(); - Assert.Equal(pendingPlain.Count, plainFrames.Length); - Assert.Equal(pendingPlain.ToArray(), plainFrames); - - // 4) 绑定消息帧不重复(同一标签只允许出现一次——错发/重放探测器) - var boundFrames = peerFrames.Where(f => f.StartsWith("b", StringComparison.Ordinal)).ToArray(); - Assert.Equal(boundFrames.Length, boundFrames.Distinct().Count()); - } - - /// - /// 重连竞速的长期观察(默认不运行):用户显式 Reconnect() 与自动重连(对端死亡)高占比 - /// 并发竞速——双 StartAsync 接受循环、Connected→Connected 直连迁移、会话编号线性化等 - /// #47 排查中判定"不可复现但结构脆弱"的路径,靠长时间高压力 + 全程不变式校验盯着。 - /// 动作分布:~40% 竞速重连(杀对端与 Reconnect 并发)、~20% 纯用户重连、~30% 普通发送 - /// (验证跨会话续发的 FIFO 完整性)、~10% 短等待;对端仅在无存活连接时重接入 - /// (单客户端模式:服务端 Connected 时监听已关闭,并发第二条连接被拒是正确行为)。 - /// - [Fact] - public async Task Soak_ReconnectRacing_LongRun() - { - var seconds = SoakSeconds(); - if (seconds <= 0) - return; // 未启用浸泡模式:跳过 - - var output = _output; - var seed = Environment.TickCount; - var random = new Random(seed); - var port = GetFreePort(); - - var recorder = new StateTransitionRecorder(); - await using var server = new StreamConnection( - new LengthPrefixFramer(), StringCodec.Instance, IPAddress.Loopback, port, - isActive: false, new StreamConnectionOptions { SendQueueCapacity = 256, AcceptRetryDelayMs = 200 }, - logger: new SoakDumpLogger(output)); - server.ConnectionChanged += (_, state) => recorder.Record(state, server.CurrentSessionId); - server.Start(CancellationToken.None); - - var peerFrames = new ConcurrentQueue(); - var pendingPlain = new ConcurrentQueue(); - var plainSeq = 0; - TcpClient? client = null; - - var deadline = TestClock.TickCount64 + (int)(seconds * 1000); - try - { - while (TestClock.TickCount64 < deadline) + else if (racing ? action < 95 : action < 70) { - var action = random.Next(100); - if (client is null) - { - // 对端(重)接入(仅在无存活对端时——单客户端模式:服务端 Connected - // 时监听已关闭,并发第二条连接被拒是正确行为) - client = new TcpClient(); - await ConnectWithRetryAsync(client, port); - var currentStream = client.GetStream(); - var readCts = new CancellationTokenSource(TimeSpan.FromSeconds(seconds + 60)); - _ = Task.Run(async () => - { - var header = new byte[4]; - try - { - while (true) - { - var read = 0; - while (read < 4) - { - var n = await currentStream.ReadAsync(header, read, 4 - read, readCts.Token); - if (n == 0) - return; - read += n; - } - - var length = BinaryPrimitives.ReadInt32BigEndian(header); - var payload = new byte[length]; - read = 0; - while (read < length) - { - var n = await currentStream.ReadAsync(payload, read, length - read, readCts.Token); - if (n == 0) - return; - read += n; - } - - peerFrames.Enqueue(Encoding.UTF8.GetString(payload)); - } - } - catch (Exception ex) when (ex is ObjectDisposedException or IOException or OperationCanceledException) - { - } - }); - } - else if (action < 50) - { - // 竞速重连:杀对端(触发自动重连)与用户显式 Reconnect() 同时打—— - // 两路 Retry 过渡竞争,可能产生双 StartAsync 接受循环 - var victim = client; - var kill = Task.Run(() => victim.Dispose()); - server.Reconnect(); - await kill; - await WaitForStateAsync(server, st => st != ConnectionState.Connected, timeoutMs: 15_000); - client = null; - } - else if (action < 70) - { - // 纯用户重连(对端保持存活:服务端主动关闭路径) - server.Reconnect(); - await WaitForStateAsync(server, st => st != ConnectionState.Connected, timeoutMs: 15_000); - client?.Dispose(); - client = null; - } - else if (action < 95) - { - // 普通发送:跨会话续发的 FIFO 完整性探针 - for (var i = 0; i < random.Next(1, 4); i++) - { - var label = $"p{plainSeq++}"; - pendingPlain.Enqueue(label); - await server.SendAsync(label); - } - } - else - { - await Task.Delay(random.Next(20, 100)); - } - - await Task.Delay(random.Next(10, 60)); + var count = random.Next(1, 4); + for (var i = 0; i < count; i++) await run.SendBusinessAsync(plain++); } - } - finally - { - client?.Dispose(); - var exitDeadline = TestClock.TickCount64 + 15_000; - while (TestClock.TickCount64 < exitDeadline && server.State == ConnectionState.Connected) - await Task.Delay(20); + else if (!racing) + await peer.InjectPartialAsync(); + else + await Task.Delay(random.Next(20, 100)); + await Task.Delay(racing ? random.Next(10, 60) : random.Next(20, 120)); } - // 终局:干净对端接入,等普通消息全部续发送达 - using (var finalClient = new TcpClient()) + if (peer is not null) await run.ClosePeerAsync(peer); + await run.WaitAsync(() => run.Server.State != ConnectionState.Connected, "final disconnect"); + await run.ObserveSendsAsync(); + + // A FIFO sentinel proves the surviving queue is drained BEFORE business retries. + // Every unclaimed attempt must have reached this peer; retries cannot hide a lost queue entry. + var unclaimed = run.UnclaimedPlain(); + peer = await run.ConnectAsync(); + await run.BarrierAsync(peer); + run.AssertQueueContinuation(unclaimed); + run.ReportUnconfirmed(); + for (var retry = 0; retry < 3 && run.Acked.Count < plain; retry++) { - await ConnectWithRetryAsync(finalClient, port); - await WaitForStateAsync(server, st => st == ConnectionState.Connected, timeoutMs: 30_000); - var finalStream = finalClient.GetStream(); - var finalCts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); - var finalReader = Task.Run(async () => - { - var header = new byte[4]; - try - { - while (true) - { - var read = 0; - while (read < 4) - { - var n = await finalStream.ReadAsync(header, read, 4 - read, finalCts.Token); - if (n == 0) - return; - read += n; - } - - var length = BinaryPrimitives.ReadInt32BigEndian(header); - var payload = new byte[length]; - read = 0; - while (read < length) - { - var n = await finalStream.ReadAsync(payload, read, length - read, finalCts.Token); - if (n == 0) - return; - read += n; - } - - peerFrames.Enqueue(Encoding.UTF8.GetString(payload)); - } - } - catch (Exception ex) when (ex is ObjectDisposedException or IOException or OperationCanceledException) - { - } - }); - - var finalDeadline = TestClock.TickCount64 + 30_000; - while (TestClock.TickCount64 < finalDeadline) - { - var plainSeen = 0; - foreach (var f in peerFrames) - if (f.StartsWith("p", StringComparison.Ordinal)) - plainSeen++; - if (plainSeen >= pendingPlain.Count) - break; - await Task.Delay(100); - } - - finalCts.Cancel(); - await finalReader; + for (var id = 0; id < plain; id++) + if (!run.Acked.ContainsKey(id)) await run.SendBusinessAsync(id); + await run.BarrierAsync(peer); } - - await server.DisposeAsync(); - - // 全程不变式:合法转移、编号线性化、终态收敛 + await run.WaitAsync(() => run.Acked.Count == plain, "all business ACKs"); + await run.ClosePeerAsync(peer); + await run.StopAsync(); recorder.Validate(); - - // FIFO 完整性:普通消息无丢失无重复 - var plainFrames = peerFrames.Where(f => f.StartsWith("p", StringComparison.Ordinal)).ToArray(); - Assert.Equal(pendingPlain.Count, plainFrames.Length); - Assert.Equal(pendingPlain.ToArray(), plainFrames); - - output.WriteLine($"[soak-racing] seed={seed}, 普通消息 {plainSeq} 条全部送达;" + - $"Connected→Connected 直连迁移 {recorder.DirectConnectedMigrations} 次(竞速观测指标)"); + Assert.Equal(Enumerable.Range(0, plain), run.Acked.Keys.OrderBy(x => x)); + Assert.Equal(Enumerable.Range(0, plain), run.Delivered.Keys.OrderBy(x => x)); + var observed = run.Received.ToArray(); + // Retried business IDs may repeat, but each unique wire attempt and each bound ID is sent once. + Assert.Equal(observed.Length, observed.Select(x => x.Wire).Distinct().Count()); + foreach (var session in observed.GroupBy(x => x.SessionId)) + { + var order = session.Select(x => run.Attempts[x.Wire].Sequence).ToArray(); + Assert.Equal(order.OrderBy(x => x), order); + } + foreach (var frame in observed.Where(x => x.Wire.StartsWith("b", StringComparison.Ordinal))) + Assert.Equal(long.Parse(frame.Wire.Split('/')[1], System.Globalization.CultureInfo.InvariantCulture), frame.SessionId); + trace.Log($"PASS seed={seed} acked={plain} attempts={run.Attempts.Count} bound={bound} direct-migrations={recorder.DirectConnectedMigrations}"); } }