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

namespace XUnitTest.Net;

/// <summary>发送背压(TrySend 暂停拒绝 / SendAsync(IPacket) 挂起等待)测试</summary>
[Collection("Net")]
public class SendBackpressureTests
{
    [Fact]
    [DisplayName("背压发送_TrySend与SendAsync_暂停拒绝_恢复完成")]
    public async Task Pause_RejectsTrySend_SuspendsSendAsync()
    {
        // 裸监听套接字:接受连接但不读取,使内核缓冲写满后发送泵阻塞、管道积压达到暂停水位
        var listener = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
        listener.Bind(new IPEndPoint(IPAddress.Loopback, 0));
        listener.Listen(1);
        var port = ((IPEndPoint)listener.LocalEndPoint!).Port;

        using var client = new TcpSession { Remote = new NetUri($"tcp://127.0.0.1:{port}") };
        try
        {
            client.Open();

            // 收缩客户端发送缓冲,加快内核缓冲饱和
            client.Client!.SendBufferSize = 16 * 1024;

            var pipe = client.SendPipe;
            pipe.PauseThreshold = 32 * 1024;
            pipe.ResumeThreshold = 16 * 1024;

            // 推入数据:泵写入内核的部分有限,应用层管道积压很快达到暂停水位(同步 Send 走管道为入队,不阻塞)
            var payload = new Byte[64 * 1024];
            var sw = Stopwatch.StartNew();
            while (!pipe.IsPaused && sw.Elapsed < TimeSpan.FromSeconds(10))
            {
                client.Send(payload);
                await Task.Delay(5);
            }
            Assert.True(pipe.IsPaused, "未达到暂停水位");

            // 超量灌注,确保积压远超恢复水位,对端不读时泵无法消化
            for (var i = 0; i < 8; i++) client.Send(payload);

            // 暂停水位:TrySend 拒绝;SendAsync 挂起等待
            Assert.False(client.TrySend(new ArrayPacket(new Byte[128])));

            var sending = client.SendAsync(new ArrayPacket(new Byte[128]));
            await Task.Delay(200);
            Assert.False(sending.IsCompleted, "暂停水位下 SendAsync 不应完成");

            // 对端开始读取:积压消化到恢复水位以下,挂起提交被唤醒
            var conn = listener.Accept();
            _ = Task.Run(() =>
            {
                try
                {
                    var buf = new Byte[64 * 1024];
                    var total = 0L;
                    while (total < 4L * 1024 * 1024)
                    {
                        var n = conn.Receive(buf);
                        if (n <= 0) break;

                        total += n;
                    }
                }
                catch
                {
                    // 测试收尾关闭连接时正常退出
                }
            });

            var rs = await sending.AsTask().WaitAsync(TimeSpan.FromSeconds(10));
            Assert.Equal(128, rs);
            Assert.False(pipe.IsPaused);

            conn.Dispose();
        }
        finally
        {
            listener.Dispose();
        }
    }

    [Fact]
    [DisplayName("背压发送_正常路径_TrySend与SendAsync送达")]
    public async Task Normal_SendDeliversInOrder()
    {
        using var server = new NetServer { Port = 0 };
        server.Start();

        var wait = new ManualResetEventSlim();
        var received = new List<Byte>();
        server.NewSession += (s, e) =>
        {
            e.Session.Session.Received += (ss, ee) =>
            {
                var pk = ee.Packet;
                if (pk != null) lock (received) received.AddRange(pk.ToArray());
            };
            wait.Set();
        };

        using var client = new TcpSession { Remote = new NetUri($"tcp://127.0.0.1:{server.Port}") };
        client.Open();

        // 会话建立且订阅就绪后再发送,避免首轮数据早于订阅
        Assert.True(wait.Wait(3_000));

        // 正常水位:TrySend 接受
        Assert.True(client.TrySend(new ArrayPacket(new Byte[] { 1, 2, 3, 4 })));

        // SendAsync(IPacket):未暂停时立即完成
        var rs = await client.SendAsync(new ArrayPacket(new Byte[] { 5, 6, 7 }));
        Assert.Equal(3, rs);

        // 送达且保序
        var sw = Stopwatch.StartNew();
        while (sw.Elapsed < TimeSpan.FromSeconds(5))
        {
            lock (received) { if (received.Count >= 7) break; }
            await Task.Delay(10);
        }

        lock (received) Assert.Equal(new Byte[] { 1, 2, 3, 4, 5, 6, 7 }, received);
    }

    [Fact]
    [DisplayName("背压发送_挂起期间关闭_发送被唤醒而不挂死")]
    public async Task Suspend_ThenClose_WakesUp()
    {
        // 裸监听套接字:接受连接但不读取,使发送泵阻塞、管道积压达到暂停水位
        var listener = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
        listener.Bind(new IPEndPoint(IPAddress.Loopback, 0));
        listener.Listen(1);
        var port = ((IPEndPoint)listener.LocalEndPoint!).Port;

        try
        {
            using var client = new TcpSession { Remote = new NetUri($"tcp://127.0.0.1:{port}") };
            client.Open();
            client.Client!.SendBufferSize = 16 * 1024;

            var pipe = client.SendPipe;
            pipe.PauseThreshold = 32 * 1024;
            pipe.ResumeThreshold = 16 * 1024;

            // 灌入数据直到达到暂停水位
            var payload = new Byte[64 * 1024];
            var sw = Stopwatch.StartNew();
            while (!pipe.IsPaused && sw.Elapsed < TimeSpan.FromSeconds(10))
            {
                client.Send(payload);
                await Task.Delay(5);
            }
            Assert.True(pipe.IsPaused, "未达到暂停水位");
            for (var i = 0; i < 8; i++) client.Send(payload);

            // 暂停水位下 SendAsync 挂起
            var sending = client.SendAsync(new ArrayPacket(new Byte[128]));
            await Task.Delay(200);
            Assert.False(sending.IsCompleted, "暂停水位下 SendAsync 不应完成");

            // 挂起期间关闭会话:关闭不应被背压阻塞,挂起发送应被唤醒并失败
            var closeTask = Task.Run(() => client.Close("close-while-paused"));
            var closeCompleted = await Task.WhenAny(closeTask, Task.Delay(8_000)) == closeTask;
            Assert.True(closeCompleted, "关闭会话不应因背压挂起而被阻塞");

            var task = sending.AsTask();
            var sendCompleted = await Task.WhenAny(task, Task.Delay(5_000)) == task;
            Assert.True(sendCompleted, "关闭会话后挂起的 SendAsync 应被唤醒,而不是永久挂起");

            try
            {
                var rs = await task;
                Assert.True(rs <= 0, $"关闭后发送不应报告成功:{rs}");
            }
            catch
            {
                // 明确异常亦可:挂起者被唤醒并失败即达成目标
            }
        }
        finally
        {
            listener.Dispose();
        }
    }
}