Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .github/workflows/soak.yml
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,10 @@ on:
description: '浸泡时长(秒)'
required: false
default: '600'
seed:
description: '可选随机 seed(重放动作序列)'
required: false
default: ''

permissions:
contents: read
Expand Down Expand Up @@ -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 }}
Expand Down
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

### 新增
Expand Down
27 changes: 27 additions & 0 deletions docs/SOAK.md
Original file line number Diff line number Diff line change
@@ -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 已合并。
8 changes: 7 additions & 1 deletion src/StreamFrame/Connection/StreamConnection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
29 changes: 29 additions & 0 deletions test/StreamFrame.Tests/DeliveryProtocolTests.cs
Original file line number Diff line number Diff line change
@@ -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.
}
}
95 changes: 95 additions & 0 deletions test/StreamFrame.Tests/SendQueueContinuationTests.cs
Original file line number Diff line number Diff line change
@@ -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<string>(new LengthPrefixFramer(), StringCodec.Instance,
IPAddress.Loopback, port, isActive: false);
var written = new TaskCompletionSource<bool>(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<string>).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<bool> condition)
{
var end = TestClock.TickCount64 + 5000;
while (!condition() && TestClock.TickCount64 < end) await Task.Delay(10);
Assert.True(condition());
}

private static async Task<string> ReadAsync(TcpClient client)
{
async Task<byte[]> 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));
}
}
Loading
Loading