From e8983f2e95b8aa82eed8072c61052bc66f0c83fd Mon Sep 17 00:00:00 2001 From: Cxx0822 <1556464090@qq.com> Date: Wed, 22 Apr 2026 20:01:55 +0800 Subject: [PATCH] =?UTF-8?q?feat:=E5=A2=9E=E5=8A=A0Mqtt=E5=AE=A2=E6=88=B7?= =?UTF-8?q?=E7=AB=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- SocketHub/Endpoints/MqttService.cs | 69 +++++++++ SocketHub/Endpoints/TcpService.cs | 16 ++ SocketHub/Filter/BinaryPipelineFilter.cs | 74 +++++++++ SocketHub/Handlers/TcpPackageHandler.cs | 23 +++ SocketHub/Handlers/WebSocketMessageHandler.cs | 22 +++ SocketHub/Models/MqttSettings.cs | 16 ++ SocketHub/Models/ProtocolFrame.cs | 23 +++ SocketHub/NLog.config | 2 +- SocketHub/Program.cs | 145 ++---------------- SocketHub/SocketHub.csproj | 1 + SocketHub/appsettings.json | 9 ++ 11 files changed, 266 insertions(+), 134 deletions(-) create mode 100644 SocketHub/Endpoints/MqttService.cs create mode 100644 SocketHub/Endpoints/TcpService.cs create mode 100644 SocketHub/Filter/BinaryPipelineFilter.cs create mode 100644 SocketHub/Handlers/TcpPackageHandler.cs create mode 100644 SocketHub/Handlers/WebSocketMessageHandler.cs create mode 100644 SocketHub/Models/MqttSettings.cs create mode 100644 SocketHub/Models/ProtocolFrame.cs diff --git a/SocketHub/Endpoints/MqttService.cs b/SocketHub/Endpoints/MqttService.cs new file mode 100644 index 0000000..9c3fd9f --- /dev/null +++ b/SocketHub/Endpoints/MqttService.cs @@ -0,0 +1,69 @@ +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Options; +using MQTTnet; +using MQTTnet.Protocol; +using SocketHub.Models; +using System.Text; + +namespace SocketHub.Endpoints +{ + internal class MqttService : BackgroundService + { + private readonly IMqttClient _mqttClient; + private readonly MqttSettings _settings; + private MqttClientOptions? _mqttOptions; + + // 直接注入配置 + MQTT 客户端 + public MqttService(IMqttClient mqttClient, IOptions settings) + { + _mqttClient = mqttClient; + _settings = settings.Value; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + // 构建连接配置 + _mqttOptions = new MqttClientOptionsBuilder() + .WithClientId($"{_settings.ClientId}_{Guid.NewGuid():N}") + .WithTcpServer(_settings.Host, _settings.Port) + .WithCleanSession() + .Build(); + + _mqttClient.ConnectedAsync += ConnectedAsync; + _mqttClient.ApplicationMessageReceivedAsync += HandleMessage; + + // 等待程序停止 + await Task.Delay(Timeout.Infinite, stoppingToken); + } + + // 连接成功事件 + public async Task ConnectedAsync(MqttClientConnectedEventArgs arg) + { + // 连接成功后订阅主题 + foreach (var topic in _settings.Topics) + { + await _mqttClient.SubscribeAsync(topic, MqttQualityOfServiceLevel.AtLeastOnce); + Console.WriteLine($"[MQTT] 已订阅主题:{topic}"); + } + + await Task.CompletedTask; + } + + public async Task HandleMessage(MqttApplicationMessageReceivedEventArgs arg) + { + var topic = arg.ApplicationMessage.Topic; + var payload = Encoding.UTF8.GetString(arg.ApplicationMessage.Payload); + + Console.WriteLine($"\n[MQTT] 收到消息\n主题:{topic}\n内容:{payload}\n"); + + await Task.CompletedTask; + } + + // 优雅停止 + public override async Task StopAsync(CancellationToken stoppingToken) + { + await _mqttClient.DisconnectAsync(); + await base.StopAsync(stoppingToken); + } + } +} diff --git a/SocketHub/Endpoints/TcpService.cs b/SocketHub/Endpoints/TcpService.cs new file mode 100644 index 0000000..ef2a606 --- /dev/null +++ b/SocketHub/Endpoints/TcpService.cs @@ -0,0 +1,16 @@ +using Microsoft.Extensions.Options; +using SocketHub.Models; +using SuperSocket.Server; +using SuperSocket.Server.Abstractions; + +namespace SocketHub.Endpoints +{ + internal class TcpService : SuperSocketService + { + public TcpService(IServiceProvider serviceProvider, IOptions serverOptions) + : base(serviceProvider, serverOptions) + { + + } + } +} diff --git a/SocketHub/Filter/BinaryPipelineFilter.cs b/SocketHub/Filter/BinaryPipelineFilter.cs new file mode 100644 index 0000000..06890c3 --- /dev/null +++ b/SocketHub/Filter/BinaryPipelineFilter.cs @@ -0,0 +1,74 @@ +using SocketHub.Models; +using SuperSocket.ProtoBase; +using System.Buffers; + +namespace SocketHub.Filter +{ + internal class BinaryPipelineFilter : FixedHeaderPipelineFilter + { + /// + /// 固定头长度 + /// + public BinaryPipelineFilter() : base(6) + { + } + + /// + /// 从包头中解析出 Body 长度 + /// + /// 字节流 + /// Body 长度 + protected override int GetBodyLengthFromHeader(ref ReadOnlySequence buffer) + { + var reader = new SequenceReader(buffer); + + // 读取前2字节 → magic + reader.TryReadBigEndian(out ushort magic); + + // 再读2字节 → Length + reader.TryReadBigEndian(out ushort length); + + // 再读2字节 → Type + reader.TryReadBigEndian(out ushort type); + + // 校验帧头 + if (magic != 0x656D) + { + throw new Exception("非法帧头"); + } + + // 限制长度 + if (length == 0 || length > 8192) + { + throw new Exception("非法长度"); + } + + // 返回 Payload 长度 + return length; + } + + /// + /// 把完整字节包 → 转换成 ProtocolFrame + /// + /// 字节流 + /// ProtocolFrame + protected override ProtocolFrame DecodePackage(ref ReadOnlySequence buffer) + { + var reader = new SequenceReader(buffer); + + // 读取前2字节 → magic + reader.TryReadBigEndian(out ushort magic); + + // 再读2字节 → Length + reader.TryReadBigEndian(out ushort length); + + // 再读2字节 → Type + reader.TryReadBigEndian(out ushort type); + + var payload = buffer.Slice(6, length).ToArray(); + + // 构造ProtocolFrame + return new ProtocolFrame(magic, length, type, payload); + } + } +} diff --git a/SocketHub/Handlers/TcpPackageHandler.cs b/SocketHub/Handlers/TcpPackageHandler.cs new file mode 100644 index 0000000..3a4dfe0 --- /dev/null +++ b/SocketHub/Handlers/TcpPackageHandler.cs @@ -0,0 +1,23 @@ +using Microsoft.Extensions.Logging; +using SocketHub.Models; +using SuperSocket.Server.Abstractions; +using SuperSocket.Server.Abstractions.Session; + +namespace SocketHub.Handler +{ + internal class TcpPackageHandler : IPackageHandler + { + private readonly ILogger _logger; + + public TcpPackageHandler(ILogger logger) + { + _logger = logger; + } + + public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken) + { + _logger.LogInformation($"Magic={package.Magic:X4}, Type={package.Type}, Len={package.Length}"); + await session.SendAsync(package.Payload, cancellationToken); + } + } +} diff --git a/SocketHub/Handlers/WebSocketMessageHandler.cs b/SocketHub/Handlers/WebSocketMessageHandler.cs new file mode 100644 index 0000000..593c4e9 --- /dev/null +++ b/SocketHub/Handlers/WebSocketMessageHandler.cs @@ -0,0 +1,22 @@ +using Microsoft.Extensions.Logging; +using SuperSocket.WebSocket; +using SuperSocket.WebSocket.Server; + +namespace SocketHub.Handler +{ + internal class WebSocketMessageHandler + { + private readonly ILogger _logger; + + public WebSocketMessageHandler(ILogger logger) + { + _logger = logger; + } + + public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package) + { + _logger.LogInformation($"[WebSocket] {package.Message}"); + await session.SendAsync("ok"); + } + } +} diff --git a/SocketHub/Models/MqttSettings.cs b/SocketHub/Models/MqttSettings.cs new file mode 100644 index 0000000..bb66ceb --- /dev/null +++ b/SocketHub/Models/MqttSettings.cs @@ -0,0 +1,16 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace SocketHub.Models +{ + internal class MqttSettings + { + public string Host { get; set; } = string.Empty; + public int Port { get; set; } = 1883; + public string ClientId { get; set; } = string.Empty; + public List Topics { get; set; } = new(); + } +} diff --git a/SocketHub/Models/ProtocolFrame.cs b/SocketHub/Models/ProtocolFrame.cs new file mode 100644 index 0000000..4910bee --- /dev/null +++ b/SocketHub/Models/ProtocolFrame.cs @@ -0,0 +1,23 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace SocketHub.Models +{ + /// + /// 协议数据结构 示例:65 6D 00 05 00 01 68 65 6C 6C 6F + /// + /// 帧头标识 + /// 数据长度 + /// 数据类型 + /// 数据内容 + /// + internal record ProtocolFrame( + ushort Magic, + ushort Length, + ushort Type, + byte[] Payload + ); +} diff --git a/SocketHub/NLog.config b/SocketHub/NLog.config index 1067898..26ca53b 100644 --- a/SocketHub/NLog.config +++ b/SocketHub/NLog.config @@ -27,6 +27,6 @@ - + \ No newline at end of file diff --git a/SocketHub/Program.cs b/SocketHub/Program.cs index 2f10ae0..14c304d 100644 --- a/SocketHub/Program.cs +++ b/SocketHub/Program.cs @@ -1,22 +1,28 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; +using MQTTnet; using NLog.Extensions.Hosting; -using SuperSocket.ProtoBase; -using SuperSocket.Server; +using SocketHub.Endpoints; +using SocketHub.Filter; +using SocketHub.Handler; +using SocketHub.Models; using SuperSocket.Server.Abstractions; -using SuperSocket.Server.Abstractions.Session; using SuperSocket.Server.Host; -using SuperSocket.WebSocket; using SuperSocket.WebSocket.Server; -using System.Buffers; var host = Host.CreateDefaultBuilder(args) .ConfigureServices((context, services) => { services.AddSingleton, TcpPackageHandler>(); services.AddSingleton(); + + // 1. 绑定 MQTT 配置 + services.Configure(context.Configuration.GetSection("Mqtt")); + // 2. 注册 MQTT 客户端 + // services.AddSingleton(sp => new MqttFactory().CreateMqttClient()); + // 3. 注册 MQTT 后台服务 + services.AddHostedService(); }) .AsMultipleServerHostBuilder() .AddServer(builder => @@ -50,130 +56,3 @@ var host = Host.CreateDefaultBuilder(args) .Build(); await host.RunAsync(); - -public class TcpService : SuperSocketService -{ - public TcpService(IServiceProvider serviceProvider, IOptions serverOptions) - : base(serviceProvider, serverOptions) - { - - } -} - -public class TcpPackageHandler : IPackageHandler -{ - private readonly ILogger _logger; - - public TcpPackageHandler(ILogger logger) - { - _logger = logger; - } - - public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken) - { - _logger.LogInformation($"Magic={package.Magic:X4}, Type={package.Type}, Len={package.Length}"); - await session.SendAsync(package.Payload, cancellationToken); - } -} - -public class WebSocketMessageHandler -{ - private readonly ILogger _logger; - - public WebSocketMessageHandler(ILogger logger) - { - _logger = logger; - } - - public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package) - { - _logger.LogInformation($"[WebSocket] {package.Message}"); - await session.SendAsync("ok"); - } -} - -/// -/// 协议数据结构 示例:65 6D 00 05 00 01 68 65 6C 6C 6F -/// -/// 帧头标识 -/// 数据长度 -/// 数据类型 -/// 数据内容 -/// -public record ProtocolFrame( - ushort Magic, - ushort Length, - ushort Type, - byte[] Payload -); - -/// -/// 二进制解析器: 把 TCP 字节流 → 拆包 → 转成 ProtocolFrame -/// -public class BinaryPipelineFilter : FixedHeaderPipelineFilter -{ - /// - /// 固定头长度 - /// - public BinaryPipelineFilter() : base(6) - { - } - - /// - /// 从包头中解析出 Body 长度 - /// - /// 字节流 - /// Body 长度 - protected override int GetBodyLengthFromHeader(ref ReadOnlySequence buffer) - { - var reader = new SequenceReader(buffer); - - // 读取前2字节 → magic - reader.TryReadBigEndian(out ushort magic); - - // 再读2字节 → Length - reader.TryReadBigEndian(out ushort length); - - // 再读2字节 → Type - reader.TryReadBigEndian(out ushort type); - - // 校验帧头 - if (magic != 0x656D) - { - throw new Exception("非法帧头"); - } - - // 限制长度 - if (length == 0 || length > 8192) - { - throw new Exception("非法长度"); - } - - // 返回 Payload 长度 - return length; - } - - /// - /// 把完整字节包 → 转换成 ProtocolFrame - /// - /// 字节流 - /// ProtocolFrame - protected override ProtocolFrame DecodePackage(ref ReadOnlySequence buffer) - { - var reader = new SequenceReader(buffer); - - // 读取前2字节 → magic - reader.TryReadBigEndian(out ushort magic); - - // 再读2字节 → Length - reader.TryReadBigEndian(out ushort length); - - // 再读2字节 → Type - reader.TryReadBigEndian(out ushort type); - - var payload = buffer.Slice(6, length).ToArray(); - - // 构造ProtocolFrame - return new ProtocolFrame(magic, length, type, payload); - } -} \ No newline at end of file diff --git a/SocketHub/SocketHub.csproj b/SocketHub/SocketHub.csproj index 1cff380..a7066a9 100644 --- a/SocketHub/SocketHub.csproj +++ b/SocketHub/SocketHub.csproj @@ -8,6 +8,7 @@ + diff --git a/SocketHub/appsettings.json b/SocketHub/appsettings.json index 4cd0fee..abdabc2 100644 --- a/SocketHub/appsettings.json +++ b/SocketHub/appsettings.json @@ -18,5 +18,14 @@ } ] } + }, + "Mqtt": { + "Host": "127.0.0.1", + "Port": 1883, + "ClientId": "NetCoreMqttService", + "Topics": [ + "test/topic1", + "test/topic2" + ] } } \ No newline at end of file