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

namespace NewLife.Net;

/// <summary>增强TCP客户端(协议模式的消息泵与帧定界交付)</summary>
/// <remarks>
/// <para>协议模式下数据经数据管道定界,由消息泵逐帧交付。流式专有配置(最大残留、单帧上限、整帧模式)与并发派发信号量随本类,
/// 不再放到 <see cref="SessionBase"/>:数据报会话没有粘包与定界问题,这些成员在它那里只会是"未使用"。</para>
/// </remarks>
public partial class TcpSession
{
    #region 属性

    /// <summary>消息泵最大缓存字节数(协议模式下无法定界的残余上限),默认 1M。0 表示不限制</summary>
    /// <remarks>残余达到上限说明对端数据与协议不匹配或已损坏,消息泵随即报错并关闭会话,避免连接僵死</remarks>
    public Int32 MaxCache { get; set; } = 1024 * 1024;

    /// <summary>整帧模式下的单帧长度上限,默认 16M。0 表示不限制</summary>
    /// <remarks>仅 <see cref="RequireFullFrame"/> 为 true 时生效:整帧解析要求整个帧驻留内存,本上限是单连接的内存安全阀</remarks>
    public Int32 MaxFrameSize { get; set; } = 16 * 1024 * 1024;

    /// <summary>消息泵是否要求整帧完整才产出。默认 false(头部到齐即交付,体可为流式)</summary>
    /// <remarks>在消息泵任务上同步消费流式体的实现应重写为 true:否则大帧会在事件内同步等待后续数据,占用线程池线程</remarks>
    protected virtual Boolean RequireFullFrame => false;

    /// <summary>关闭时等待在途帧交付完成的上限(毫秒),默认 1000。0 表示不等待</summary>
    /// <remarks>
    /// <para>数据由接收环写入管道、消息泵异步交付,而关闭(含对端 FIN 的被动关闭)在接收环上同步推进:
    /// 不等在途交付直接关闭,会把“已读入管道但尚未交给上层”的帧连同关闭一起丢掉。
    /// 对端“发完即关”时最后一条消息最先受损(如 MQTT 的 DISCONNECT 丢失,上层把正常断开误判为异常断开并误发遗嘱)。</para>
    /// <para>关闭流程先取消挂起读取,再等待泵退出(泵退出前会把已读出的帧交付完),最多等待本上限;
    /// 处理器长时间不返回时不再等待,按剩余流程关闭。</para>
    /// </remarks>
    public Int32 PumpDrainTimeout { get; set; } = 1000;

    #endregion

    #region 接收流程

    /// <summary>启动接收处理。协议模式先启动消息泵(先于接收环,首个数据到达前就绪)</summary>
    /// <remarks>仅流式会话有消息泵:数据经数据管道定界后逐帧交付,接收环收到的原始字节不再进事件(见 <see cref="RawReceiveTaken"/>)</remarks>
    protected override void StartReceiveFlow()
    {
        if (Protocol == null) return;

        if (AutoReceive) StartMessagePump();
        else WriteLog("协议模式需要自动接收(AutoReceive),拉取模式下消息泵未启动,收到的是原始字节");
    }

    /// <summary>停止接收处理。关闭流程开头先停消息泵(取消挂起读取),随后关闭处理器链与会话</summary>
    protected override void StopReceiveFlow() => StopMessagePump();

    /// <summary>停止接收处理(异步)。取消挂起读取并等待在途帧交付完成,保证在途帧先于关闭交付给上层</summary>
    protected override Task StopReceiveFlowAsync() => StopMessagePumpAsync();

    /// <summary>本轮原始数据是否已被消息泵接管</summary>
    /// <remarks>协议模式下数据经数据管道定界交付,接收环不再把原始字节抛给事件;拉取模式(无泵)时仍走事件</remarks>
    protected override Boolean RawReceiveTaken => _pumpTask != null;

    /// <summary>销毁资源</summary>
    /// <param name="disposing">是否释放托管资源</param>
    protected override void Dispose(Boolean disposing)
    {
        base.Dispose(disposing);

        // 只释放、不置空:置空后并发任务的 ??= 会新建“满额”信号量,其 Release 直接抛 SemaphoreFullException。
        // 保留已释放实例,迟到的 Release 得到 ObjectDisposedException,由收尾逻辑按正常时序忽略
        _concurrency?.Dispose();
    }

    #endregion

    #region 消息泵

    private CancellationTokenSource? _pumpCts;

    /// <summary>泵任务的观察续体。存的是 ContinueWith 续体而非泵本体:仅用于观察失败,以及标记“协议模式已启动”(停止时清空)</summary>
    private Task? _pumpTask;

    /// <summary>泵本体任务。关闭时等待在途帧交付完成用(观察续体与泵本体完成时序不同,不可用于等待)</summary>
    private Task? _pumpRunTask;

    /// <summary>当前异步流是否处于消息泵内。用于识别泵内关闭(处理器内主动关闭、泵内报错关闭),避免等待泵自身退出</summary>
#if NET45
    // net45 无 AsyncLocal,退回线程本地(与 TokenHttpFilter 同一降级口径);异步续体切线程后标记可能丢失,此时退化为等到上限
    private readonly ThreadLocal<Boolean> _inPump = new();
#else
    private readonly AsyncLocal<Boolean> _inPump = new();
#endif

    /// <summary>启动消息泵。协议模式(<see cref="SessionBase.Protocol"/> 非空)下由打开流程与服务端会话启动流程调用</summary>
    /// <remarks>首次访问数据管道确保其在接收环启动前就绪。调用方须保证时序(每次会话生命周期至多一次)</remarks>
    protected void StartMessagePump()
    {
        var codec = Protocol;
        if (codec == null) return;

        // 触建数据管道(接收环启动前就绪,首个数据到达即可投递)
        var reader = Pipe.Reader;

        var cts = new CancellationTokenSource();
        _pumpCts = cts;

        WriteLog("启动消息泵:{0}", codec);

        var task = PumpAsync(new MessagePump(codec) { MaxCache = MaxCache, RequireFullFrame = RequireFullFrame, MaxFrameSize = MaxFrameSize }, reader, cts.Token);
        _pumpRunTask = task;

        // 观察泵任务:泵内异常若无人观察会被静默吞掉,表现为“连接还在但再也收不到消息”,极难定位
        _pumpTask = task.ContinueWith(
            t =>
            {
                try
                {
                    if (t.IsFaulted) OnError("MessagePump", t.Exception!.GetBaseException());
                }
                catch { }
            },
            CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default);
    }

    /// <summary>停止消息泵。取消挂起读取,泵任务随后自行退出(数据管道完成同样唤醒读取)</summary>
    private void StopMessagePump()
    {
        var cts = Interlocked.Exchange(ref _pumpCts, null);
        if (cts != null)
        {
            cts.Cancel();
            cts.Dispose();
        }

        // 清空“泵已启动”标记:RawReceiveTaken 以本生命周期是否已起泵为准。
        // 不清会让标志在泵退出后恒真——后续接收判断仍按“数据已交给泵”静默丢弃原始字节
        _pumpTask = null;
        _pumpRunTask = null;
    }

    /// <summary>停止消息泵(异步)。取消挂起读取,并等待在途帧交付完成,避免关闭前已读出的帧被丢弃</summary>
    /// <remarks>
    /// <para>取消只作用于挂起读取:读取已完成的帧仍会被泵交付,因此等待泵退出即可保证这些帧先于关闭交付给上层。</para>
    /// <para>泵内主动关闭(处理器内、泵内报错)不得等待自身退出,否则要等到上限,直接返回。</para>
    /// </remarks>
    /// <returns>泵收尾任务</returns>
    protected virtual async Task StopMessagePumpAsync()
    {
        var run = Interlocked.Exchange(ref _pumpRunTask, null);

        StopMessagePump();

        var timeout = PumpDrainTimeout;
        if (run == null || run.IsCompleted || timeout <= 0 || _inPump.Value) return;

        // 在途帧交付通常微秒级完成;本等待是处理器长时间不返回时的兜底,超时后不再等待
        var done = await Task.WhenAny(run, Task.Delay(timeout)).ConfigureAwait(false);
        if (done != run) return;

        // 观察泵任务的异常,避免无人消费的故障任务堆积
        try { await run.ConfigureAwait(false); } catch { }
    }

    /// <summary>消息泵循环。定界消息帧并逐帧交付;同连接消息串行处理</summary>
    /// <param name="pump">消息帧泵</param>
    /// <param name="reader">数据管道读取器</param>
    /// <param name="cancellationToken">取消通知(会话关闭)</param>
    private async Task PumpAsync(MessagePump pump, PipeReader reader, CancellationToken cancellationToken)
    {
        // 先让出一次再置标记:异步流标记若在同步前缀赋值,会随执行上下文泄漏到启动会话的调用方,
        // 之后在同一线程上发生的关闭会被误判为“泵内关闭”而跳过等待,丢帧照旧
        await Task.Yield();

        // 标记当前异步流处于泵内:泵内(含处理器)主动关闭会话时不得等待泵自身退出
        _inPump.Value = true;

        // 不按取消令牌提前退出:取消只应终止“等待新数据”,已到达并写入管道的帧仍要交付完再退出。
        // 循环靠读取结果终止——无数据可读时的取消以 OperationCanceledException 逃出,由下方捕获后退出。
        // 反之(读到即退)会把“对端发完即关”时已到达的最后一条消息连同关闭一起丢掉。
        while (true)
        {
            IMessage? message;
            try
            {
                message = await pump.ReadAsync(reader, cancellationToken).ConfigureAwait(false);
            }
            catch (OperationCanceledException)
            {
                break;
            }
            catch (Exception ex)
            {
                // 关闭已在进行(令牌已取消):读取失效随关闭一并收尾,不再当作数据流故障报错与关闭,
                // 否则关闭流程释放管道后泵的收尾读取会刷出无意义的 MessagePumpError
                if (cancellationToken.IsCancellationRequested) break;

                OnError("MessagePump", ex);

                // 数据流不可恢复(协议错误/IO 异常):关闭会话,避免半开连接僵死
                Close("MessagePumpError");
                break;
            }

            // 数据管道完成(连接关闭):退出
            if (message == null) break;

            // 有匹配队列时,入站消息都可能被配对交付给等待方:等待方在事件链之后还要异步消费同一份体,
            // 而流式体不能被二次读(泵会继续读同一 PipeReader,导致单读者冲突或把下一帧字节当成体),因此先物化为内存体。
            // 不能依赖 message.Reply——无方向位协议(如 LengthFieldCodec 的恒真 matcher)恒为 false,会让等待方拿到流式体而串包
            if (MatchQueue != null && message.Body is { IsStreaming: true })
            {
                try
                {
                    var all = await message.Body.ReadAllAsync(cancellationToken).ConfigureAwait(false);
                    message.SetBody(all);
                }
                catch (OperationCanceledException)
                {
                    message.TryDispose();
                    continue;
                }
                catch (Exception ex)
                {
                    // 流式体读满失败(对端半途断开/管道故障):数据流已不可恢复,与帧层读失败同口径关闭会话。
                    // 异常若向上逃出泵任务,续体只记日志不关会话,表现为“连接还在、消息不再到达、也不触发关闭”的半死态
                    message.TryDispose();
                    OnError("MessagePump", ex);
                    Close("MessagePumpError");
                    break;
                }
            }

            // 并行模式:先物化流式体(一次拷贝换并行安全),信号量约束并发后派发;处理顺序不定(SRMP 按序列号配对)
            if (MaxConcurrency > 1)
            {
                if (message.Body is { IsStreaming: true })
                {
                    try
                    {
                        var all = await message.Body.ReadAllAsync(cancellationToken).ConfigureAwait(false);
                        message.SetBody(all);
                    }
                    catch (OperationCanceledException)
                    {
                        message.TryDispose();
                        continue;
                    }
                    catch (Exception ex)
                    {
                        // 同匹配队列分支:流式体读满失败即数据流不可恢复,丢弃消息并关闭会话
                        message.TryDispose();
                        OnError("MessagePump", ex);
                        Close("MessagePumpError");
                        break;
                    }
                }

                try
                {
                    await Concurrency.WaitAsync(cancellationToken).ConfigureAwait(false);
                }
                catch (OperationCanceledException)
                {
                    message.TryDispose();
                    break;
                }
                catch (ObjectDisposedException)
                {
                    // 会话销毁:并发信号量已释放,消息无人处理
                    message.TryDispose();
                    break;
                }

                // Task.Run 派发:async 方法首段同步执行,须真正切换线程池,否则同步处理器会阻塞泵循环
                _ = Task.Run(() => ProcessMessageAsync(message, cancellationToken, true));
                continue;
            }

            // 串行模式:同连接依次处理(前一条收尾后才读下一帧)
            await ProcessMessageAsync(message, cancellationToken, false).ConfigureAwait(false);
        }
    }

    /// <summary>处理单个消息:统一进入事件链(可观测)→ 响应尝试匹配交付 → 收尾</summary>
    /// <param name="message">消息</param>
    /// <param name="cancellationToken">取消通知</param>
    /// <param name="releaseSlot">完成后是否释放并发信号量(并行派发为 true)</param>
    /// <remarks>
    /// <para>交付收尾:未交付等待方的消息丢弃未读负载对齐帧尾,随后释放;命中交付的消息由等待方释放。</para>
    /// <para>并行模式(<see cref="SessionBase.MaxConcurrency"/> 大于1)下由独立任务调用,异常统一经 <see cref="SessionBase.OnError"/> 上报,不向任务外部抛出。</para>
    /// </remarks>
    private async Task ProcessMessageAsync(IMessage message, CancellationToken cancellationToken, Boolean releaseSlot)
    {
        try
        {
            var delivered = false;
            try
            {
                await OnMessageAsync(message).ConfigureAwait(false);

                // 事件链(可观测)可能已在处理器内读空消息体;命中配对的等待方还要消费同一份体,
                // 交付前把内存模式体复位到起点(流式体不参与配对交付,不可重放)
                if (message.Body is { IsStreaming: false } body) body.Reset();

                delivered = TryMatchResponse(message);
            }
            catch (OperationCanceledException) { }
            catch (Exception ex)
            {
                OnError("OnMessage", ex);
            }
            finally
            {
                // 交付收尾:未交付等待方的消息丢弃未读负载对齐帧尾,随后释放
                if (!delivered)
                {
                    try
                    {
                        await MessagePump.DiscardAsync(message, cancellationToken).ConfigureAwait(false);
                    }
                    catch (OperationCanceledException) { }
                    finally
                    {
                        message.TryDispose();
                    }
                }
            }
        }
        catch (Exception ex)
        {
            // 兕底:防止 fire-and-forget 任务异常未观察
            OnError("ProcessMessage", ex);
            message.TryDispose();
        }
        finally
        {
            // 会话销毁时信号量可能已释放:任务收尾晚于 Dispose 属正常时序,忽略该异常
            if (releaseSlot)
            {
                try { Concurrency.Release(); }
                catch (ObjectDisposedException) { }
            }
        }
    }

    private SemaphoreSlim? _concurrency;

    /// <summary>并发信号量。并行模式(<see cref="SessionBase.MaxConcurrency"/> 大于1)下约束同连接并发处理数,等待时形成背压</summary>
    /// <remarks>惰性创建后<b>只释放、不置空</b>:置空后迟到的归还会经 <c>??=</c> 新建出“满额”信号量,
    /// 其 <c>Release</c> 抛 <see cref="SemaphoreFullException"/>;保留已释放实例则只抛
    /// <see cref="ObjectDisposedException"/>,由收尾逻辑按正常时序忽略。</remarks>
    private SemaphoreSlim Concurrency => _concurrency ??= new SemaphoreSlim(MaxConcurrency, MaxConcurrency);

    #endregion
}