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

namespace NewLife.Messaging;

/// <summary>事件分发枢纽。按主题路由网络消息或本地事件到各事件总线</summary>
/// <remarks>
/// <para><b>核心职责</b>:</para>
/// <list type="number">
/// <item><description><b>主题路由</b>:维护 <c>topic → IEventBus&lt;TEvent&gt;</c> 映射,提供 <see cref="GetEventBus"/>/<see cref="RegisterBus"/>。</description></item>
/// <item><description><b>协议适配</b>:默认解析 <c>event#topic#clientId#message</c> 文本协议;
/// 派生类可重写 <see cref="TryDecode(IPacket,out EventEnvelope)"/> / <see cref="TryDecode(String,out EventEnvelope)"/> 替换为二进制、MQTT 等协议。</description></item>
/// <item><description><b>控制面 + 数据面</b>:解析后区分动作(订阅/取消订阅)与数据消息,分别走 <see cref="SubscribeAsync"/>/<see cref="UnsubscribeAsync"/> 或 <see cref="PublishAsync"/>。</description></item>
/// </list>
/// <para><b>典型应用</b>:作为 IoT 网关或 MQ Bridge,将网络层(如 WebSocket/TCP)收到的消息按 topic 投递到进程内订阅者。</para>
/// </remarks>
/// <typeparam name="TEvent">事件类型</typeparam>
public class EventHub<TEvent> : IEventHandler<IPacket>, IEventHandler<String>, ILogFeature, ITracerFeature
{
    #region 嵌套:事件信封
    /// <summary>协议解码的中间结果,承载主题/客户端/动作或事件三种语义</summary>
    protected readonly struct EventEnvelope
    {
        /// <summary>主题</summary>
        public String Topic { get; }

        /// <summary>客户端标识。发送方,路由分发时用于排除回环</summary>
        public String ClientId { get; }

        /// <summary>动作。订阅/取消订阅控制指令;为空表示数据消息</summary>
        public String? Action { get; }

        /// <summary>已解码的事件实例</summary>
        public TEvent? Event { get; }

        /// <summary>是否动作信封</summary>
        public Boolean IsAction => !Action.IsNullOrEmpty();

        /// <summary>构造动作信封</summary>
        /// <param name="topic">主题</param>
        /// <param name="clientId">客户端标识</param>
        /// <param name="action">动作指令</param>
        public static EventEnvelope ForAction(String topic, String clientId, String action) => new(topic, clientId, action, default);

        /// <summary>构造事件信封</summary>
        /// <param name="topic">主题</param>
        /// <param name="clientId">客户端标识</param>
        /// <param name="event">事件实例</param>
        public static EventEnvelope ForEvent(String topic, String clientId, TEvent @event) => new(topic, clientId, null, @event);

        private EventEnvelope(String topic, String clientId, String? action, TEvent? @event)
        {
            Topic = topic;
            ClientId = clientId;
            Action = action;
            Event = @event;
        }
    }
    #endregion

    #region 属性
    /// <summary>事件总线工厂。用于按需创建各 topic 对应的事件总线</summary>
    public IEventBusFactory? Factory { get; set; }

    /// <summary>JSON 主机。用于编解码事件体</summary>
    public IJsonHost JsonHost { get; set; } = JsonHelper.Default;

    /// <summary>JSON序列化选项,影响复杂对象的编码和解码行为</summary>
    public JsonOptions? JsonOptions { get; set; }

    /// <summary>动作判定阈值。消息体短于该长度且不以 <c>{</c> 开头时视为控制动作指令</summary>
    public Int32 ActionMaxLength { get; set; } = 32;

    /// <summary>链路追踪</summary>
    public ITracer? Tracer { get; set; }

    private readonly ConcurrentDictionary<String, IEventBus<TEvent>> _eventBuses = new();
    /// <summary>已注册的事件总线集合。Key 为 topic</summary>
    public IDictionary<String, IEventBus<TEvent>> EventBuses => _eventBuses;

    private static readonly Byte[] _prefixBytes = Encoding.ASCII.GetBytes("event#");
    private static readonly Char[] _prefixChars = "event#".ToCharArray();
    #endregion

    #region 注册与获取
    /// <summary>注册主题对应的事件总线</summary>
    /// <param name="topic">事件主题</param>
    /// <param name="eventBus">事件总线实例</param>
    public void RegisterBus(String topic, IEventBus<TEvent> eventBus) => _eventBuses[topic] = eventBus;

    /// <summary>尝试获取主题对应的事件总线(不会创建新的)</summary>
    /// <param name="topic">事件主题</param>
    /// <param name="eventBus">事件总线</param>
    /// <returns>是否存在</returns>
    public Boolean TryGetBus(String topic, out IEventBus<TEvent> eventBus) => _eventBuses.TryGetValue(topic, out eventBus!);

    /// <summary>获取或创建主题对应的事件总线</summary>
    /// <param name="topic">事件主题</param>
    /// <param name="clientId">客户标识,仅在使用 <see cref="Factory"/> 创建时传递</param>
    /// <returns>事件总线实例</returns>
    public virtual IEventBus<TEvent> GetEventBus(String topic, String clientId = "")
    {
        if (_eventBuses.TryGetValue(topic, out var bus)) return bus;
        bus = Factory?.CreateEventBus<TEvent>(topic, clientId) ?? new EventBus<TEvent>();
        return _eventBuses.GetOrAdd(topic, bus);
    }
    #endregion

    #region 发布与订阅
    /// <summary>向指定主题发布事件</summary>
    /// <remarks>
    /// 发送方信息(用于排除回环)通过 <see cref="EventContext.ClientId"/> 传递。
    /// 无需排除回环时可传入 <see langword="null"/>,由总线自动创建上下文。
    /// </remarks>
    /// <param name="topic">主题</param>
    /// <param name="event">事件</param>
    /// <param name="context">事件上下文;为空时由总线自行创建</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>成功处理该事件的处理器数量</returns>
    public virtual Task<Int32> PublishAsync(String topic, TEvent @event, IEventContext? context = null, CancellationToken cancellationToken = default)
    {
        var bus = GetEventBus(topic);

        // 自动注入 Topic 到上下文,便于订阅者读取
        if (context is EventContext ctx) ctx.Topic ??= topic;

        return bus.PublishAsync(@event, context, cancellationToken);
    }

    /// <summary>向指定主题订阅事件</summary>
    /// <param name="topic">主题</param>
    /// <param name="clientId">客户标识</param>
    /// <param name="handler">事件处理器</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>是否订阅成功</returns>
    public virtual Task<Boolean> SubscribeAsync(String topic, String clientId, IEventHandler<TEvent> handler, CancellationToken cancellationToken = default) => GetEventBus(topic, clientId).SubscribeAsync(handler, clientId, cancellationToken);

    /// <summary>取消指定主题/客户端的订阅</summary>
    /// <param name="topic">主题</param>
    /// <param name="clientId">客户标识</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>是否成功取消订阅</returns>
    public virtual async Task<Boolean> UnsubscribeAsync(String topic, String clientId, CancellationToken cancellationToken = default)
    {
        if (!_eventBuses.TryGetValue(topic, out var bus)) return false;

        var ok = await bus.UnsubscribeAsync(clientId, cancellationToken).ConfigureAwait(false);

        // 取消后若已无订阅者且为默认实现,移除总线避免泄漏。
        // 采用"先移除再复查放回"缩小检查-再行动窗口:若移除后恰好有并发订阅注册到该总线,则重新放回
        if (bus is EventBus<TEvent> eb && eb.Handlers.Count == 0)
        {
            if (_eventBuses.TryRemove(topic, out _) && eb.Handlers.Count > 0)
                _eventBuses.TryAdd(topic, eb);
        }

        return ok;
    }
    #endregion

    #region 协议解码(可重载)
    /// <summary>尝试从二进制数据包解码事件信封。派生类可重写以替换协议实现</summary>
    /// <param name="data">网络数据包</param>
    /// <param name="envelope">解码出的事件信封</param>
    /// <returns>是否解码成功</returns>
    protected virtual Boolean TryDecode(IPacket data, out EventEnvelope envelope)
    {
        envelope = default;
        if (data == null) return false;

        if (data.Next == null)
        {
            if (!TryParseHeader(data.GetSpan(), out var topic, out var clientId, out var headerLen)) return false;

            return DecodePayload(data, topic, clientId, headerLen, out envelope);
        }
        else
        {
            // 头部 event#topic#clientId# 可能跨节点(跨接收轮组链的帧):跨段扫描定位头部末端,物化头部字节后复用跨度解析
            var headLength = FindHeaderLength(data);
            if (headLength <= 0) return false;

            if (!TryParseHeader(data.ReadBytes(0, headLength), out var topic, out var clientId, out var headerLen)) return false;

            return DecodePayload(data, topic, clientId, headerLen, out envelope);
        }
    }

    /// <summary>跨段扫描头部末端。event#topic#clientId# 的分隔符为单字节 '#',用序列读取器跨段扫描</summary>
    /// <param name="data">数据包链</param>
    /// <returns>头部字节数(含末尾 '#');分隔符不足 3 个时返回 0</returns>
    private static Int32 FindHeaderLength(IPacket data)
    {
        var reader = new SequenceReader<Byte>(data.AsReadOnlySequence());
        for (var i = 0; i < 3; i++)
        {
            if (!reader.TryAdvanceTo((Byte)'#', true)) return 0;
        }

        return (Int32)reader.Consumed;
    }

    /// <summary>从头部之后的负载解码事件信封</summary>
    /// <param name="data">网络数据包</param>
    /// <param name="topic">主题</param>
    /// <param name="clientId">客户端标识</param>
    /// <param name="headerLen">头部字节数</param>
    /// <param name="envelope">解码出的事件信封</param>
    /// <returns>是否解码成功</returns>
    /// <remarks>
    /// <para>当 <typeparamref name="TEvent"/> 本身是数据包时,负载切片所有权随之转移给订阅者(订阅者负责释放);
    /// 其它事件类型则切片仅供取值,本方法内部归还引用计数。</para>
    /// </remarks>
    private Boolean DecodePayload(IPacket data, String topic, String clientId, Int32 headerLen, out EventEnvelope envelope)
    {
        envelope = default;

        var msg = data.Slice(headerLen);
        if (msg.Length == 0)
        {
            msg.TryDispose();

            return false;
        }

        // TEvent 本身是 IPacket,直接零拷贝构造事件信封(切片随信封交给订阅者释放)
        if (msg is TEvent evt) { envelope = EventEnvelope.ForEvent(topic, clientId, evt); return true; }

        // 其它事件类型:切片只用于取值,取值后立即归还引用计数(不归还则缓冲永远回不了池)
        try
        {
            var msgStr = msg.ToStr();

            return BuildEnvelopeFromString(topic, clientId, msgStr, out envelope);
        }
        finally
        {
            msg.TryDispose();
        }
    }

    /// <summary>尝试从字符串解码事件信封。派生类可重写以替换协议实现</summary>
    /// <param name="data">字符串消息</param>
    /// <param name="envelope">解码出的事件信封</param>
    /// <returns>是否解码成功</returns>
    protected virtual Boolean TryDecode(String data, out EventEnvelope envelope)
    {
        envelope = default;
        if (data.IsNullOrEmpty()) return false;
        if (!TryParseHeader(data.AsSpan(), out var topic, out var clientId, out var headerLen)) return false;

        var msg = data[headerLen..];
        if (msg.Length == 0) return false;

        return BuildEnvelopeFromString(topic, clientId, msg, out envelope);
    }

    /// <summary>编码事件为字符串(便于通过文本协议发送)</summary>
    /// <param name="topic">主题</param>
    /// <param name="clientId">客户端标识</param>
    /// <param name="event">事件实例</param>
    /// <returns>编码后的字符串</returns>
    protected virtual String EncodeEvent(String topic, String clientId, TEvent @event)
    {
        var body = @event is String s ? s : JsonHost.Write(@event!, JsonOptions);
        return $"event#{topic}#{clientId}#{body}";
    }

    /// <summary>编码控制动作为字符串(订阅/取消订阅等)</summary>
    /// <param name="topic">主题</param>
    /// <param name="clientId">客户端标识</param>
    /// <param name="action">动作指令</param>
    /// <returns>编码后的字符串</returns>
    protected virtual String EncodeAction(String topic, String clientId, String action) => $"event#{topic}#{clientId}#{action}";

    /// <summary>根据消息体字符串构造事件信封:动作 / 字符串事件 / JSON 解码</summary>
    private Boolean BuildEnvelopeFromString(String topic, String clientId, String msg, out EventEnvelope envelope)
    {
        // TEvent = String,整条消息体始终视为事件,不识别动作(避免短消息被误判为控制指令)
        if (msg is TEvent strEvt)
        {
            envelope = EventEnvelope.ForEvent(topic, clientId, strEvt);
            return true;
        }

        // 短字符串且不以 { 开头 → 视为动作指令
        if (msg[0] != '{' && msg.Length < ActionMaxLength)
        {
            envelope = EventEnvelope.ForAction(topic, clientId, msg);
            return true;
        }

        // JSON 反序列化
        var evt = JsonHost.Read<TEvent>(msg, JsonOptions);
        if (evt == null)
        {
            envelope = default;
            return false;
        }
        envelope = EventEnvelope.ForEvent(topic, clientId, evt);
        return true;
    }

    /// <summary>解析 <c>event#topic#clientId#</c> 二进制头部</summary>
    /// <param name="data">输入数据</param>
    /// <param name="topic">输出主题</param>
    /// <param name="clientId">输出客户端标识</param>
    /// <param name="headerLength">头部字节数(含末尾 #)</param>
    /// <returns>是否解析成功</returns>
    public static Boolean TryParseHeader(ReadOnlySpan<Byte> data, out String topic, out String clientId, out Int32 headerLength)
    {
        topic = clientId = String.Empty;
        headerLength = 0;
        if (!data.StartsWith(_prefixBytes)) return false;

        var p = data.IndexOf((Byte)'#');
        var rest = data[(p + 1)..];
        var p2 = rest.IndexOf((Byte)'#');
        if (p2 <= 0) return false;
        topic = rest[..p2].ToStr();

        var rest2 = rest[(p2 + 1)..];
        var p3 = rest2.IndexOf((Byte)'#');
        if (p3 <= 0) return false;
        clientId = rest2[..p3].ToStr();

        headerLength = p + 1 + p2 + 1 + p3 + 1;
        return true;
    }

    /// <summary>解析 <c>event#topic#clientId#</c> 字符串头部</summary>
    /// <param name="data">输入数据</param>
    /// <param name="topic">输出主题</param>
    /// <param name="clientId">输出客户端标识</param>
    /// <param name="headerLength">头部字符数(含末尾 #)</param>
    /// <returns>是否解析成功</returns>
    public static Boolean TryParseHeader(ReadOnlySpan<Char> data, out String topic, out String clientId, out Int32 headerLength)
    {
        topic = clientId = String.Empty;
        headerLength = 0;
        if (!data.StartsWith(_prefixChars)) return false;

        var p = data.IndexOf('#');
        var rest = data[(p + 1)..];
        var p2 = rest.IndexOf('#');
        if (p2 <= 0) return false;
        topic = rest[..p2].ToString();

        var rest2 = rest[(p2 + 1)..];
        var p3 = rest2.IndexOf('#');
        if (p3 <= 0) return false;
        clientId = rest2[..p3].ToString();

        headerLength = p + 1 + p2 + 1 + p3 + 1;
        return true;
    }
    #endregion

    #region 网络消息接收
    /// <summary>接收网络字节流消息并按协议解码后路由</summary>
    /// <param name="data">网络数据包</param>
    /// <param name="context">事件上下文</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>分发到本地订阅者的处理器数量</returns>
    public virtual async Task<Int32> OnReceiveAsync(IPacket data, IEventContext? context = null, CancellationToken cancellationToken = default)
    {
        if (data == null) return 0;
        if (!TryDecode(data, out var envelope)) return 0;

        return await DispatchEnvelopeAsync(envelope, data, context, cancellationToken).ConfigureAwait(false);
    }

    /// <summary>接收字符串消息并按协议解码后路由</summary>
    /// <param name="data">字符串消息</param>
    /// <param name="context">事件上下文</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>分发到本地订阅者的处理器数量</returns>
    public virtual async Task<Int32> OnReceiveAsync(String data, IEventContext? context = null, CancellationToken cancellationToken = default)
    {
        if (data.IsNullOrEmpty()) return 0;
        if (!TryDecode(data, out var envelope)) return 0;

        return await DispatchEnvelopeAsync(envelope, data, context, cancellationToken).ConfigureAwait(false);
    }

    Task IEventHandler<IPacket>.HandleAsync(IPacket @event, IEventContext? context, CancellationToken cancellationToken) => OnReceiveAsync(@event, context, cancellationToken);
    Task IEventHandler<String>.HandleAsync(String @event, IEventContext? context, CancellationToken cancellationToken) => OnReceiveAsync(@event, context, cancellationToken);

    /// <summary>把解码后的事件信封路由到控制面或数据面</summary>
    private async Task<Int32> DispatchEnvelopeAsync(EventEnvelope envelope, Object raw, IEventContext? context, CancellationToken cancellationToken)
    {
        // 把原始数据透传到扩展项,便于网络层做后续路由(如转发到其他客户端)
        if (context is IExtend ext) ext["Raw"] = raw;

        // 控制面:订阅/取消订阅
        if (envelope.IsAction)
        {
            var action = envelope.Action!;
            if (action.EqualIgnoreCase("subscribe"))
            {
                if ((context as IExtend)?["Handler"] is IEventHandler<TEvent> handler)
                {
                    await SubscribeAsync(envelope.Topic, envelope.ClientId, handler, cancellationToken).ConfigureAwait(false);
                    return 1;
                }
                return 0;
            }
            if (action.EqualIgnoreCase("unsubscribe"))
            {
                await UnsubscribeAsync(envelope.Topic, envelope.ClientId, cancellationToken).ConfigureAwait(false);
                return 1;
            }
            return 0;
        }

        // 数据面:发布事件
        if (envelope.Event is null) return 0;

        // 将发送方 ClientId 注入上下文,便于本地总线排除回环
        if (context is EventContext ec)
        {
            ec.ClientId ??= envelope.ClientId;
        }
        else if (context == null && !envelope.ClientId.IsNullOrEmpty())
        {
            context = new EventContext { Topic = envelope.Topic, ClientId = envelope.ClientId };
        }

        return await PublishAsync(envelope.Topic, envelope.Event, context, cancellationToken).ConfigureAwait(false);
    }
    #endregion

    #region 等待接收
    /// <summary>异步等待指定主题的第一条事件</summary>
    /// <param name="topic">主题</param>
    /// <param name="cancellationToken">取消令牌</param>
    /// <returns>事件实例</returns>
    public async Task<TEvent> ReceiveAsync(String topic, CancellationToken cancellationToken = default)
    {
        var bus = GetEventBus(topic);
        try
        {
            return await bus.ReceiveAsync(cancellationToken).ConfigureAwait(false);
        }
        finally
        {
            // 等待完成后若无其他订阅者,移除总线避免泄漏。
            // 采用"先移除再复查放回"缩小检查-再行动窗口,避免并发订阅丢失
            if (bus is EventBus<TEvent> eb && eb.Handlers.Count == 0)
            {
                if (_eventBuses.TryRemove(topic, out _) && eb.Handlers.Count > 0)
                    _eventBuses.TryAdd(topic, eb);
            }
        }
    }

    /// <summary>异步等待指定主题的第一条事件,超时后抛出 <see cref="OperationCanceledException"/></summary>
    /// <param name="topic">主题</param>
    /// <param name="timeout">超时时间</param>
    /// <returns>事件实例</returns>
    public async Task<TEvent> ReceiveAsync(String topic, TimeSpan timeout)
    {
        using var cts = new CancellationTokenSource(timeout);
        return await ReceiveAsync(topic, cts.Token).ConfigureAwait(false);
    }
    #endregion

    #region 日志
    /// <summary>日志</summary>
    public ILog Log { get; set; } = Logger.Null;

    /// <summary>写日志</summary>
    /// <param name="format">格式串</param>
    /// <param name="args">参数</param>
    public void WriteLog(String format, params Object[] args) => Log?.Info(format, args);
    #endregion
}