using MQTTnet; using MQTTnet.Client; using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading.Tasks; using MQTTnet; using MQTTnet.Client; using MQTTnet.Protocol; namespace JY.Inspection.Common { public class MqttSingleton { // 单例实例 private static readonly Lazy _instance = new Lazy(() => new MqttSingleton()); public static MqttSingleton Instance => _instance.Value; // 内部成员 private IMqttClient _client; private MqttClientOptions _options; private bool _isInitialized = false; // 事件:收到消息时触发 public event Action MessageReceived; // 私有构造函数 private MqttSingleton() { } /// /// 初始化 MQTT 客户端(只需调用一次) /// /// MQTT 服务端地址 /// 端口,默认1883 /// 客户端ID,为空则自动生成 /// 用户名(可选) /// 密码(可选) /// 是否启用TLS加密 public async Task InitializeAsync( string brokerHost, int port = 1883, string clientId = null, string username = null, string password = null, bool useTls = false) { if (_isInitialized) return; // 创建客户端 var factory = new MqttFactory(); _client = factory.CreateMqttClient(); // 构建连接选项 var optionsBuilder = new MqttClientOptionsBuilder() .WithTcpServer(brokerHost, port) .WithClientId(string.IsNullOrEmpty(clientId) ? $"WinForm_{Guid.NewGuid()}" : clientId) .WithCleanSession(); if (!string.IsNullOrEmpty(username)) { optionsBuilder.WithCredentials(username, password); } if (useTls) { optionsBuilder.WithTls(); } _options = optionsBuilder.Build(); // 绑定事件 _client.ConnectedAsync += async e => { Console.WriteLine("MQTT 已连接"); }; _client.DisconnectedAsync += async e => { Console.WriteLine("MQTT 已断开,尝试重连..."); // 自动重连 await Task.Delay(3000).ContinueWith(_ => ConnectAsync()); }; _client.ApplicationMessageReceivedAsync += async e => { var topic = e.ApplicationMessage.Topic; var payload = e.ApplicationMessage.PayloadSegment == null ? null : Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment.Array); MessageReceived?.Invoke(topic, payload); }; // 连接 await ConnectAsync(); _isInitialized = true; } /// /// 连接MQTT服务端 /// private async Task ConnectAsync() { try { if (!_client.IsConnected) { await _client.ConnectAsync(_options, System.Threading.CancellationToken.None); } } catch (Exception ex) { Console.WriteLine($"MQTT连接失败: {ex.Message}"); } } /// /// 发送消息 /// /// 主题 /// 消息内容 /// QoS等级,默认1 public async Task SendAsync(string topic, string payload, int qos = 1) { if (!_client.IsConnected) { await ConnectAsync(); } if (_client.IsConnected) { var message = new MqttApplicationMessageBuilder() .WithTopic(topic) .WithPayload(payload) .WithQualityOfServiceLevel((MqttQualityOfServiceLevel)qos) .WithRetainFlag(false) .Build(); await _client.PublishAsync(message, System.Threading.CancellationToken.None); } } /// /// 订阅主题 /// public async Task SubscribeAsync(string topic) { if (_client.IsConnected) { await _client.SubscribeAsync(new MqttTopicFilterBuilder() .WithTopic(topic) .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce) .Build()); } } /// /// 断开连接 /// public async Task DisconnectAsync() { if (_client != null && _client.IsConnected) { await _client.DisconnectAsync(); } } /// /// 当前是否已连接 /// public bool IsConnected => _client?.IsConnected ?? false; } }