From f06d0139dc4b3b7a549d2e9cf63d18f2057ec23b Mon Sep 17 00:00:00 2001 From: Cxx0822 <1556464090@qq.com> Date: Wed, 22 Apr 2026 22:56:41 +0800 Subject: [PATCH] =?UTF-8?q?feat:=E6=9B=B4=E6=96=B0Mqtt=E5=AE=A2=E6=88=B7?= =?UTF-8?q?=E7=AB=AF=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- SocketHub/Endpoints/MqttService.cs | 60 ++++++++++++++----- SocketHub/Endpoints/TcpService.cs | 7 +-- SocketHub/Handlers/TcpPackageHandler.cs | 9 +-- SocketHub/Handlers/WebSocketMessageHandler.cs | 9 +-- SocketHub/Models/MqttSettings.cs | 10 +--- SocketHub/Program.cs | 20 ++----- 6 files changed, 56 insertions(+), 59 deletions(-) diff --git a/SocketHub/Endpoints/MqttService.cs b/SocketHub/Endpoints/MqttService.cs index 9c3fd9f..54512b4 100644 --- a/SocketHub/Endpoints/MqttService.cs +++ b/SocketHub/Endpoints/MqttService.cs @@ -1,25 +1,21 @@ using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using MQTTnet; using MQTTnet.Protocol; +using SocketHub.Handler; using SocketHub.Models; using System.Text; namespace SocketHub.Endpoints { - internal class MqttService : BackgroundService + internal class MqttService(IMqttClient mqttClient, IOptions settings, ILogger logger) : BackgroundService { - private readonly IMqttClient _mqttClient; - private readonly MqttSettings _settings; + private readonly IMqttClient _mqttClient = mqttClient; + private readonly MqttSettings _settings = settings.Value; + private readonly ILogger _logger = logger; private MqttClientOptions? _mqttOptions; - // 直接注入配置 + MQTT 客户端 - public MqttService(IMqttClient mqttClient, IOptions settings) - { - _mqttClient = mqttClient; - _settings = settings.Value; - } - protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // 构建连接配置 @@ -29,37 +25,69 @@ namespace SocketHub.Endpoints .WithCleanSession() .Build(); - _mqttClient.ConnectedAsync += ConnectedAsync; + _mqttClient.ConnectedAsync += OnConnectedAsync; + _mqttClient.DisconnectedAsync += OnDisconnectedAsync; _mqttClient.ApplicationMessageReceivedAsync += HandleMessage; + _logger.LogInformation("[MQTT] 正在连接到服务器 {Host}:{Port}...", _settings.Host, _settings.Port); + + await _mqttClient.ConnectAsync(_mqttOptions, stoppingToken); + // 等待程序停止 await Task.Delay(Timeout.Infinite, stoppingToken); } - // 连接成功事件 - public async Task ConnectedAsync(MqttClientConnectedEventArgs arg) + public async Task OnConnectedAsync(MqttClientConnectedEventArgs arg) { + _logger.LogInformation("[MQTT] 已成功连接到服务器"); + // 连接成功后订阅主题 foreach (var topic in _settings.Topics) { await _mqttClient.SubscribeAsync(topic, MqttQualityOfServiceLevel.AtLeastOnce); - Console.WriteLine($"[MQTT] 已订阅主题:{topic}"); + _logger.LogInformation($"[MQTT] 已订阅主题:{topic}"); } await Task.CompletedTask; } + public async Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs arg) + { + _logger.LogError("[MQTT] 连接失败,原因:{Reason}", arg.Reason); + + 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"); + _logger.LogInformation($"\n[MQTT] 收到消息\n主题:{topic}\n内容:{payload}\n"); await Task.CompletedTask; } - // 优雅停止 + public async Task PublishAsync(string topic, string payload, CancellationToken cancellationToken = default) + { + try + { + var message = new MqttApplicationMessageBuilder() + .WithTopic(topic) + .WithPayload(payload) + .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce) + .Build(); + + await _mqttClient.PublishAsync(message, cancellationToken); + return true; + } + catch (Exception ex) + { + _logger.LogError(ex, "[MQTT] 发布消息到主题 {Topic} 失败", topic); + return false; + } + } + public override async Task StopAsync(CancellationToken stoppingToken) { await _mqttClient.DisconnectAsync(); diff --git a/SocketHub/Endpoints/TcpService.cs b/SocketHub/Endpoints/TcpService.cs index ef2a606..011d4bb 100644 --- a/SocketHub/Endpoints/TcpService.cs +++ b/SocketHub/Endpoints/TcpService.cs @@ -5,12 +5,7 @@ using SuperSocket.Server.Abstractions; namespace SocketHub.Endpoints { - internal class TcpService : SuperSocketService + internal class TcpService(IServiceProvider serviceProvider, IOptions serverOptions) : SuperSocketService(serviceProvider, serverOptions) { - public TcpService(IServiceProvider serviceProvider, IOptions serverOptions) - : base(serviceProvider, serverOptions) - { - - } } } diff --git a/SocketHub/Handlers/TcpPackageHandler.cs b/SocketHub/Handlers/TcpPackageHandler.cs index 3a4dfe0..dfe303f 100644 --- a/SocketHub/Handlers/TcpPackageHandler.cs +++ b/SocketHub/Handlers/TcpPackageHandler.cs @@ -5,14 +5,9 @@ using SuperSocket.Server.Abstractions.Session; namespace SocketHub.Handler { - internal class TcpPackageHandler : IPackageHandler + internal class TcpPackageHandler(ILogger logger) : IPackageHandler { - private readonly ILogger _logger; - - public TcpPackageHandler(ILogger logger) - { - _logger = logger; - } + private readonly ILogger _logger = logger; public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken) { diff --git a/SocketHub/Handlers/WebSocketMessageHandler.cs b/SocketHub/Handlers/WebSocketMessageHandler.cs index 593c4e9..cb47ea5 100644 --- a/SocketHub/Handlers/WebSocketMessageHandler.cs +++ b/SocketHub/Handlers/WebSocketMessageHandler.cs @@ -4,14 +4,9 @@ using SuperSocket.WebSocket.Server; namespace SocketHub.Handler { - internal class WebSocketMessageHandler + internal class WebSocketMessageHandler(ILogger logger) { - private readonly ILogger _logger; - - public WebSocketMessageHandler(ILogger logger) - { - _logger = logger; - } + private readonly ILogger _logger = logger; public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package) { diff --git a/SocketHub/Models/MqttSettings.cs b/SocketHub/Models/MqttSettings.cs index bb66ceb..fc7bcdf 100644 --- a/SocketHub/Models/MqttSettings.cs +++ b/SocketHub/Models/MqttSettings.cs @@ -1,16 +1,10 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; -using System.Threading.Tasks; - -namespace SocketHub.Models +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(); + public List Topics { get; set; } = []; } } diff --git a/SocketHub/Program.cs b/SocketHub/Program.cs index 14c304d..163a2a0 100644 --- a/SocketHub/Program.cs +++ b/SocketHub/Program.cs @@ -17,21 +17,17 @@ var host = Host.CreateDefaultBuilder(args) services.AddSingleton, TcpPackageHandler>(); services.AddSingleton(); - // 1. 绑定 MQTT 配置 + // 1. 注入 MQTT 配置 services.Configure(context.Configuration.GetSection("Mqtt")); // 2. 注册 MQTT 客户端 - // services.AddSingleton(sp => new MqttFactory().CreateMqttClient()); + services.AddSingleton(serviceProvider => new MqttClientFactory().CreateMqttClient()); // 3. 注册 MQTT 后台服务 services.AddHostedService(); }) .AsMultipleServerHostBuilder() .AddServer(builder => { - builder - .ConfigureServerOptions((ctx, config) => - { - return config.GetSection("TcpServer"); - }); + builder.ConfigureServerOptions((ctx, config) => config.GetSection("TcpServer")); }) .AddWebSocketServer(builder => { @@ -43,15 +39,9 @@ var host = Host.CreateDefaultBuilder(args) await handler.HandleAsync(session, package); }) - .ConfigureServerOptions((ctx, config) => - { - return config.GetSection("WebSocketServer"); - }); - }) - .ConfigureLogging(logging => - { - logging.ClearProviders(); + .ConfigureServerOptions((ctx, config) => config.GetSection("WebSocketServer")); }) + .ConfigureLogging(logging => logging.ClearProviders()) .UseNLog() .Build();