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]

### 修复
- **监听绑定失败恢复([#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 必需检查名不变。

## [2.6.0] - 2026-08-29
Expand Down
84 changes: 58 additions & 26 deletions src/StreamFrame/Connection/StreamConnection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -77,8 +77,10 @@ public string? RemoteIpAddress
private readonly StreamConnectionOptions _options;

#if NET9_0_OR_GREATER
private readonly Lock _listenerGate = new(); // 监听器发布与停机释放互斥,不跨 await 持有。
private readonly Lock _sessionGate = new(); // System.Threading.Lock(net9+,比 monitor 锁更轻量)
#else
private readonly object _listenerGate = new();
private readonly object _sessionGate = new();
#endif
private Pipe? _pipe;
Expand Down Expand Up @@ -263,8 +265,16 @@ private async Task StartAsync(CancellationToken ct, int loopId)
private static Socket CreateTcpSocket()
{
var socket = new Socket(SocketType.Stream, ProtocolType.Tcp);
socket.DualMode = true;
return socket;
try
{
socket.DualMode = true;
return socket;
}
catch
{
socket.Dispose();
throw;
}
}

/// <summary>双栈监听地址归一:0.0.0.0 绑定到 IPv6 的 ::(v4 流量经映射地址到达)。</summary>
Expand Down Expand Up @@ -317,17 +327,15 @@ private async Task<Socket> ConnectAsync(CancellationToken ct)
if (Volatile.Read(ref _acceptLoopId) != loopId)
return null;

if (_server == null)
InitServer();
var listener = InitServer(ct);

try
{
#if NETSTANDARD2_0
var listener = _server!;
// FromAsync 无法取消;停机时 listener 被释放,挂起的 accept 以异常收尾
accepted = await Task.Factory.FromAsync(listener.BeginAccept, listener.EndAccept, null).ConfigureAwait(false);
#else
accepted = await _server!.AcceptAsync(ct).ConfigureAwait(false);
accepted = await listener.AcceptAsync(ct).ConfigureAwait(false);
#endif
}
catch (Exception ex) when (!ct.IsCancellationRequested && !IsDisposed)
Expand All @@ -344,16 +352,23 @@ private async Task<Socket> ConnectAsync(CancellationToken ct)
// 单客户端模式:accept 到第一个客户端后关闭监听 socket,
// 后续连接在 TCP 层被立即拒绝。代次门控保证此处 _server 属于当代循环,
// 不会误关/漏关其它循环的监听器。
if (_options.AcceptFirstClientOnly && _server != null)
if (_options.AcceptFirstClientOnly)
{
_server.Dispose();
_server = null;
lock (_listenerGate)
{
_server?.Dispose();
_server = null;
}
}
}
}
finally
{
_acceptLock.Release();
lock (_listenerGate)
{
if (!IsDisposed)
_acceptLock.Release();
}
}

if (accepted is not null)
Expand All @@ -363,22 +378,36 @@ private async Task<Socket> ConnectAsync(CancellationToken ct)
}
}

private void InitServer()
private Socket InitServer(CancellationToken ct)
{
if (_server != null)
lock (_listenerGate)
{
_server.Dispose();
_server = null;
}
ct.ThrowIfCancellationRequested();
// Shutdown 先标记终态,再取得同一把锁释放;迟到的初始化不得发布监听器。
if (IsDisposed)
throw new OperationCanceledException(ct);
if (_server != null)
return _server;

_server = CreateTcpSocket();
_server.Blocking = false;
// 允许重绑覆盖遗留的 TIME_WAIT:服务端主动关闭(用户 Reconnect/停机)后立即重新
// 监听不受 2MSL 限制(Linux 上没有该选项会遇到 EADDRINUSE;Windows 实测宽松,
// 一并设置保持跨平台行为一致——#47 防御性修复)
_server.SetSocketOption(SocketOptionLevel.Socket, SocketOptionName.ReuseAddress, true);
_server.Bind(new IPEndPoint(NormalizeListenAddress(IpAddress), Port));
_server.Listen(0);
var listener = CreateTcpSocket();
try
{
listener.Blocking = false;
// 允许重绑覆盖遗留的 TIME_WAIT:服务端主动关闭(用户 Reconnect/停机)后立即重新
// 监听不受 2MSL 限制(Linux 上没有该选项会遇到 EADDRINUSE;Windows 实测宽松,
// 一并设置保持跨平台行为一致——#47 防御性修复)
listener.SetSocketOption(SocketOptionLevel.Socket, SocketOptionName.ReuseAddress, true);
listener.Bind(new IPEndPoint(NormalizeListenAddress(IpAddress), Port));
listener.Listen(0);
_server = listener; // 全部成功才转移所有权;失败保持字段为空,下一轮可重试。
return listener;
}
catch
{
listener.Dispose();
throw;
}
}
}

/// <summary>统一配置已连接 socket:非阻塞、接收缓冲、可选 TCP KeepAlive(半开连接探测)。</summary>
Expand Down Expand Up @@ -1149,11 +1178,14 @@ private void Shutdown()
_socket = null;
}

_server?.Dispose();
_server = null;
lock (_listenerGate)
{
_server?.Dispose();
_server = null;
_acceptLock.Dispose();
}

_sendLock.Dispose();
_acceptLock.Dispose();
_lifetimeCts?.Dispose();
_lifetimeCts = null;
}
Expand Down
212 changes: 212 additions & 0 deletions test/StreamFrame.Tests/ListenerRetryTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,212 @@
using System.Net;
using System.Net.Sockets;
using System.Reflection;
using Microsoft.Extensions.Logging;

namespace StreamFrame.Tests;

public class ListenerRetryTests
{
private static readonly TimeSpan Budget = TimeSpan.FromSeconds(10);

// 日志在 Bind 失败已退出初始化之后发出;阻塞此处可精确控制释放端口/停机的顺序。
private sealed class FailureGate : ILogger, IDisposable
{
public TaskCompletionSource<bool> Failed { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
public ManualResetEventSlim Continue { get; } = new(false);
public IDisposable? BeginScope<TState>(TState state) where TState : notnull => null;
public bool IsEnabled(LogLevel logLevel) => true;
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception,
Func<TState, Exception?, string> formatter)
{
if (!formatter(state, exception).StartsWith("Connect to ", StringComparison.Ordinal))
return;
Failed.TrySetResult(true);
if (!Continue.Wait(Budget))
throw new TimeoutException("测试未释放绑定失败同步点");
}
public void Dispose() => Continue.Dispose();
}

private static TcpListener Occupy(int port = 0)
{
var listener = new TcpListener(IPAddress.Loopback, port);
listener.Server.ExclusiveAddressUse = true;
listener.Start();
return listener;
}

private static StreamConnection<string> Create(int port, bool active, ILogger? logger = null)
=> new(new LengthPrefixFramer(), StringCodec.Instance, IPAddress.Loopback, port, active,
new StreamConnectionOptions { AcceptRetryDelayMs = 25, ConnectRetryDelayMs = 25 }, logger);

private static object? RetainedListener(StreamConnection<string> server)
=> typeof(StreamConnection<string>).GetField("_server", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(server);

private static async Task<string> Receive(StreamConnection<string> connection)
{
using var timeout = new CancellationTokenSource(Budget);
await foreach (var message in connection.GetMessages(timeout.Token))
return message;
throw new InvalidOperationException("消息流提前结束");
}

[Theory]
[InlineData(false)]
[InlineData(true)]
public async Task PublishedListener_StopClosesSocketAndReleasesPort(bool cancel)
{
var occupied = Occupy();
var port = ((IPEndPoint)occupied.LocalEndpoint).Port;
occupied.Stop();
using var lifetime = new CancellationTokenSource();
await using var server = Create(port, false);
// Start 同步运行到首次未完成的 Accept;此时监听已发布、尚无客户端。
server.Start(lifetime.Token);
var listener = Assert.IsType<Socket>(RetainedListener(server));
Assert.True(listener.IsBound);
if (cancel)
lifetime.Cancel();
else
await server.DisposeAsync();
Assert.Null(RetainedListener(server));
Assert.Throws<ObjectDisposedException>(() => listener.GetSocketOption(SocketOptionLevel.Socket, SocketOptionName.ReuseAddress));
var replacement = Occupy(port);
replacement.Stop();
}

[Theory]
[InlineData(false)]
[InlineData(true)]
public async Task Initialization_QueuedBehindShutdown_CannotPublish(bool cancel)
{
var occupied = Occupy();
var port = ((IPEndPoint)occupied.LocalEndpoint).Port;
using var lifetime = new CancellationTokenSource();
using var failure = new FailureGate();
await using var server = Create(port, false, failure);
var start = Task.Run(() => server.Start(lifetime.Token));
try
{
await failure.Failed.Task.WaitAsync(Budget);
var type = typeof(StreamConnection<string>);
var listenerGate = type.GetField("_listenerGate", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(server)!;
var initialize = type.GetMethod("InitServer", BindingFlags.Instance | BindingFlags.NonPublic)!;
using var attempting = new ManualResetEventSlim();
using var stopping = new ManualResetEventSlim();
server.ConnectionChanged += (_, state) =>
{
if (state == ConnectionState.Disconnected)
stopping.Set();
};
Task lateInitialization;
Task shutdown;
// 白盒门控准确覆盖:重试已排队,停机标志已发布,随后才允许初始化取锁。
#if NET9_0_OR_GREATER
lock ((Lock)listenerGate)
#else
lock (listenerGate)
#endif
{
lateInitialization = Task.Run(() =>
{
attempting.Set();
var error = Assert.Throws<TargetInvocationException>(() => initialize.Invoke(server, new object[] { CancellationToken.None }));
Assert.IsAssignableFrom<OperationCanceledException>(error.InnerException);
});
Assert.True(attempting.Wait(Budget));
shutdown = Task.Run(async () =>
{
if (cancel)
lifetime.Cancel();
else
await server.DisposeAsync();
});
Assert.True(stopping.Wait(Budget));
occupied.Stop();
}
await Task.WhenAll(lateInitialization, shutdown).WaitAsync(Budget);
failure.Continue.Set();
await start.WaitAsync(Budget);
Assert.Null(RetainedListener(server));
var replacement = Occupy(port);
replacement.Stop();
}
finally
{
failure.Continue.Set();
occupied.Stop();
await start.WaitAsync(Budget);
}
}

[Fact]
public async Task BindFailure_ReleasedPort_AutomaticallyConnectsAndExchangesMessages()
{
var occupied = Occupy();
var port = ((IPEndPoint)occupied.LocalEndpoint).Port;
using var gate = new FailureGate();
await using var server = Create(port, false, gate);
var start = Task.Run(() => server.Start(CancellationToken.None));
try
{
await gate.Failed.Task.WaitAsync(Budget);
occupied.Stop();
gate.Continue.Set();
await start.WaitAsync(Budget);
await using var client = Create(port, true);
client.Start(CancellationToken.None);
await server.WaitForConnectedAsync().WaitAsync(Budget);
await client.WaitForConnectedAsync().WaitAsync(Budget);
await client.SendAsync("after-bind-failure");
Assert.Equal("after-bind-failure", await Receive(server));
await server.SendAsync("reply");
Assert.Equal("reply", await Receive(client));
Assert.Null(RetainedListener(server)); // 单客户端接受后必须关闭监听器。
using var second = new TcpClient();
await Assert.ThrowsAsync<SocketException>(() => second.ConnectAsync(IPAddress.Loopback, port).WaitAsync(Budget));
}
finally
{
gate.Continue.Set();
occupied.Stop();
await start.WaitAsync(Budget);
}
}

[Theory]
[InlineData(false)]
[InlineData(true)]
public async Task BindFailure_StopBeforeRetry_DoesNotRetainOrRepublishListener(bool cancel)
{
var occupied = Occupy();
var port = ((IPEndPoint)occupied.LocalEndpoint).Port;
using var lifetime = new CancellationTokenSource();
using var gate = new FailureGate();
await using var server = Create(port, false, gate);
var start = Task.Run(() => server.Start(lifetime.Token));
try
{
await gate.Failed.Task.WaitAsync(Budget);
Assert.Null(RetainedListener(server));
if (cancel)
lifetime.Cancel();
else
await server.DisposeAsync();
occupied.Stop();
gate.Continue.Set();
await start.WaitAsync(Budget);
Assert.True(server.IsDisposed);
Assert.Equal(ConnectionState.Disconnected, server.State);
Assert.Null(RetainedListener(server));
var replacement = Occupy(port);
replacement.Stop();
}
finally
{
gate.Continue.Set();
occupied.Stop();
await start.WaitAsync(Budget);
}
}
}
Loading