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

namespace XUnitTest.Net;

/// <summary>UDP 协议模式(数据报定界 + 消息分发)实网测试</summary>
[Collection("Net")]
public class UdpMessageSessionTests
{
    #region 宿主
    /// <summary>UDP 消息宿主:收集收到的行内容</summary>
    public class UdpMessageServer : NetServer<UdpMessageSession>
    {
        /// <summary>收到的消息内容,线程安全</summary>
        public ConcurrentQueue<String> ReceivedLines { get; } = new();
    }

    /// <summary>UDP 消息会话:读出行内容并回显</summary>
    public class UdpMessageSession : NetSession<UdpMessageServer>
    {
        protected override void OnReceive(ReceivedEventArgs e)
        {
            if (e.Message is not Message msg) return;

            var line = ReadBody(msg);
            if (line == null) return;

            Host.ReceivedLines.Enqueue(line);

            // 回显:协议构建后经底层会话发送
            var reply = new Message();
            reply.SetBody(new ArrayPacket(Encoding.UTF8.GetBytes("echo:" + line)));
            (Session as UdpSession)?.SendMessage(reply);
        }

        internal static String? ReadBody(Message msg)
        {
            var body = msg.Body;
            if (body == null) return null;

            // UDP 消息体为内存模式,读满立即完成
            var data = body.ReadAllAsync().AsTask().GetAwaiter().GetResult();
            var line = Encoding.UTF8.GetString(data.AsReadOnlySequence().ToArray());
            data.TryDispose();

            return line;
        }
    }
    #endregion

    #region 工具
    private static UdpMessageServer NewServer()
    {
        var server = new UdpMessageServer
        {
            Port = 0,
            ProtocolType = NetType.Udp,
            AddressFamily = AddressFamily.InterNetwork,
        };
        server.Protocol = new SplitDataCodec();
        server.Start();

        return server;
    }

    private static NetClient NewClient(UdpMessageServer server)
    {
        var client = new NetClient($"udp://127.0.0.1:{server.Port}")
        {
            Protocol = new SplitDataCodec(),
            AutoReconnect = false,
        };
        client.Open();

        return client;
    }

    private static async Task<Boolean> WaitUntilAsync(Func<Boolean> condition, Int32 timeoutMs = 5_000)
    {
        var start = Environment.TickCount64;
        while (Environment.TickCount64 - start < timeoutMs)
        {
            if (condition()) return true;

            await Task.Delay(20);
        }

        return condition();
    }
    #endregion

    [Fact]
    [DisplayName("UDP协议_客户端到服务端_行协议消息")]
    public async Task ClientToServer()
    {
        using var server = NewServer();
        using var client = NewClient(server);

        var msg = new Message();
        msg.SetBody(new ArrayPacket(Encoding.UTF8.GetBytes("hello")));
        client.SendMessage(msg);

        Assert.True(await WaitUntilAsync(() => server.ReceivedLines.Contains("hello")));
    }

    [Fact]
    [DisplayName("UDP协议_服务端回显_客户端收到消息")]
    public async Task ServerToClient()
    {
        using var server = NewServer();
        using var client = NewClient(server);

        var received = new TaskCompletionSource<String>();
        client.Received += (s, e) =>
        {
            if (e.Message is not Message msg) return;

            var line = UdpMessageSession.ReadBody(msg);
            if (line != null) received.TrySetResult(line);
        };

        var msg = new Message();
        msg.SetBody(new ArrayPacket(Encoding.UTF8.GetBytes("ping")));
        client.SendMessage(msg);

        var task = await Task.WhenAny(received.Task, Task.Delay(5_000));
        Assert.True(task == received.Task, "客户端应在超时内收到回显消息");
        Assert.Equal("echo:ping", await received.Task);
    }

    [Fact]
    [DisplayName("UDP协议_单数据报多帧_逐条分发")]
    public async Task MultiFramesInOneDatagram()
    {
        using var server = NewServer();
        using var client = NewClient(server);

        // 原始字节发送:一个数据报内两帧行协议
        client.Send(Encoding.UTF8.GetBytes("one\r\ntwo\r\n"));

        Assert.True(await WaitUntilAsync(() => server.ReceivedLines.Contains("one") && server.ReceivedLines.Contains("two")));
    }

    [Fact]
    [DisplayName("UDP协议_请求响应等待_明确拒绝")]
    public async Task RequestResponse_NotSupported()
    {
        using var server = NewServer();
        using var client = NewClient(server);

        var msg = new Message();
        msg.SetBody(new ArrayPacket(Encoding.UTF8.GetBytes("rpc")));

        await Assert.ThrowsAsync<NotSupportedException>(() => client.SendMessageAsync(msg).AsTask());
    }

    [Fact]
    [DisplayName("UDP协议_压缩协议_解压结果不被线上压缩字节覆盖")]
    public async Task CompressedCodec_BodyNotOverwritten()
    {
        // 压缩协议要求整帧(UDP 数据报自带完整帧),可验证“协议预绑定体不被二次绑定覆盖”
        using var server = new UdpMessageServer
        {
            Port = 0,
            ProtocolType = NetType.Udp,
            AddressFamily = AddressFamily.InterNetwork,
        };
        server.Protocol = new CompressedCodec(new SrmpCodec());
        server.Start();

        using var client = new NetClient($"udp://127.0.0.1:{server.Port}")
        {
            Protocol = new CompressedCodec(new SrmpCodec()),
            AutoReconnect = false,
        };
        client.Open();

        var msg = new DefaultMessage { Sequence = 1 };
        msg.SetBody(new ArrayPacket(Encoding.UTF8.GetBytes("compressed-body")));
        client.SendMessage(msg);

        // 服务端应读到解压后的原文;
        // 修复前会被线上压缩字节覆盖(乱码),Contains 永不成立
        Assert.True(await WaitUntilAsync(() => server.ReceivedLines.Contains("compressed-body")));
    }
}