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

namespace NewLife.Messaging;

/// <summary>压缩消息编解码器。装饰内层协议,只对消息负载做 Deflate 压缩/解压(帧头与帧定界仍由内层协议明文封装)</summary>
/// <remarks>
/// <para>组合示例:<c>server.Protocol = new CompressedCodec(new SrmpCodec());</c>——同一 <see cref="IMessageCodec"/> 接口可任意嵌套组合(压缩/加密等变换层)。</para>
/// <para><b>连接级约定</b>:收发两端须同时配置本编码器,线格式不含压缩标志。</para>
/// <para><b>接收</b>:要求完整帧到齐后整载解压(压缩体无法流式解压),解压后经 <see cref="Message.SetBody"/> 预绑定,帧泵检测到已绑定体时直接消费整帧;坏数据抛出异常交由帧泵关闭会话。</para>
/// <para><b>发送</b>:整帧构建时压缩内存体再交内层协议;空体不压缩。流式发送不支持(<see cref="BuildHeader"/> 抛异常)。</para>
/// </remarks>
public class CompressedCodec(IMessageCodec inner) : IMessageCodec, IMessageCodecDecorator
{
    /// <summary>内层协议</summary>
    public IMessageCodec Inner { get; } = inner;

    /// <summary>解压后最大字节数,默认 16M。0 表示不限制</summary>
    /// <remarks>
    /// <para>线上帧长上限(如 <see cref="MessagePump.MaxFrameSize"/>)约束的是压缩后大小,解压后大小完全由对端控制(Deflate 可达 1000:1)。
    /// 无上限时对端只需发一个 1KB 的帧,即可让接收方分配 GB 级内存(压缩炸弹),且解压发生在接收路径同步执行。</para>
    /// <para>超过上限时抛出异常,与坏帧同样交由帧泵关闭会话。</para>
    /// </remarks>
    public Int32 MaxDecompressedSize { get; set; } = 16 * 1024 * 1024;

    /// <summary>定界并构造消息。内层定界后,完整帧解压消息体并预绑定;帧未完整返回 null 等更多数据</summary>
    /// <param name="buffer">帧首窗口(只读序列,可跨段)</param>
    /// <returns>解析结果;头部不足或帧未完整返回 null</returns>
    public ParseResult? TryParse(ReadOnlySequence<Byte> buffer)
    {
        var rs = Inner.TryParse(buffer);
        if (rs == null) return null;

        var r = rs.Value;
        var msg = r.Message;
        if (msg == null) return rs;   // 无消息帧直通

        // 压缩体必须整载解压:帧未完整时声明“已定界待整帧”,让帧泵保留窗口等更多数据(不消费、也不当成无法定界的残余)
        if (r.HeaderSize + r.BodyLength > buffer.Length)
        {
            msg.Dispose();
            return new ParseResult { NeedFullFrame = true, HeaderSize = r.HeaderSize, BodyLength = r.BodyLength };
        }

        if (r.BodyLength > 0)
        {
            // 压缩体整载取出:单段直接引用帧窗口(零拷贝),跨段才拼一份;解压结果窃取内存流缓冲
            var body = buffer.Slice(r.HeaderSize, r.BodyLength);
            using var input = body.IsSingleSegment && MemoryMarshal.TryGetArray(body.First, out var segment)
                ? new MemoryStream(segment.Array!, segment.Offset, segment.Count, false)
                : new MemoryStream(body.ToArray());

            var ms = new MemoryStream();
            try
            {
                using (var inflate = new DeflateStream(input, CompressionMode.Decompress, true)) CopyToBounded(inflate, ms, MaxDecompressedSize);
            }
            catch
            {
                // 解压失败(含超限):释放已构造的消息,避免泄漏
                msg.Dispose();
                throw;
            }

            ms.Position = 0;
            msg.SetBody(new ArrayPacket(ms));
        }

        return r;
    }

    /// <summary>带长度上限的流拷贝。超过上限抛异常,防御解压炸弹</summary>
    /// <param name="source">源流</param>
    /// <param name="destination">目标流</param>
    /// <param name="maxSize">最大字节数,0 表示不限制</param>
    private static void CopyToBounded(Stream source, Stream destination, Int32 maxSize)
    {
        var buf = ArrayPool<Byte>.Shared.Rent(8192);
        try
        {
            var total = 0;
            while (true)
            {
                var n = source.Read(buf, 0, buf.Length);
                if (n <= 0) break;

                total += n;
                if (maxSize > 0 && total > maxSize) throw new InvalidDataException($"解压后长度超过上限 {maxSize},拒绝继续解压(疑似压缩炸弹)");

                destination.Write(buf, 0, n);
            }
        }
        finally
        {
            ArrayPool<Byte>.Shared.Return(buf);
        }
    }

    /// <summary>整帧构建。压缩消息体后交内层协议构建,构建后还原消息负载</summary>
    /// <param name="message">消息</param>
    /// <returns>整帧拥有句柄,调用方负责 Dispose</returns>
    public IOwnerPacket? Build(IMessage message)
    {
        var body = message.Payload;
        if (body == null || body.Total == 0) return Inner.Build(message);

        // 逐段写入压缩流,链式负载不必先聚合;压缩结果窃取内存流缓冲
        var ms = new MemoryStream();
        using (var deflate = new DeflateStream(ms, CompressionLevel.Optimal, true)) body.CopyTo(deflate);

        ms.Position = 0;

        // 内层协议从 Payload 取体:先摘除原负载(放弃持有、不归还)再顶上压缩体,构建后原样换回
        message.SetBody(null);
        message.SetBody(new ArrayPacket(ms));

        try
        {
            return Inner.Build(message);
        }
        finally
        {
            message.SetBody(body);
        }
    }

    /// <summary>仅构建头部。压缩体长度无法预声明,不支持流式发送</summary>
    /// <param name="message">消息</param>
    /// <param name="bodyLength">消息体字节数</param>
    /// <returns>头部数据包</returns>
    /// <exception cref="NotSupportedException">压缩协议无法预声明压缩后长度</exception>
    public IOwnerPacket BuildHeader(IMessage message, Int64 bodyLength) => throw new NotSupportedException("压缩协议无法预声明压缩后长度,请使用整帧构建(Build)");
}