feat:增加Mqtt客户端

This commit is contained in:
2026-04-22 20:01:55 +08:00
parent 964bbcb51d
commit e8983f2e95
11 changed files with 266 additions and 134 deletions

View File

@@ -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<MqttSettings> 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);
}
}
}

View File

@@ -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<ProtocolFrame>
{
public TcpService(IServiceProvider serviceProvider, IOptions<ServerOptions> serverOptions)
: base(serviceProvider, serverOptions)
{
}
}
}

View File

@@ -0,0 +1,74 @@
using SocketHub.Models;
using SuperSocket.ProtoBase;
using System.Buffers;
namespace SocketHub.Filter
{
internal class BinaryPipelineFilter : FixedHeaderPipelineFilter<ProtocolFrame>
{
/// <summary>
/// 固定头长度
/// </summary>
public BinaryPipelineFilter() : base(6)
{
}
/// <summary>
/// 从包头中解析出 Body 长度
/// </summary>
/// <param name="buffer">字节流</param>
/// <returns>Body 长度</returns>
protected override int GetBodyLengthFromHeader(ref ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(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;
}
/// <summary>
/// 把完整字节包 → 转换成 ProtocolFrame
/// </summary>
/// <param name="buffer">字节流</param>
/// <returns>ProtocolFrame</returns>
protected override ProtocolFrame DecodePackage(ref ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(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);
}
}
}

View File

@@ -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<ProtocolFrame>
{
private readonly ILogger<TcpPackageHandler> _logger;
public TcpPackageHandler(ILogger<TcpPackageHandler> 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);
}
}
}

View File

@@ -0,0 +1,22 @@
using Microsoft.Extensions.Logging;
using SuperSocket.WebSocket;
using SuperSocket.WebSocket.Server;
namespace SocketHub.Handler
{
internal class WebSocketMessageHandler
{
private readonly ILogger<WebSocketMessageHandler> _logger;
public WebSocketMessageHandler(ILogger<WebSocketMessageHandler> logger)
{
_logger = logger;
}
public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package)
{
_logger.LogInformation($"[WebSocket] {package.Message}");
await session.SendAsync("ok");
}
}
}

View File

@@ -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<string> Topics { get; set; } = new();
}
}

View File

@@ -0,0 +1,23 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
namespace SocketHub.Models
{
/// <summary>
/// 协议数据结构 示例65 6D 00 05 00 01 68 65 6C 6C 6F
/// </summary>
/// <param name="Magic">帧头标识</param>
/// <param name="Length">数据长度</param>
/// <param name="Type">数据类型</param>
/// <param name="Payload">数据内容</param>
///
internal record ProtocolFrame(
ushort Magic,
ushort Length,
ushort Type,
byte[] Payload
);
}

View File

@@ -27,6 +27,6 @@
</targets> </targets>
<rules> <rules>
<logger name="*" minlevel="Info" writeTo="stdout"/> <logger name="*" minlevel="Info" writeTo="stdout, fluentbit"/>
</rules> </rules>
</nlog> </nlog>

View File

@@ -1,22 +1,28 @@
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options; using MQTTnet;
using NLog.Extensions.Hosting; using NLog.Extensions.Hosting;
using SuperSocket.ProtoBase; using SocketHub.Endpoints;
using SuperSocket.Server; using SocketHub.Filter;
using SocketHub.Handler;
using SocketHub.Models;
using SuperSocket.Server.Abstractions; using SuperSocket.Server.Abstractions;
using SuperSocket.Server.Abstractions.Session;
using SuperSocket.Server.Host; using SuperSocket.Server.Host;
using SuperSocket.WebSocket;
using SuperSocket.WebSocket.Server; using SuperSocket.WebSocket.Server;
using System.Buffers;
var host = Host.CreateDefaultBuilder(args) var host = Host.CreateDefaultBuilder(args)
.ConfigureServices((context, services) => .ConfigureServices((context, services) =>
{ {
services.AddSingleton<IPackageHandler<ProtocolFrame>, TcpPackageHandler>(); services.AddSingleton<IPackageHandler<ProtocolFrame>, TcpPackageHandler>();
services.AddSingleton<WebSocketMessageHandler>(); services.AddSingleton<WebSocketMessageHandler>();
// 1. 绑定 MQTT 配置
services.Configure<MqttSettings>(context.Configuration.GetSection("Mqtt"));
// 2. 注册 MQTT 客户端
// services.AddSingleton<IMqttClient>(sp => new MqttFactory().CreateMqttClient());
// 3. 注册 MQTT 后台服务
services.AddHostedService<MqttService>();
}) })
.AsMultipleServerHostBuilder() .AsMultipleServerHostBuilder()
.AddServer<TcpService, ProtocolFrame, BinaryPipelineFilter>(builder => .AddServer<TcpService, ProtocolFrame, BinaryPipelineFilter>(builder =>
@@ -50,130 +56,3 @@ var host = Host.CreateDefaultBuilder(args)
.Build(); .Build();
await host.RunAsync(); await host.RunAsync();
public class TcpService : SuperSocketService<ProtocolFrame>
{
public TcpService(IServiceProvider serviceProvider, IOptions<ServerOptions> serverOptions)
: base(serviceProvider, serverOptions)
{
}
}
public class TcpPackageHandler : IPackageHandler<ProtocolFrame>
{
private readonly ILogger<TcpPackageHandler> _logger;
public TcpPackageHandler(ILogger<TcpPackageHandler> 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<WebSocketMessageHandler> _logger;
public WebSocketMessageHandler(ILogger<WebSocketMessageHandler> logger)
{
_logger = logger;
}
public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package)
{
_logger.LogInformation($"[WebSocket] {package.Message}");
await session.SendAsync("ok");
}
}
/// <summary>
/// 协议数据结构 示例65 6D 00 05 00 01 68 65 6C 6C 6F
/// </summary>
/// <param name="Magic">帧头标识</param>
/// <param name="Length">数据长度</param>
/// <param name="Type">数据类型</param>
/// <param name="Payload">数据内容</param>
///
public record ProtocolFrame(
ushort Magic,
ushort Length,
ushort Type,
byte[] Payload
);
/// <summary>
/// 二进制解析器: 把 TCP 字节流 → 拆包 → 转成 ProtocolFrame
/// </summary>
public class BinaryPipelineFilter : FixedHeaderPipelineFilter<ProtocolFrame>
{
/// <summary>
/// 固定头长度
/// </summary>
public BinaryPipelineFilter() : base(6)
{
}
/// <summary>
/// 从包头中解析出 Body 长度
/// </summary>
/// <param name="buffer">字节流</param>
/// <returns>Body 长度</returns>
protected override int GetBodyLengthFromHeader(ref ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(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;
}
/// <summary>
/// 把完整字节包 → 转换成 ProtocolFrame
/// </summary>
/// <param name="buffer">字节流</param>
/// <returns>ProtocolFrame</returns>
protected override ProtocolFrame DecodePackage(ref ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(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);
}
}

View File

@@ -8,6 +8,7 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="MQTTnet" Version="5.1.0.1559" />
<PackageReference Include="NLog.Extensions.Hosting" Version="6.1.2" /> <PackageReference Include="NLog.Extensions.Hosting" Version="6.1.2" />
<PackageReference Include="NLog.Targets.Network" Version="6.0.4" /> <PackageReference Include="NLog.Targets.Network" Version="6.0.4" />
<PackageReference Include="SuperSocket" Version="2.0.2" /> <PackageReference Include="SuperSocket" Version="2.0.2" />

View File

@@ -18,5 +18,14 @@
} }
] ]
} }
},
"Mqtt": {
"Host": "127.0.0.1",
"Port": 1883,
"ClientId": "NetCoreMqttService",
"Topics": [
"test/topic1",
"test/topic2"
]
} }
} }