feat:更新Mqtt客户端功能

This commit is contained in:
2026-04-22 22:56:41 +08:00
parent e8983f2e95
commit f06d0139dc
6 changed files with 56 additions and 59 deletions

View File

@@ -1,25 +1,21 @@
using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
using MQTTnet; using MQTTnet;
using MQTTnet.Protocol; using MQTTnet.Protocol;
using SocketHub.Handler;
using SocketHub.Models; using SocketHub.Models;
using System.Text; using System.Text;
namespace SocketHub.Endpoints namespace SocketHub.Endpoints
{ {
internal class MqttService : BackgroundService internal class MqttService(IMqttClient mqttClient, IOptions<MqttSettings> settings, ILogger<TcpPackageHandler> logger) : BackgroundService
{ {
private readonly IMqttClient _mqttClient; private readonly IMqttClient _mqttClient = mqttClient;
private readonly MqttSettings _settings; private readonly MqttSettings _settings = settings.Value;
private readonly ILogger<TcpPackageHandler> _logger = logger;
private MqttClientOptions? _mqttOptions; private MqttClientOptions? _mqttOptions;
// 直接注入配置 + MQTT 客户端
public MqttService(IMqttClient mqttClient, IOptions<MqttSettings> settings)
{
_mqttClient = mqttClient;
_settings = settings.Value;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken) protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{ {
// 构建连接配置 // 构建连接配置
@@ -29,37 +25,69 @@ namespace SocketHub.Endpoints
.WithCleanSession() .WithCleanSession()
.Build(); .Build();
_mqttClient.ConnectedAsync += ConnectedAsync; _mqttClient.ConnectedAsync += OnConnectedAsync;
_mqttClient.DisconnectedAsync += OnDisconnectedAsync;
_mqttClient.ApplicationMessageReceivedAsync += HandleMessage; _mqttClient.ApplicationMessageReceivedAsync += HandleMessage;
_logger.LogInformation("[MQTT] 正在连接到服务器 {Host}:{Port}...", _settings.Host, _settings.Port);
await _mqttClient.ConnectAsync(_mqttOptions, stoppingToken);
// 等待程序停止 // 等待程序停止
await Task.Delay(Timeout.Infinite, stoppingToken); await Task.Delay(Timeout.Infinite, stoppingToken);
} }
// 连接成功事件 public async Task OnConnectedAsync(MqttClientConnectedEventArgs arg)
public async Task ConnectedAsync(MqttClientConnectedEventArgs arg)
{ {
_logger.LogInformation("[MQTT] 已成功连接到服务器");
// 连接成功后订阅主题 // 连接成功后订阅主题
foreach (var topic in _settings.Topics) foreach (var topic in _settings.Topics)
{ {
await _mqttClient.SubscribeAsync(topic, MqttQualityOfServiceLevel.AtLeastOnce); await _mqttClient.SubscribeAsync(topic, MqttQualityOfServiceLevel.AtLeastOnce);
Console.WriteLine($"[MQTT] 已订阅主题:{topic}"); _logger.LogInformation($"[MQTT] 已订阅主题:{topic}");
} }
await Task.CompletedTask; await Task.CompletedTask;
} }
public async Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs arg)
{
_logger.LogError("[MQTT] 连接失败,原因:{Reason}", arg.Reason);
await Task.CompletedTask;
}
public async Task HandleMessage(MqttApplicationMessageReceivedEventArgs arg) public async Task HandleMessage(MqttApplicationMessageReceivedEventArgs arg)
{ {
var topic = arg.ApplicationMessage.Topic; var topic = arg.ApplicationMessage.Topic;
var payload = Encoding.UTF8.GetString(arg.ApplicationMessage.Payload); 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; await Task.CompletedTask;
} }
// 优雅停止 public async Task<bool> 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) public override async Task StopAsync(CancellationToken stoppingToken)
{ {
await _mqttClient.DisconnectAsync(); await _mqttClient.DisconnectAsync();

View File

@@ -5,12 +5,7 @@ using SuperSocket.Server.Abstractions;
namespace SocketHub.Endpoints namespace SocketHub.Endpoints
{ {
internal class TcpService : SuperSocketService<ProtocolFrame> internal class TcpService(IServiceProvider serviceProvider, IOptions<ServerOptions> serverOptions) : SuperSocketService<ProtocolFrame>(serviceProvider, serverOptions)
{ {
public TcpService(IServiceProvider serviceProvider, IOptions<ServerOptions> serverOptions)
: base(serviceProvider, serverOptions)
{
}
} }
} }

View File

@@ -5,14 +5,9 @@ using SuperSocket.Server.Abstractions.Session;
namespace SocketHub.Handler namespace SocketHub.Handler
{ {
internal class TcpPackageHandler : IPackageHandler<ProtocolFrame> internal class TcpPackageHandler(ILogger<TcpPackageHandler> logger) : IPackageHandler<ProtocolFrame>
{ {
private readonly ILogger<TcpPackageHandler> _logger; private readonly ILogger<TcpPackageHandler> _logger = logger;
public TcpPackageHandler(ILogger<TcpPackageHandler> logger)
{
_logger = logger;
}
public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken) public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken)
{ {

View File

@@ -4,14 +4,9 @@ using SuperSocket.WebSocket.Server;
namespace SocketHub.Handler namespace SocketHub.Handler
{ {
internal class WebSocketMessageHandler internal class WebSocketMessageHandler(ILogger<WebSocketMessageHandler> logger)
{ {
private readonly ILogger<WebSocketMessageHandler> _logger; private readonly ILogger<WebSocketMessageHandler> _logger = logger;
public WebSocketMessageHandler(ILogger<WebSocketMessageHandler> logger)
{
_logger = logger;
}
public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package) public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package)
{ {

View File

@@ -1,16 +1,10 @@
using System; namespace SocketHub.Models
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
namespace SocketHub.Models
{ {
internal class MqttSettings internal class MqttSettings
{ {
public string Host { get; set; } = string.Empty; public string Host { get; set; } = string.Empty;
public int Port { get; set; } = 1883; public int Port { get; set; } = 1883;
public string ClientId { get; set; } = string.Empty; public string ClientId { get; set; } = string.Empty;
public List<string> Topics { get; set; } = new(); public List<string> Topics { get; set; } = [];
} }
} }

View File

@@ -17,21 +17,17 @@ var host = Host.CreateDefaultBuilder(args)
services.AddSingleton<IPackageHandler<ProtocolFrame>, TcpPackageHandler>(); services.AddSingleton<IPackageHandler<ProtocolFrame>, TcpPackageHandler>();
services.AddSingleton<WebSocketMessageHandler>(); services.AddSingleton<WebSocketMessageHandler>();
// 1. 绑定 MQTT 配置 // 1. 注入 MQTT 配置
services.Configure<MqttSettings>(context.Configuration.GetSection("Mqtt")); services.Configure<MqttSettings>(context.Configuration.GetSection("Mqtt"));
// 2. 注册 MQTT 客户端 // 2. 注册 MQTT 客户端
// services.AddSingleton<IMqttClient>(sp => new MqttFactory().CreateMqttClient()); services.AddSingleton<IMqttClient>(serviceProvider => new MqttClientFactory().CreateMqttClient());
// 3. 注册 MQTT 后台服务 // 3. 注册 MQTT 后台服务
services.AddHostedService<MqttService>(); services.AddHostedService<MqttService>();
}) })
.AsMultipleServerHostBuilder() .AsMultipleServerHostBuilder()
.AddServer<TcpService, ProtocolFrame, BinaryPipelineFilter>(builder => .AddServer<TcpService, ProtocolFrame, BinaryPipelineFilter>(builder =>
{ {
builder builder.ConfigureServerOptions((ctx, config) => config.GetSection("TcpServer"));
.ConfigureServerOptions((ctx, config) =>
{
return config.GetSection("TcpServer");
});
}) })
.AddWebSocketServer(builder => .AddWebSocketServer(builder =>
{ {
@@ -43,15 +39,9 @@ var host = Host.CreateDefaultBuilder(args)
await handler.HandleAsync(session, package); await handler.HandleAsync(session, package);
}) })
.ConfigureServerOptions((ctx, config) => .ConfigureServerOptions((ctx, config) => config.GetSection("WebSocketServer"));
{
return config.GetSection("WebSocketServer");
});
})
.ConfigureLogging(logging =>
{
logging.ClearProviders();
}) })
.ConfigureLogging(logging => logging.ClearProviders())
.UseNLog() .UseNLog()
.Build(); .Build();