解决MySql布尔型新旧版本兼容问题,采用枚举来表示布尔型的数据表。由正向工程赋值
大石头 authored at 2018-05-15 21:21:05
11.87 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>发送出口测试:直发(无发送队列)——同步写完才返回,多入口共用一把写锁不交错</summary>
/// <remarks>
/// <para>架构:所有发送入口(Send/SendAsync/SendAsync(Stream)/SendFileAsync)共用一把写锁,任一时刻只有一位写者。</para>
/// <para>观测手段:对端不读且收缩收发缓冲时,直发会堵在内核写阻塞上(同步 Send 不返回);
/// 对端开始读后数据完整到达。并发用例用“块内字节一致”检出交错。</para>
/// </remarks>
[Collection("Net.D")]
public class SendEntryTests
{
    #region 工具
    /// <summary>指定时间内是否已完成</summary>
    private static async Task<Boolean> CompletedWithinAsync(Task task, Int32 timeoutMs)
    {
        var finished = await Task.WhenAny(task, Task.Delay(timeoutMs));
        return ReferenceEquals(finished, task);
    }

    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);
        }
    }
    #endregion

    [Fact]
    [DisplayName("直发_Send 同步写出_慢对端阻塞_读后续达")]
    public async Task Send_SyncWriteBlocksOnSlowPeer()
    {
        const Int32 chunkSize = 64 * 1024;
        const Int32 chunkCount = 64;                        // 合计 4MB,超过内核收发缓冲
        var body = new Byte[chunkSize];
        body.AsSpan().Fill(0x5A);

        // 裸监听套接字:接受连接但不读取;收缩收发缓冲,让直发必然堵在内核写阻塞上
        var listener = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
        listener.Bind(new IPEndPoint(IPAddress.Loopback, 0));
        listener.ReceiveBufferSize = 8 * 1024;
        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;

            using var peer = await listener.AcceptWithinAsync();
            peer.ReceiveBufferSize = 8 * 1024;
            peer.ReceiveTimeout = 30_000;

            // 直发:对端不读时同步 Send 必然阻塞(内存由内核缓冲天然限住,不需要发送队列)
            var sending = Task.Factory.StartNew(() =>
            {
                var total = 0;
                for (var i = 0; i < chunkCount; i++)
                {
                    var rs = client.Send(body, 0, chunkSize);
                    if (rs <= 0) break;

                    total += rs;
                }

                return total;
            }, TaskCreationOptions.LongRunning);
            Assert.False(await CompletedWithinAsync(sending, 300), "前置条件不成立:直发应被内核缓冲阻塞");

            // 对端开始读:直发跑到写完,数据完整到达
            var received = new MemoryStream();
            var reading = Task.Factory.StartNew(() =>
            {
                var buf = new Byte[64 * 1024];
                while (received.Length < chunkSize * chunkCount)
                {
                    var n = peer.Receive(buf);
                    if (n <= 0) break;

                    received.Write(buf, 0, n);
                }
            }, TaskCreationOptions.LongRunning);

            Assert.Equal(chunkSize * chunkCount, await sending.WaitAsync(TimeSpan.FromSeconds(30)));
            await reading.WaitAsync(TimeSpan.FromSeconds(30));

            var bytes = received.ToArray();
            Assert.Equal(chunkSize * chunkCount, bytes.Length);
            Assert.DoesNotContain(bytes, b => b != 0x5A);
        }
        finally
        {
            listener.Dispose();
        }
    }

    [Fact]
    [DisplayName("直发_并发写入_各块不交错(写锁契约)")]
    public async Task Send_ConcurrentWriters_NoInterleave()
    {
        const Int32 writers = 8;
        const Int32 blocks = 200;
        const Int32 blockSize = 512;

        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();

            using var peer = await listener.AcceptWithinAsync();
            peer.ReceiveTimeout = 30_000;

            var total = writers * blocks * blockSize;
            var received = new MemoryStream();
            var reading = Task.Factory.StartNew(() =>
            {
                var buf = new Byte[64 * 1024];
                while (received.Length < total)
                {
                    var n = peer.Receive(buf);
                    if (n <= 0) break;

                    received.Write(buf, 0, n);
                }
            }, TaskCreationOptions.LongRunning);

            // 多线程同时发送:写锁保证任一时刻只有一位写者,所以每个 512 字节块内部必然同字节
            var tasks = new Task[writers];
            for (var w = 0; w < writers; w++)
            {
                var id = (Byte)(w + 1);
                var block = new Byte[blockSize];
                block.AsSpan().Fill(id);

                tasks[w] = Task.Factory.StartNew(() =>
                {
                    for (var i = 0; i < blocks; i++) client.Send(block, 0, blockSize);
                }, TaskCreationOptions.LongRunning);
            }

            await Task.WhenAll(tasks).WaitAsync(TimeSpan.FromSeconds(30));
            await reading.WaitAsync(TimeSpan.FromSeconds(30));

            var bytes = received.ToArray();
            Assert.Equal(total, bytes.Length);

            // 交错会把别的写者的字节混进块内(块内出现两种字节)
            for (var i = 0; i < total; i += blockSize)
            {
                var first = bytes[i];
                for (var j = 1; j < blockSize; j++)
                {
                    Assert.Equal(first, bytes[i + j]);
                }
            }
        }
        finally
        {
            listener.Dispose();
        }
    }

    [Fact]
    [DisplayName("直发_Send与SendAsync保序送达")]
    public async Task SendAndSendAsync_DeliverInOrder()
    {
        using var server = new NetServer { Port = 0, ProtocolType = NetType.Tcp };
        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.GetSpan().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));

        // 同步直发与异步直发交替,共用写锁 → 前缀严格保序
        Assert.Equal(4, client.Send(new Byte[] { 1, 2, 3, 4 }));
        Assert.Equal(3, await client.SendAsync(new ArrayPacket(new Byte[] { 5, 6, 7 })));
        Assert.Equal(2, client.Send(new Byte[] { 8, 9 }));

        var sw = Stopwatch.StartNew();
        while (sw.Elapsed < TimeSpan.FromSeconds(5))
        {
            lock (received) { if (received.Count >= 9) break; }
            await Task.Delay(10);
        }

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

    [Fact]
    [DisplayName("文件发送_非SSL_内核零拷贝_内容逐字节一致")]
    public async Task SendFile_ZeroCopy_ContentMatches()
    {
        const Int32 size = 300 * 1024;
        var payload = new Byte[size];
        Random.Shared.NextBytes(payload);

        var file = Path.Combine(Path.GetTempPath(), "nl_sendfile_" + Guid.NewGuid().ToString("N") + ".bin");
        await File.WriteAllBytesAsync(file, payload);

        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();

            using var peer = await listener.AcceptWithinAsync();
            peer.ReceiveTimeout = 30_000;

            var received = new MemoryStream();
            var reading = Task.Factory.StartNew(() =>
            {
                var buf = new Byte[64 * 1024];
                while (received.Length < size)
                {
                    var n = peer.Receive(buf);
                    if (n <= 0) break;

                    received.Write(buf, 0, n);
                }
            }, TaskCreationOptions.LongRunning);

            // 非 SSL会话:整个文件交给内核推送(sendfile/TransmitFile),应用层不经读块
            var rs = await client.SendFileAsync(file);
            Assert.Equal(size, rs);

            await reading.WaitAsync(TimeSpan.FromSeconds(30));
            Assert.Equal(payload, received.ToArray());
        }
        finally
        {
            listener.Dispose();
            File.Delete(file);
        }
    }

    [Fact]
    [DisplayName("发送失败_错误回调内再次 Send_不自死锁")]
    public async Task SendFailed_ErrorCallbackSendsAgain_NoDeadlock()
    {
        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();

            using var peer = await listener.AcceptWithinAsync();
            peer.ReceiveTimeout = 30_000;

            // 确定性制造发送失败:关掉发送方向后任何 Send 都立刻失败,
            // 而套接字仍处已绑定(Open 仍为真)、接收方向不受影响,不会顺带把会话关掉
            client.Client!.Shutdown(SocketShutdown.Send);

            // 错误回调内同步再次 Send:上报若发生在写锁内,这里会二次等待同一把 SemaphoreSlim 而自死锁
            var innerDone = new ManualResetEventSlim();
            var fired = 0;
            client.Error += (s, e) =>
            {
                // 只在内层发一次,避免“失败 → 上报 → 再失败”的无限递归
                if (Interlocked.Increment(ref fired) > 1) return;

                client.Send(new Byte[] { 9, 9 });
                innerDone.Set();
            };

            var rs = 0;
            Exception? error = null;
            var sending = Task.Factory.StartNew(() =>
            {
                try { rs = client.SendAsync(new ArrayPacket(new Byte[64])).AsTask().GetAwaiter().GetResult(); }
                catch (Exception ex) { error = ex; }
            }, TaskCreationOptions.LongRunning);

            Assert.True(await CompletedWithinAsync(sending, 10_000), "发送失败时,错误回调内再次 Send 发生了自死锁");
            Assert.True(innerDone.IsSet, $"发送失败没有触发错误回调(rs={rs}, error={error?.GetType().Name}: {error?.Message})");
        }
        finally
        {
            listener.Dispose();
        }
    }
}