using System.Buffers;
using System.ComponentModel;
using System.Diagnostics;
using System.Net;
using System.Net.Sockets;
using System.Runtime.CompilerServices;
using System.Security.Authentication;
using System.Security.Cryptography.X509Certificates;
using System.Text;
using NewLife;
using NewLife.Data;
using NewLife.Messaging;
using NewLife.Net;
using Xunit;
namespace XUnitTest.Net;
// 协议交付契约:处理器返回后消息即收尾,处理器内需同步完成负载消费(慢路径等待为同连接串行语义),不适用 xUnit1031
#pragma warning disable xUnit1031
/// <summary>协议模式(Protocol 属性 + 消息泵)实网回环测试</summary>
[Collection("Net")]
public class MessageSessionTests
{
#region 工具
private static async Task<T> WithTimeout<T>(Task<T> task, Int32 timeoutMs = 10_000)
{
var completed = await Task.WhenAny(task, Task.Delay(timeoutMs)).ConfigureAwait(false);
if (completed != task) throw new TimeoutException("等待消息超时");
return await task.ConfigureAwait(false);
}
private static Byte[] MakePayload(Int32 count)
{
var buf = new Byte[count];
for (var i = 0; i < count; i++) buf[i] = (Byte)(i * 31 + 7);
return buf;
}
private static TaskCompletionSource<T> NewTcs<T>() => new(TaskCreationOptions.RunContinuationsAsynchronously);
/// <summary>加载测试自签名证书(内嵌 pfx)</summary>
private static X509Certificate2 LoadTestCert()
{
var pfx = typeof(MessageSessionTests).Assembly.GetManifestResourceStream("XUnitTest.certs.newlifex.com.pfx")!.ReadBytes(-1);
#if NET9_0_OR_GREATER
return X509CertificateLoader.LoadPkcs12(pfx, "123456");
#else
return new X509Certificate2(pfx, "123456", X509KeyStorageFlags.DefaultKeySet);
#endif
}
#endregion
#region 匹配队列
[Fact]
[DisplayName("匹配队列_超时_取消池化等待源")]
public async Task MatchQueue_Timeout_CancelsSource()
{
var queue = new DefaultMatchQueue();
var source = PooledValueTaskSource<Message>.Rent();
queue.Add(this, new DefaultMessage { Sequence = 1 }, 300, source);
var task = source.ValueTask.AsTask();
// 超时取消由队列内部约 1 秒周期的检查定时器驱动,回调还要经线程池调度,故上限给足余量。
// 队列在有请求未结束期间被定时器持有,用例不必额外保活(见 MatchQueue_QueueNotReferenced_StillCancels)
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => WithTimeout(task, 60_000));
}
[Fact]
[DisplayName("匹配队列_入队后失去引用_超时仍取消等待方")]
public async Task MatchQueue_QueueNotReferenced_StillCancels()
{
// 入队后队列即不可达,等待期间持续强制 GC 让它必定被回收
var task = EnqueueWithoutRoot(300);
for (var i = 0; i < 40 && !task.IsCompleted; i++)
{
// 只用第0代回收:队列是刚分配的年轻对象,足以回收仅被弱引用持有的队列,又不必整轮阻塞回收
GC.Collect(0);
await Task.Delay(50);
}
// 队列被自身定时器持有,超时照常取消;若定时器随队列一起被回收,这里会一直等到上限才以 TimeoutException 失败
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => WithTimeout(task, 30_000));
}
/// <summary>入队后就不再引用队列,返回等待方任务</summary>
[MethodImpl(MethodImplOptions.NoInlining)]
private static Task<Message> EnqueueWithoutRoot(Int32 msTimeout)
{
var queue = new DefaultMatchQueue();
var source = PooledValueTaskSource<Message>.Rent();
queue.Add(null, new DefaultMessage { Sequence = 9 }, msTimeout, source);
return source.ValueTask.AsTask();
}
[Fact]
[DisplayName("匹配队列_清空_取消池化等待源")]
public async Task MatchQueue_Clear_CancelsSource()
{
var queue = new DefaultMatchQueue();
var source = PooledValueTaskSource<Message>.Rent();
queue.Add(this, new DefaultMessage { Sequence = 1 }, 300, source);
var task = source.ValueTask.AsTask();
// 与超时路径共用同一套取消逻辑,但同步触发、不依赖定时器调度,因此不会因机器负载而飘
queue.Clear();
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => task);
}
[Fact]
[DisplayName("匹配队列_回调命中_结果交付池化源")]
public void MatchQueue_Match_DeliversSource()
{
var queue = new DefaultMatchQueue();
var source = PooledValueTaskSource<Message>.Rent();
var matcher = new SrmpCodec();
var req = new DefaultMessage { Sequence = 1 };
queue.Add(this, req, 10_000, source);
var resp = new DefaultMessage { Sequence = 1, Kind = MessageKinds.Response };
var ok = queue.Match(this, resp, resp, (rq, rs) => rq is Message a && rs is Message b && matcher.Match(a, b));
Assert.True(ok);
Assert.True(source.ValueTask.IsCompleted);
Assert.Same(resp, source.ValueTask.Result);
}
[Fact]
[DisplayName("协议匹配_序列号超过255_按低8位配对")]
public void SrmpMatch_LargeSequence_MatchesLowByte()
{
var matcher = new SrmpCodec();
var req = new DefaultMessage { Sequence = 300 };
// 线格式只带 1 字节序列号,对端回显的是低 8 位
var resp = new DefaultMessage { Sequence = 300 & 0xFF, Kind = MessageKinds.Response };
Assert.NotEqual(req.Sequence, resp.Sequence);
Assert.True(matcher.Match(req, resp));
// 低 8 位不同则不配对
resp.Sequence = (300 & 0xFF) + 1;
Assert.False(matcher.Match(req, resp));
}
[Fact]
[DisplayName("匹配队列_等待方已取消_完成失败按未命中返回")]
public void MatchQueue_WaiterCanceled_ReturnsFalse()
{
var queue = new DefaultMatchQueue();
var source = PooledValueTaskSource<Message>.Rent();
var matcher = new SrmpCodec();
var req = new DefaultMessage { Sequence = 8 };
queue.Add(this, req, 10_000, source);
// 等待方先行放弃(取消),而队列项仍在;此时迟到响应到达
Assert.True(source.TrySetCanceled());
var resp = new DefaultMessage { Sequence = 8, Kind = MessageKinds.Response };
var ok = queue.Match(this, resp, resp, (rq, rs) => rq is Message a && rs is Message b && matcher.Match(a, b));
// 完成失败必须按“未命中”返回:调用方据此丢弃负载并释放消息,否则消息与池化缓冲泄漏
Assert.False(ok);
}
[Fact]
[DisplayName("匹配队列_池化源回收复用_残留项不得完成新请求")]
public async Task MatchQueue_StaleItem_DoesNotCompleteReusedSource()
{
var queue = new DefaultMatchQueue();
var matcher = new SrmpCodec();
// 第一次等待:入队后取消,等待方结束并归还池
var first = PooledValueTaskSource<Message>.Rent();
var oldRequest = new DefaultMessage { Sequence = 0x21 };
queue.Add(this, oldRequest, 10_000, first);
var oldTask = first.ValueTask.AsTask();
Assert.True(first.TrySetCanceled());
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => oldTask);
// 归还后同一实例通常被下一个请求借出复用(池为 LIFO)。不硬断言 Same:xUnit 并行下
// 可能被其它用例抢走,断言失败会变成与产品缺陷无关的假红;两种情形都必须“不得完成新等待源”
var second = PooledValueTaskSource<Message>.Rent();
var newRequest = new DefaultMessage { Sequence = 0x22 };
queue.Add(this, newRequest, 10_000, second);
// 迟到响应只能匹配上残留的旧队列项:完成必须失败,且不得完成复用后的新请求
var lateResponse = new DefaultMessage { Sequence = 0x21, Kind = MessageKinds.Response };
var ok = queue.Match(this, lateResponse, lateResponse, (rq, rs) => rq is Message a && rs is Message b && matcher.Match(a, b));
Assert.False(ok);
Assert.False(second.ValueTask.IsCompleted);
}
#endregion
[Fact]
[DisplayName("协议模式_单帧分两段到达_接收方粘包重组")]
public async Task SplitFrame_Reassembled()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var got = NewTcs<Byte[]>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage msg) return;
// 先到一半的帧不会交付(头部定界后等体到齐),到齐后才进本事件;
// 交付契约:处理器返回后消息收尾,负载需在本方法内读完
var all = msg.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
got.TrySetResult(body);
};
// 裸套接字把一整帧分两段发出:接收方必须自行重组,不得依赖发送方“一段写完”
using var client = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
await client.ConnectAsync(IPAddress.Loopback, server.Port);
var codec = new SrmpCodec();
var msg = new DefaultMessage();
msg.SetBody(new ArrayPacket("Split Frame Body"u8.ToArray()));
using var frame = codec.Build(msg);
var bytes = frame.ToArray();
Assert.True(bytes.Length > 6, "帧至少包含头部与部分负载,才能切成两段");
// 第一段:头部 + 部分负载,此时接收方只能定界不能成帧;延迟后补发剩余负载
client.Send(bytes, 0, 5, SocketFlags.None);
await Task.Delay(50);
client.Send(bytes, 5, bytes.Length - 5, SocketFlags.None);
var body = await WithTimeout(got.Task, 5_000);
Assert.Equal("Split Frame Body"u8.ToArray(), body);
}
[Fact]
[DisplayName("协议模式_小消息往返_服务端应答")]
public async Task SmallMessage_RoundTrip()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var serverGot = NewTcs<(Int32 Sequence, Boolean OneWay, Byte[] Body)>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
// 交付契约:处理器返回后消息收尾(未读体丢弃、消息释放);
// 需要负载时在本方法内读完(数据未到齐会等待——同连接消息串行)
var all = req.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
serverGot.TrySetResult((req.Sequence, req.Kind == MessageKinds.OneWay, body));
// 服务端应答(经协议构建整帧发送)
var reply = (DefaultMessage)req.CreateReply();
reply.SetBody(new ArrayPacket(new Byte[] { 0x0B, 0x0C }));
((INetSession)s!).SendMessage(reply);
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
var clientGot = NewTcs<(Int32 Sequence, Boolean Reply, Byte[] Body)>();
client.Received += (s, e) =>
{
if (e.Message is not DefaultMessage msg || msg.Kind < MessageKinds.Response) return;
// 处理器内读完负载(交付契约:返回后消息收尾)
var all = msg.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
clientGot.TrySetResult((msg.Sequence, msg.Kind >= MessageKinds.Response, body));
};
Assert.True(client.Open());
var reqMsg = new DefaultMessage { Sequence = 0x35 };
reqMsg.SetBody(new ArrayPacket(new Byte[] { 0x01, 0x02, 0x03 }));
client.SendMessage(reqMsg);
// 服务端收到请求:字段与负载完整
var (seq, oneWay, body) = await WithTimeout(serverGot.Task);
Assert.Equal(0x35, seq);
Assert.False(oneWay);
Assert.Equal(new Byte[] { 0x01, 0x02, 0x03 }, body);
// 客户端收到应答:Reply 且序列号配对
var (rseq, replyFlag, replyBody) = await WithTimeout(clientGot.Task);
Assert.True(replyFlag);
Assert.Equal(0x35, rseq);
Assert.Equal(new Byte[] { 0x0B, 0x0C }, replyBody);
}
[Fact]
[DisplayName("协议模式_大帧_流式读取完整")]
public async Task LargeMessage_Streaming()
{
var payload = MakePayload(300_000);
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var got = NewTcs<Byte[]>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
try
{
// 头部到齐即已交付,负载流式读取;大帧数据未到齐时等待(同连接消息串行)
var all = req.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
got.TrySetResult(body);
}
catch (Exception ex)
{
got.TrySetException(ex);
}
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
Assert.True(client.Open());
var reqMsg = new DefaultMessage { Sequence = 0x11 };
reqMsg.SetBody(new ArrayPacket(payload));
client.SendMessage(reqMsg);
var body = await WithTimeout(got.Task);
Assert.Equal(payload, body);
}
[Fact]
[DisplayName("协议模式_单向消息_不等待响应")]
public async Task OneWayMessage_NoReply()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var got = NewTcs<Boolean>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
// 仅头部语义即可处理:不读负载,未读体由消息泵在处理器返回后丢弃对齐
got.TrySetResult(req.Kind == MessageKinds.OneWay);
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
Assert.True(client.Open());
var reqMsg = new DefaultMessage { Sequence = 0x22, Kind = MessageKinds.OneWay };
reqMsg.SetBody(new ArrayPacket(new Byte[] { 1 }));
client.SendMessage(reqMsg);
Assert.True(await WithTimeout(got.Task));
}
[Fact]
[DisplayName("协议模式_请求响应_客户端等待匹配")]
public async Task RequestResponse_ClientAwaits()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var serverGot = NewTcs<Int32>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
// 回显应答:物化请求负载作为响应体(所有权转移)
var all = req.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var reply = (DefaultMessage)req.CreateReply();
reply.SetBody(all);
serverGot.TrySetResult(((INetSession)s!).SendMessage(reply));
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
var pushed = 0;
client.Received += (s, e) => Interlocked.Increment(ref pushed);
Assert.True(client.Open());
var req = new DefaultMessage { Sequence = 0x42 };
req.SetBody(new ArrayPacket(new Byte[] { 5, 6, 7 }));
var reqTask = client.SendMessageAsync(req).AsTask();
// 服务端已收到请求并完成应答发送
var sent = await WithTimeout(serverGot.Task, 5_000);
Assert.True(sent > 0, "服务端应答发送失败");
var resp = (DefaultMessage)await WithTimeout(reqTask, 8_000);
Assert.Equal(MessageKinds.Response, resp.Kind);
Assert.Equal(0x42, resp.Sequence);
// 响应先进入事件链(可观测),随后匹配交付等待方;交付后跳过收尾,消息由等待方释放
Assert.Equal(1, pushed);
// 交付消息体为内存模式(流式负载已物化),等待方可异步消费
Assert.False(resp.Body!.IsStreaming);
var body = await resp.Body.ReadAllAsync();
Assert.Equal(new Byte[] { 5, 6, 7 }, body.AsReadOnlySequence().ToArray());
body.TryDispose();
// 等待方负责释放响应消息
resp.Dispose();
}
[Fact]
[DisplayName("协议模式_长度字段_大响应流式分片_不串包")]
public async Task LengthFieldCodec_LargeStreamingResponse_NoCorruption()
{
// 服务端:回一个大负载(大于单次接收缓冲 → 分片到达 → 流式体)
var payload = MakePayload(200_000);
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new LengthFieldCodec { Size = 4 } };
server.Start();
var serverGot = NewTcs<Int32>();
server.Received += (s, e) =>
{
if (e.Message == null) return;
var reply = new Message();
reply.SetBody(new ArrayPacket(payload));
serverGot.TrySetResult(((INetSession)s!).SendMessage(reply));
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}")
{
Protocol = new LengthFieldCodec { Size = 4 },
MatchTimeout = 10_000,
};
Assert.True(client.Open());
var req = new Message();
req.SetBody(new ArrayPacket(new Byte[] { 1, 2, 3 }));
var reqTask = client.SendMessageAsync(req).AsTask();
var sent = await WithTimeout(serverGot.Task, 5_000);
Assert.True(sent > 0, "服务端应答发送失败");
// 无方向位协议(LengthFieldCodec 的恒真 matcher)消息 Kind 恒为 Request:等待方必须拿到已物化的完整负载。
// 修复前物化门槛用 message.Reply,流式体不物化,等待方与帧泵争读同一读取器 → 抛单读者异常或把下一帧字节当成体
var resp = await WithTimeout(reqTask, 15_000);
Assert.False(resp.Body!.IsStreaming);
var body = await resp.Body.ReadAllAsync();
Assert.Equal(payload, body.AsReadOnlySequence().ToArray());
body.TryDispose();
resp.Dispose();
}
[Fact]
[DisplayName("协议模式_请求响应_序列号超过255仍能配对")]
public async Task RequestResponse_SequenceOver255_StillMatches()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var serverGot = NewTcs<Int32>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
// 回显应答:物化请求负载作为响应体(所有权转移)
var all = req.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var reply = (DefaultMessage)req.CreateReply();
reply.SetBody(all);
serverGot.TrySetResult(((INetSession)s!).SendMessage(reply));
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}")
{
Protocol = new SrmpCodec(),
MatchTimeout = 3_000,
};
Assert.True(client.Open());
// 客户端用自增计数器:序列号超过 255 时线格式只保留低 8 位。
// 若按完整 Int32 比较则恒不配对,只能等 MatchTimeout 超时取消
var req = new DefaultMessage { Sequence = 300 };
req.SetBody(new ArrayPacket(new Byte[] { 7, 8 }));
var reqTask = client.SendMessageAsync(req).AsTask();
var sent = await WithTimeout(serverGot.Task, 5_000);
Assert.True(sent > 0, "服务端应答发送失败");
// 配对超时内完成交付:修复前此处必然抛超时取消
var resp = (DefaultMessage)await WithTimeout(reqTask, 8_000);
Assert.Equal(MessageKinds.Response, resp.Kind);
Assert.Equal(300 & 0xFF, resp.Sequence);
var body = await resp.Body!.ReadAllAsync();
Assert.Equal(new Byte[] { 7, 8 }, body.AsReadOnlySequence().ToArray());
body.TryDispose();
resp.Dispose();
}
[Fact]
[DisplayName("协议模式_请求响应_无应答超时取消")]
public async Task RequestResponse_Timeout()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
// 服务端收到请求但不作答
var got = NewTcs<Boolean>();
server.Received += (s, e) => got.TrySetResult(true);
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}")
{
Protocol = new SrmpCodec(),
MatchTimeout = 300,
};
Assert.True(client.Open());
var req = new DefaultMessage { Sequence = 0x77 };
req.SetBody(new ArrayPacket(new Byte[] { 1 }));
var reqTask = client.SendMessageAsync(req).AsTask();
// 服务端应能收到请求
Assert.True(await WithTimeout(got.Task, 5_000), "服务端未收到请求");
// 匹配队列超时(定时器秒级精度)取消等待方
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => WithTimeout(reqTask, 10_000));
}
[Fact]
[DisplayName("协议模式_超限残余_会话被关闭")]
public async Task GarbageData_ExceedsMaxCache_ClosesConnection()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec(), MaxCache = 256 };
server.Start();
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { AutoReconnect = false };
client.Open();
var closed = NewTcs<Boolean>();
client.Closed += (s, e) => closed.TrySetResult(true);
// 原始字节发送:0xFFFF 扩展头声明负数长度,永远无法定界
var garbage = new Byte[1024];
for (var i = 0; i < garbage.Length; i++) garbage[i] = 0xFF;
client.Send(garbage);
// 服务端泵判定协议错误并关闭会话,客户端感知断开
var task = await Task.WhenAny(closed.Task, Task.Delay(5_000));
Assert.True(task == closed.Task, "客户端应在超时内感知服务端关闭连接");
Assert.True(await closed.Task);
}
[Fact]
[DisplayName("协议模式_对端发完即关_已到达帧先于会话关闭交付")]
public async Task PeerSendThenClose_FrameDeliveredBeforeClose()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var codec = new SrmpCodec();
var watch = Stopwatch.StartNew();
var busy = new ManualResetEventSlim(false);
var entered = NewTcs<Boolean>();
var second = NewTcs<Int64>();
var closed = NewTcs<Int64>();
// 关闭信号取底层套接字会话的销毁事件:上层正是在它之后才判定“异常断开”并做善后
//(如协议层按是否收到正常断开报文决定是否发布遗嘱),故“已到达帧先于它交付”是可观察契约
server.NewSession += (s, e) =>
{
if (e.Session?.Session is SessionBase sk) sk.OnDisposed += (s2, e2) => closed.TrySetResult(watch.ElapsedTicks);
};
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage msg) return;
// 首帧卡住消息泵:制造“次帧已到达管道、泵仍在交付首帧”的窗口
if (msg.Sequence == 1)
{
entered.TrySetResult(true);
busy.Wait(5_000);
}
else if (msg.Sequence == 2)
second.TrySetResult(watch.ElapsedTicks);
};
using var client = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
await client.ConnectAsync(IPAddress.Loopback, server.Port);
// 首帧发完先等处理器确实卡住,保证次帧到达时泵仍在交付首帧(拉取时机确定,不依赖线程池调度)
client.Send(BuildFrame(codec, 1, MakePayload(16)));
await WithTimeout(entered.Task, 5_000);
// 次帧 + 立即关闭(对端发完即关,不等应答):次帧已到达管道,关闭此时才推进
client.Send(BuildFrame(codec, 2, MakePayload(16)));
client.Shutdown(SocketShutdown.Send);
client.Close();
// 留出关闭流程推进的时间,再放行首帧处理器
await Task.Delay(200);
busy.Set();
var secondAt = await WithTimeout(second.Task, 5_000);
var closedAt = await WithTimeout(closed.Task, 5_000);
Assert.True(secondAt < closedAt, $"已到达帧应在会话关闭信号前交付(交付 {secondAt},关闭 {closedAt})");
}
private static Byte[] BuildFrame(SrmpCodec codec, Int32 sequence, Byte[] body)
{
var msg = new DefaultMessage { Sequence = sequence };
msg.SetBody(new ArrayPacket(body));
using var frame = codec.Build(msg);
return frame.ToArray();
}
[Fact]
[DisplayName("协议模式_流式发送_服务端收流式体")]
public async Task StreamingSend_RoundTrip()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
var serverGot = NewTcs<Byte[]>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
var all = req.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
serverGot.TrySetResult(body);
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
client.Open();
// 流式发送 300KB:头部先行 + 流内容分块(内容不整载)
var payload = MakePayload(300_000);
using var stream = new MemoryStream(payload);
var sent = await client.SendMessageAsync(new DefaultMessage { Sequence = 0x33 }, stream, payload.Length);
Assert.Equal(payload.Length, sent);
Assert.Equal(payload, await WithTimeout(serverGot.Task));
}
#region 会话处理器
/// <summary>协议宿主:走 INetHandler 处理器分发(验证处理器可经事件参数取得消息)</summary>
public class HandlerServer : NetServer<HandlerSession>
{
/// <summary>处理器收到的消息体</summary>
public TaskCompletionSource<String> Got { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
public override INetHandler? CreateHandler(INetSession session) => new MessageNetHandler(this);
}
/// <summary>会话处理器:从事件参数转型取消息(协议模式契约)</summary>
public class MessageNetHandler : INetHandler
{
private readonly HandlerServer _server;
public MessageNetHandler(HandlerServer server) => _server = server;
public void Init(INetSession session) { }
public void Process(IData data)
{
// 契约:data 实际为 ReceivedEventArgs,协议模式下消息在 Message
if (data is not ReceivedEventArgs e || e.Message is not DefaultMessage msg) return;
var all = msg.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
_server.Got.TrySetResult(Encoding.UTF8.GetString(body));
}
}
public class HandlerSession : NetSession<HandlerServer> { }
#endregion
[Fact]
[DisplayName("协议模式_会话处理器_收到带消息的事件参数")]
public async Task Protocol_NetHandler()
{
using var server = new HandlerServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec() };
server.Start();
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
client.Open();
var msg = new DefaultMessage { Sequence = 1 };
msg.SetBody(new ArrayPacket(Encoding.UTF8.GetBytes("via-handler")));
client.SendMessage(msg);
Assert.Equal("via-handler", await WithTimeout(server.Got.Task));
}
[Fact]
[DisplayName("协议模式_SSL回环_流式消息完整")]
public async Task SslProtocol_RoundTrip()
{
using var cert = LoadTestCert();
using var server = new NetServer
{
Port = 0,
ProtocolType = NetType.Tcp,
SslProtocol = SslProtocols.Tls12,
Certificate = cert,
Protocol = new SrmpCodec(),
};
server.Start();
var serverGot = NewTcs<Byte[]>();
server.Received += (s, e) =>
{
if (e.Message is not DefaultMessage req) return;
var all = req.Body!.ReadAllAsync().AsTask().GetAwaiter().GetResult();
var body = all.AsReadOnlySequence().ToArray();
all.TryDispose();
serverGot.TrySetResult(body);
};
using var client = new TcpSession
{
Remote = new NetUri($"tcp://127.0.0.1:{server.Port}"),
SslProtocol = SslProtocols.Tls12,
Protocol = new SrmpCodec(),
};
client.Open();
// 大帧(300KB)经 SSL 流:服务端泵按流式体读取
var payload = MakePayload(300_000);
var msg = new DefaultMessage { Sequence = 0x61 };
msg.SetBody(new ArrayPacket(payload));
client.SendMessage(msg);
Assert.Equal(payload, await WithTimeout(serverGot.Task, 15_000));
}
#region 服务端并行
[Fact]
[DisplayName("协议模式_服务端并行_慢请求不阻塞快请求")]
public async Task Server_Parallel_SlowRequest_DoesNotBlock()
{
using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp, Protocol = new SrmpCodec(), MaxConcurrency = 8 };
server.Start();
var slowIn = NewTcs<Boolean>();
using var slowGate = new ManualResetEventSlim();
server.Received += (s, e) =>
{
if (s is not INetSession session || e.Message is not DefaultMessage req) return;
var reply = req.CreateReply();
switch (req.Sequence)
{
case 0x51:
// 慢请求:在处理链内阻塞,直到测试放行
slowIn.TrySetResult(true);
slowGate.Wait(10_000);
reply.SetBody(new ArrayPacket("slow"u8.ToArray()));
session.SendMessage(reply);
break;
case 0x52:
reply.SetBody(new ArrayPacket("fast"u8.ToArray()));
session.SendMessage(reply);
break;
}
};
using var client = new NetClient($"tcp://127.0.0.1:{server.Port}") { Protocol = new SrmpCodec() };
Assert.True(client.Open());
// 多路复用:并发发出慢请求(Seq 0x51)与快请求(Seq 0x52)
var slow = client.SendMessageAsync(new DefaultMessage { Sequence = 0x51 }).AsTask();
await WithTimeout(slowIn.Task);
var fast = client.SendMessageAsync(new DefaultMessage { Sequence = 0x52 }).AsTask();
// 并行证据:慢请求仍被阻塞时,快请求先完成;串行实现下 3 秒内必超时
var fastResp = await WithTimeout(fast, 3_000);
Assert.Equal("fast", fastResp.Payload!.ToStr());
// 放行慢请求,双请求均完成且各自配对
slowGate.Set();
var slowResp = await WithTimeout(slow, 5_000);
Assert.Equal("slow", slowResp.Payload!.ToStr());
var slowMsg = Assert.IsType<DefaultMessage>(slowResp);
Assert.Equal(MessageKinds.Response, slowMsg.Kind);
Assert.Equal(0x51, slowMsg.Sequence);
}
#endregion
}
#pragma warning restore xUnit1031
|