解决MySql布尔型新旧版本兼容问题,采用枚举来表示布尔型的数据表。由正向工程赋值
大石头 authored at 2018-05-15 21:21:05
10.92 KiB
X
using System.Buffers;
using System.ComponentModel;
using System.Diagnostics;
using System.Net;
using System.Net.Sockets;
using NewLife.Data;
using NewLife.Messaging;
using NewLife.Net;
using Xunit;

namespace XUnitTest.Net;

/// <summary>TcpSession 流式收发(出站单出口/流式发送/关闭排空)回环测试</summary>
[Collection("Net.A")]
public class TcpSessionStreamTests
{
    #region 工具
    private static async Task WaitUntilAsync(Func<Boolean> condition, Int32 timeoutMs = 10_000)
    {
        var sw = Stopwatch.StartNew();
        while (!condition())
        {
            if (sw.ElapsedMilliseconds > timeoutMs) throw new TimeoutException("等待条件超时");

            await Task.Delay(10);
        }
    }

    /// <summary>建立回环连接并等待服务端会话</summary>
    private static async Task<(NetServer Server, TcpClient Client, TcpSession Session)> ConnectAsync()
    {
        var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp };
        server.Start();

        var wait = new ManualResetEventSlim();
        TcpSession? session = null;
        server.NewSession += (s, e) =>
        {
            session = e.Session.Session as TcpSession;
            wait.Set();
        };

        var client = new TcpClient();
        await client.ConnectAsync(IPAddress.Loopback, server.Port);
        if (!wait.Wait(3_000) || session == null)
        {
            client.Dispose();
            server.Dispose();

            throw new TimeoutException("等待服务端会话超时");
        }

        return (server, client, session);
    }

    /// <summary>并发排空:边发边收,不依赖内核缓冲容量</summary>
    private static Task<Int32> DrainAsync(TcpClient client, Byte[] received)
    {
        var stream = client.GetStream();
        stream.ReadTimeout = 30_000;

        return Task.Run(() =>
        {
            var n = 0;
            while (n < received.Length)
            {
                var c = stream.Read(received, n, received.Length - n);
                if (c <= 0) break;

                n += c;
            }

            return n;
        });
    }
    #endregion

    #region 会话级行为
    [Fact]
    [DisplayName("直发_Send 同步写出_返回后改写原缓冲不影响已发数据")]
    public async Task DirectSend_WritesCallerBuffer()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            var src = new Byte[] { 10, 20, 30 };
            var rs = session.Send(new ArrayPacket(src));
            Assert.Equal(3, rs);

            // 直发是同步写出:返回即已写入内核,此后改写原缓冲不影响已发数据
            src[0] = 99;

            var received = new Byte[3];
            Assert.Equal(3, await DrainAsync(client, received));
            Assert.Equal(new Byte[] { 10, 20, 30 }, received);
        }
    }

    [Fact]
    [DisplayName("直发_发送后立即关闭_已写数据不丢")]
    public async Task DirectSend_ThenClose_DataNotLost()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            var p1 = new Byte[] { 1, 2, 3 };
            var p2 = new Byte[] { 4, 5 };
            var received = new Byte[p1.Length + p2.Length];
            var reader = DrainAsync(client, received);

            // 直发:两次发送返回时数据已在网络栈里,随后关闭不会丢数据
            Assert.Equal(3, session.Send(new ArrayPacket(p1)));
            Assert.Equal(2, session.Send(new ArrayPacket(p2)));

            var ok = session.Close("test");
            Assert.True(ok);

            Assert.Equal(received.Length, await reader.WaitAsync(TimeSpan.FromSeconds(5)));
            Assert.Equal(p1.Concat(p2).ToArray(), received);

            // 服务端会话关闭即销毁:后续发送抛 ObjectDisposedException
            Assert.True(session.Disposed);
            Assert.Throws<ObjectDisposedException>(() => session.Send(new ArrayPacket(new Byte[] { 6 })));
        }
    }

    [Fact]
    [DisplayName("流式会话_销毁_入站管道残余归还")]
    public async Task InboundPipe_Close_ReleasesResidue()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            // 触发入站管道创建:此后接收数据同时投递到管道(共享切片)
            var pipe = session.Pipe;

            await client.GetStream().WriteAsync(new Byte[512]);
            await WaitUntilAsync(() => pipe.UnconsumedLength > 0);

            // 关闭并销毁会话:入站管道残余必须归还。
            // 旧实现只 _pipe.Writer.Complete(),读侧链上的池缓冲一挂到会话对象不可达,由终结器兜底
            session.Close("test");

            Assert.Equal(0, pipe.UnconsumedLength);
        }
    }
    #endregion

    #region 回环
    [Fact]
    [DisplayName("直发_回环_批量下发完整到达")]
    public async Task Loopback_BatchSend_Arrives()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            var payload = new Byte[100_000];
            Random.Shared.NextBytes(payload);

            // 对端并发排空:边发边收,不依赖内核缓冲容量
            var received = new Byte[payload.Length];
            var reader = DrainAsync(client, received);

            var rs = session.Send(payload);
            Assert.Equal(payload.Length, rs);

            // 对端完整收到且内容一致
            var read = await reader.WaitAsync(TimeSpan.FromSeconds(30));
            Assert.Equal(payload.Length, read);
            Assert.Equal(payload, received);
        }
    }

    [Fact]
    [DisplayName("流式发送_回环_大文件完整到达")]
    public async Task SendStream_Loopback_LargeArrives()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            // 1MB 数据流式发送:分块读入一块、写完再读下一块,内存恒为一块
            var payload = new Byte[1_000_000];
            Random.Shared.NextBytes(payload);

            // 对端并发排空:发送受阻于内核缓冲容量时,若客户端不读将永远阻塞(并行负载下必须边发边收)
            var received = new Byte[payload.Length];
            var reader = DrainAsync(client, received);

            var total = await session.SendAsync(new MemoryStream(payload));
            Assert.Equal(payload.Length, total);

            // 对端完整收到且内容一致
            var read = await reader.WaitAsync(TimeSpan.FromSeconds(30));
            Assert.Equal(payload.Length, read);
            Assert.Equal(payload, received);
        }
    }

    [Fact]
    [DisplayName("流式发送_SRMP头+流式体_单条消息完整到达")]
    public async Task SendStream_SrmpHeaderBody_OneMessage()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            // 发送侧:头声明长度 + 流式体(等价 SRMP 大消息流式发送);同一把写锁保证头与体无交错
            var payload = new Byte[256 * 1024];
            Random.Shared.NextBytes(payload);

            // 对端并发排空:边发边收,不依赖内核缓冲容量
            var frameLen = 8 + payload.Length;
            var received = new Byte[frameLen];
            var reader = DrainAsync(client, received);

            var msg = new DefaultMessage { Sequence = 0x3C };
            var codec = new SrmpCodec();
            var rs = session.Send(codec.BuildHeader(msg, payload.Length));
            Assert.Equal(8, rs);

            var total = await session.SendAsync(new MemoryStream(payload), payload.Length);
            Assert.Equal(payload.Length, total);

            // 接收侧:按消息自解析定界,完整负载与发送一致
            var read = await reader.WaitAsync(TimeSpan.FromSeconds(30));
            Assert.Equal(frameLen, read);

            var pr = codec.TryParse(new ArrayPacket(received).AsReadOnlySequence());
            Assert.NotNull(pr);
            Assert.Equal(8, pr.Value.HeaderSize);
            Assert.Equal(payload.Length, pr.Value.BodyLength);
            Assert.Equal(0x3C, ((DefaultMessage)pr.Value.Message!).Sequence);
            Assert.Equal(payload, received[8..]);
        }
    }

    [Fact]
    [DisplayName("流式发送_调用方取消_抛取消异常且不关闭会话")]
    public async Task SendStream_CallerCancel_ThrowsAndKeepsSession()
    {
        var (server, client, session) = await ConnectAsync();
        using (server)
        using (client)
        {
            // 读块挂起直到取消触发:取消点确定,不依赖内核缓冲容量或对端读速
            using var source = new BlockingStream();
            using var cts = new CancellationTokenSource();

            var task = session.SendAsync(source, 1024, cts.Token).AsTask();
            Assert.True(source.Entered.Wait(5_000), "未进入读块,用例前置条件不成立");

            cts.Cancel();

            // 调用方取消不是发送故障:向上抛取消异常
            await Assert.ThrowsAnyAsync<OperationCanceledException>(() => task);

            // 会话不得被当作发送故障关闭(旧实现走通用 catch 上报并 Close("SendError"))
            Assert.True(session.Active, "调用方取消被误判为发送故障,会话已被关闭");
        }
    }

    /// <summary>读块挂起直到被取消的测试流</summary>
    private sealed class BlockingStream : Stream
    {
        private readonly TaskCompletionSource<Boolean> _entered = new(TaskCreationOptions.RunContinuationsAsynchronously);

        /// <summary>已进入读块</summary>
        public Task<Boolean> Entered => _entered.Task;

        public override async Task<Int32> ReadAsync(Byte[] buffer, Int32 offset, Int32 count, CancellationToken cancellationToken)
        {
            _entered.TrySetResult(true);

            await Task.Delay(System.Threading.Timeout.Infinite, cancellationToken).ConfigureAwait(false);

            return 0;
        }

        public override Boolean CanRead => true;
        public override Boolean CanSeek => false;
        public override Boolean CanWrite => false;
        public override Int64 Length => throw new NotSupportedException();
        public override Int64 Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); }
        public override void Flush() { }
        public override Int32 Read(Byte[] buffer, Int32 offset, Int32 count) => throw new NotSupportedException();
        public override Int64 Seek(Int64 offset, SeekOrigin origin) => throw new NotSupportedException();
        public override void SetLength(Int64 value) => throw new NotSupportedException();
        public override void Write(Byte[] buffer, Int32 offset, Int32 count) => throw new NotSupportedException();
    }
    #endregion
}