解决MySql布尔型新旧版本兼容问题,采用枚举来表示布尔型的数据表。由正向工程赋值
大石头 authored at 2018-05-15 21:21:05
32.29 KiB
X
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