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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
## [Unreleased]

### 修复
- **坏头后完整帧交付([#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 必需检查名不变。

Expand Down
12 changes: 11 additions & 1 deletion src/StreamFrame/Connection/FrameDecoder.cs
Original file line number Diff line number Diff line change
Expand Up @@ -92,8 +92,18 @@ public async Task RunAsync(CancellationToken ct)
var buffer = result.Buffer;

// 切尽当前缓冲内的所有完整帧;已缓冲的字节即便会话正在停止也要投递完
while (TryDecodeFrame(ref buffer, out var payload))
while (true)
{
var previousLength = buffer.Length;
if (!TryDecodeFrame(ref buffer, out var payload))
{
// false 也可能已丢弃坏头;有消费进展时继续检查剩余缓冲。
// 无进展才等待更多输入,避免半帧导致忙循环。
if (buffer.Length < previousLength)
continue;
break;
}

TMessage message;
try
{
Expand Down
2 changes: 1 addition & 1 deletion src/StreamFrame/Framing/IFrameDiscardReporting.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,6 @@ public interface IFrameDiscardReporting
/// <param name="buffer">待解析的字节流;成功时前进到下一帧起点,失败时保留未消费字节。</param>
/// <param name="payload">切出的帧内负载(不含定界字节)。</param>
/// <param name="discarded">本次调用中被定界器丢弃的字节;无丢弃时为空序列。</param>
/// <returns>成功切出一帧返回 true;数据不足(半包)返回 false。</returns>
/// <returns>成功切出一帧返回 true;未切出帧返回 false,仍可消费无效前缀,推进契约同 <see cref="IFramer.TryDecodeFrame"/>。</returns>
bool TryDecodeFrame(ref ReadOnlySequence<byte> buffer, out ReadOnlySequence<byte> payload, out ReadOnlySequence<byte> discarded);
}
7 changes: 6 additions & 1 deletion src/StreamFrame/Framing/IFramer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,11 @@ public interface IFramer
/// </summary>
/// <param name="buffer">待解析的字节流;成功时前进到下一帧起点,失败时保留未消费字节。</param>
/// <param name="payload">切出的帧内负载(不含定界字节)。</param>
/// <returns>成功切出一帧返回 true;数据不足(半包)返回 false。</returns>
/// <returns>成功切出一帧返回 true;未切出帧返回 false(数据不足或丢弃无效字节)。</returns>
/// <remarks>
/// buffer 必须保留原缓冲的未消费后缀。返回 false 时也可通过消费前缀进行重同步;
/// 连接层在长度减少时继续解析剩余缓冲,未消费任何字节时等待更多输入。
/// 返回 true 时必须消费完整帧的线上字节(包括空负载帧的定界字节)。
/// </remarks>
bool TryDecodeFrame(ref ReadOnlySequence<byte> buffer, out ReadOnlySequence<byte> payload);
}
2 changes: 1 addition & 1 deletion src/StreamFrame/Framing/LengthPrefixFramer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ public bool TryDecodeFrame(ref ReadOnlySequence<byte> buffer, out ReadOnlySequen

var length = ReadLengthPrefix(buffer);

// 非法长度(负数 / 超上限):丢弃长度头,尝试从下一字节重新同步。
// 非法长度(负数 / 超上限):丢弃整个四字节长度头,从其后重新同步。
if ((uint)length > (uint)MaxPayloadBytes)
{
discarded = buffer.Slice(0, LengthPrefixSize);
Expand Down
110 changes: 110 additions & 0 deletions test/StreamFrame.Tests/FrameDecoderTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,116 @@ namespace StreamFrame.Tests;

public class FrameDecoderTests
{
[Theory]
[InlineData(1, false, 0)]
[InlineData(3, false, 0)]
[InlineData(3, true, 0)]
[InlineData(1, false, 100)]
public async Task Resync_CompleteFrames_DeliveredWithoutFurtherInput(int badHeaders, bool eof, int timeoutMs)
{
using var h = CreateDecoder(new LengthPrefixFramer(), incompleteFrameTimeoutMs: timeoutMs, sessionId: 42);
var bytes = Enumerable.Repeat((byte)0xFF, badHeaders * 4)
.Concat(new byte[] { 0, 0, 0, 3, 65, 65, 65, 0, 0, 0, 1, 66 }).ToArray();
await h.Pipe.Writer.WriteAsync(bytes);
if (eof) h.Pipe.Writer.Complete();
var run = h.Decoder.RunAsync(CancellationToken.None);
try
{
Assert.True(h.Relay.Reader.TryRead(out var first), "完整帧必须在等待更多输入前交付");
Assert.Equal("AAA", first.Message);
Assert.Equal(42, first.SessionId);
Assert.True(h.Relay.Reader.TryRead(out var second));
Assert.Equal("B", second.Message);
Assert.False(h.Relay.Reader.TryRead(out _));
Assert.Equal(badHeaders, h.Errors.Count);
Assert.All(h.Errors, error =>
{
Assert.Equal(FrameErrorKind.DiscardedByResync, error.Kind);
Assert.Equal(42, error.SessionId);
Assert.Equal(4, error.ObservedByteCount);
Assert.Equal(new byte[] { 255, 255, 255, 255 }, error.Bytes.ToArray());
Assert.False(error.IsTruncated);
});
if (timeoutMs > 0)
{
await Task.Delay(timeoutMs * 3);
Assert.False(run.IsCompleted); // 已切尽的缓冲不得启动半帧超时
}
}
finally
{
if (!eof) h.Pipe.Writer.Complete();
await run;
}
}

[Theory]
[InlineData(false)]
[InlineData(true)]
public async Task Resync_CustomFramer_ProgressThenPartialStopsUntilMoreInput(bool eof)
{
var framer = new ProgressFramer();
using var h = CreateDecoder(framer);
await h.Pipe.Writer.WriteAsync(new byte[] { 255, 0, 0, 0, 2, 65 });
var run = h.Decoder.RunAsync(CancellationToken.None);
try
{
Assert.Equal(2, framer.Calls); // 丢弃一次、半帧一次;不忙循环
Assert.False(h.Relay.Reader.TryRead(out _));
if (!eof) await h.Pipe.Writer.WriteAsync(new byte[] { 66 });
}
finally
{
h.Pipe.Writer.Complete();
await run;
}
if (eof) Assert.False(h.Relay.Reader.TryRead(out _));
else
{
Assert.True(h.Relay.Reader.TryRead(out var message));
Assert.Equal("AB", message.Message);
}
}

[Fact]
public async Task Resync_CustomFramer_CompleteFrameDeliveredWithoutFurtherInput()
{
using var h = CreateDecoder(new ProgressFramer());
await h.Pipe.Writer.WriteAsync(new byte[] { 255, 0, 0, 0, 1, 65 });
var run = h.Decoder.RunAsync(CancellationToken.None);
try
{
Assert.True(h.Relay.Reader.TryRead(out var message));
Assert.Equal("A", message.Message);
Assert.Empty(h.Errors); // 普通 IFramer 不提供丢弃诊断
}
finally
{
h.Pipe.Writer.Complete();
await run;
}
}

private sealed class ProgressFramer : IFramer
{
private readonly LengthPrefixFramer _inner = new();
public int Calls { get; private set; }
public int MaxPayloadBytes => _inner.MaxPayloadBytes;
public void EncodeFrame(ReadOnlySpan<byte> payload, IBufferWriter<byte> writer) => _inner.EncodeFrame(payload, writer);
public bool TryDecodeFrame(ref ReadOnlySequence<byte> buffer, out ReadOnlySequence<byte> payload)
{
Calls++;
if (Calls > 10) throw new InvalidOperationException("无进展忙循环");
if (!buffer.IsEmpty && buffer.First.Span[0] == 255)
{
buffer = buffer.Slice(1);
payload = default;
return false;
}
return _inner.TryDecodeFrame(ref buffer, out payload);
}
}

/// <summary>指标构造用的占位端点计数器(仅作标签,无实际意义)。</summary>
private static int _portCounter;

Expand Down
Loading