修正命名错误
xiyunfei authored at 2023-04-09 21:10:58
5.15 KiB
NewLife.JT808
using JT808.Data;
using JT808.Server.Services;
using NewLife;
using NewLife.Caching;
using NewLife.JT808.Models;
using NewLife.JT808.Protocols;
using NewLife.Log;
using NewLife.Serialization;
using NewLife.Threading;

namespace JT808.Server.HostedServices;

/// <summary>JT808 TCP 网关服务生命周期管理</summary>
/// <remarks>
/// 参考 GPS808 的 GpsHostedService:启动 JT808Server 并管理其生命周期。
/// 同时消费 Redis command 队列,处理 Web 端下发的指令。
/// </remarks>
public class GpsHostedService : IHostedService
{
    private readonly JT808Server _server;
    private readonly ICacheProvider _cacheProvider;
    private TimerX? _commandTimer;

    public GpsHostedService(ICacheProvider cacheProvider)
    {
        _cacheProvider = cacheProvider;

        _server = new JT808Server
        {
            Port = 808,
            Log = XTrace.Log,
            SessionLog = XTrace.Log,
            Tracer = DefaultTracer.Instance,
        };

        // 注册指令处理器
        _server.AddHandler(new Handlers.NormalController());
        _server.AddHandler(new Handlers.PositionController());
        _server.AddHandler(new Handlers.InfoController());
        _server.AddHandler(new Handlers.MediaController());
        _server.AddHandler(new Handlers.FileHandler());
    }

    public Task StartAsync(CancellationToken cancellationToken)
    {
        _server.Start();
        XTrace.WriteLine("JT808 网关服务已启动,端口:{0}", _server.Port);

        // 启动 Redis 指令队列消费(每 1 秒轮询)
        _commandTimer = new TimerX(DoConsumeCommand, null, 5_000, 1_000) { Async = true };
        XTrace.WriteLine("Redis 指令队列消费已启动");

        return Task.CompletedTask;
    }

    public Task StopAsync(CancellationToken cancellationToken)
    {
        _commandTimer.TryDispose();
        _commandTimer = null;
        _server.Stop("服务停止");
        XTrace.WriteLine("JT808 网关服务已停止");
        return Task.CompletedTask;
    }

    /// <summary>消费 Redis command 队列,下发指令到终端</summary>
    private void DoConsumeCommand(Object? state)
    {
        try
        {
            var queue = _cacheProvider.GetQueue<CommandModel>("Command", "Group");
            if (queue == null) return;

            while (queue.Count > 0)
            {
                var cmd = queue.TakeOne();
                if (cmd == null) break;

                XTrace.WriteLine("收到指令队列消息:Mobile={0}, Type={1}", cmd.Mobile, cmd.BodyType);

                // 反序列化消息体
                if (cmd.BodyType.IsNullOrEmpty() || cmd.BodyData.IsNullOrEmpty()) continue;

                var type = Type.GetType(cmd.BodyType);
                if (type == null)
                {
                    XTrace.WriteLine("无法解析消息体类型:{0}", cmd.BodyType);
                    continue;
                }

                var body = JsonHelper.ToJsonEntity(cmd.BodyData, type);
                if (body == null)
                {
                    XTrace.WriteLine("无法反序列化消息体:{0}", cmd.BodyData);
                    continue;
                }

                // 查找会话并下发
                var session = _server.Manager.Find(cmd.Mobile) as JT808Session;
                if (session != null)
                {
                    // 使用 SendJT808 发送消息体(内部构造 JTMessage)
                    session.SendJT808(body);
                    XTrace.WriteLine("指令已下发到终端:Mobile={0}, Type={1}", cmd.Mobile, cmd.BodyType);

                    // 记录指令
                    var device = Device.FindByMobile(cmd.Mobile);
                    if (device != null)
                    {
                        var cmdRecord = new DeviceCommand
                        {
                            DeviceId = device.Id,
                            Mobile = cmd.Mobile,
                            Kind = 0,
                            Status = 1, // 已下发
                            BodyType = cmd.BodyType,
                            BodyData = cmd.BodyData,
                            CreateTime = DateTime.Now,
                        };
                        _ = Task.Run(() => cmdRecord.InsertAsync());
                    }
                }
                else
                {
                    XTrace.WriteLine("设备不在线,无法下发指令:Mobile={0}", cmd.Mobile);
                }
            }
        }
        catch (Exception ex)
        {
            XTrace.WriteLine("指令队列消费异常:{0}", ex.Message);
        }
    }
}

/// <summary>指令模型。与 Web 端 CommandClient 保持一致</summary>
internal class CommandModel
{
    /// <summary>目标终端手机号</summary>
    public String Mobile { get; set; } = String.Empty;

    /// <summary>消息体类型全名</summary>
    public String BodyType { get; set; } = String.Empty;

    /// <summary>消息体 JSON 数据</summary>
    public String BodyData { get; set; } = String.Empty;

    /// <summary>来源</summary>
    public String? Source { get; set; }

    /// <summary>创建时间</summary>
    public DateTime CreateTime { get; set; }
}