Files
MassSpectrometer/IOT.Services/MqttClientWrapper.cs
T
2026-09-09 18:42:38 +08:00

381 lines
12 KiB
C#

using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
using PLCCommunication;
using PLCCommunication.LogNet;
using PLCCommunication.MQTT;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using System.Xml.Linq;
namespace WindowsService1
{
public class MqttClientWrapper
{
private MqttClient mqttClient;
private bool mqttIsConnected;
private const int MaxReconnectAttempts = 5;
private const int ReconnectDelay = 5000; // 5 seconds
private MqttConnectionOptions options;
public event Action<string, string, string> OnMessageReceived;
public long ReceiveCount { get; private set; }
public PLCCommunication.LogNet.ILogNet LogNet = new LogNetDateTime("D:\\APILog\\Logs\\MQTTLogs", GenerateMode.ByEveryDay, 20);
public enum MessageFormat
{
Binary,
Text,
Xml,
Json
}
public MqttClientWrapper(string ipAddress, int port, string clientId, MqttCredential credentials)
{
InitializeMqttOptions(ipAddress, port, clientId, credentials);
}
private void InitializeMqttOptions(string ipAddress, int port, string clientId, MqttCredential credentials)
{
options = new MqttConnectionOptions
{
IpAddress = ipAddress,
Port = port,
ClientId = clientId,
Credentials = credentials,
UseRSAProvider = false,
CleanSession = false,
UseSSL = false,
SSLSecure = false,
CertificateFile = "",
KeepAlivePeriod = TimeSpan.FromSeconds(30),
KeepAliveSendInterval = TimeSpan.FromSeconds(10),
ConnectTimeout = 5,
WillMessage = null
};
}
public async Task InitializeAndConnectAsync()
{
try
{
mqttClient = new MqttClient(options);
Logger.WriteMqtt("MQTT 初始化成功");
await ConnectMqttAsync();
}
catch (Exception ex)
{
Logger.WriteError($"MQTT 初始化失败: {ex.Message}");
}
}
private async Task ConnectMqttAsync()
{
try
{
JYResult connect = await mqttClient.ConnectServerAsync();
if (connect.IsSuccess)
{
mqttIsConnected = true;
Logger.WriteMqtt("MQTT 连接服务器成功");
}
else
{
Logger.WriteError("MQTT 无法连接到服务器");
}
}
catch (Exception ex)
{
Logger.WriteError($"MQTT 无法连接到服务器: {ex.Message}");
}
}
public async Task SendMessageAsync(string source, string publishTopic, string type, JObject messageObj)
{
if (!mqttClient.IsConnected)
{
mqttIsConnected = false;
Logger.WriteMqtt("客户端未连接到服务器。正在尝试重新连接...");
await ReconnectWithRetry();
if (!mqttClient.IsConnected)
{
throw new InvalidOperationException("无法与 MQTT 代理建立连接.");
}
}
try
{
if (messageObj != null)
{
string json = messageObj.ToString();
//string json = JsonConvert.SerializeObject(messageObj);
JYResult result = await mqttClient.PublishMessageAsync(new MqttApplicationMessage
{
Topic = publishTopic,
QualityOfServiceLevel = MqttQualityOfServiceLevel.AtMostOnce,
Payload = Encoding.UTF8.GetBytes(json),
Retain = false,
});
if (result.IsSuccess)
{
if (Service1.SwitchLog)
{
Logger.WriteMqtt($"线程 [{source}] 发送 [{type}] 消息到主题 [{publishTopic}]: {json}");
}
}
else
{
Logger.WriteMqtt($"线程 [{source}] 发送 [{type}] 消息到主题 [{publishTopic}] 失败: {result.Message}");
}
}
}
catch (Exception ex)
{
Logger.WriteError($"发送消息错误: {ex.Message}");
}
}
public async Task SendMessageAsync(string source, string publishTopic, string type, string messageObj)
{
if (!mqttClient.IsConnected)
{
mqttIsConnected = false;
Logger.WriteMqtt("客户端未连接到服务器。正在尝试重新连接...");
await ReconnectWithRetry();
if (!mqttClient.IsConnected)
{
throw new InvalidOperationException("无法与 MQTT 代理建立连接.");
}
}
try
{
if (messageObj != null)
{
JYResult result = await mqttClient.PublishMessageAsync(new MqttApplicationMessage
{
Topic = publishTopic,
QualityOfServiceLevel = MqttQualityOfServiceLevel.AtMostOnce,
Payload = Encoding.UTF8.GetBytes(messageObj),
Retain = false,
});
if (result.IsSuccess)
{
if (Service1.SwitchLog)
{
Logger.WriteMqtt($"线程 [{source}] 发送 [{type}] 消息到主题 [{publishTopic}]: {messageObj}");
}
}
else
{
Logger.WriteMqtt($"线程 [{source}] 发送 [{type}] 消息到主题 [{publishTopic}] 失败: {result.Message}");
}
}
}
catch (Exception ex)
{
Logger.WriteError($"发送消息错误: {ex.Message}");
}
}
public void SubscribeMessage(string topic, MessageFormat format = MessageFormat.Text, bool truncateLongMessage = false)
{
if (!mqttIsConnected)
{
ConnectMqttAsync().Wait();
}
try
{
JYResult operateResult = mqttClient.SubscribeMessage(new string[] { topic });
if (!operateResult.IsSuccess)
{
Logger.WriteMqtt($"订阅消息失败: {operateResult.Message}");
return;
}
SubscribeTopic subscribeTopic = mqttClient.GetSubscribeTopic(topic);
if (subscribeTopic != null)
{
subscribeTopic.OnMqttMessageReceived += (sender, args) =>
{
HandleReceivedMessage(args.Topic, args.Payload, format, truncateLongMessage);
};
Logger.WriteMqtt($"成功订阅主题: {topic}");
}
else
{
Logger.WriteMqtt($"无法获取订阅主题 {topic}");
}
}
catch (Exception ex)
{
Logger.WriteError($"订阅主题时出错 {topic}: {ex.Message}");
}
}
private void HandleReceivedMessage(string topic, byte[] payload, MessageFormat format, bool truncateLongMessage)
{
ReceiveCount++;
string text = string.Empty;
try
{
switch (format)
{
case MessageFormat.Binary:
text = BitConverter.ToString(payload).Replace("-", " ");
break;
case MessageFormat.Text:
text = Encoding.UTF8.GetString(payload);
break;
case MessageFormat.Xml:
try
{
text = XElement.Parse(Encoding.UTF8.GetString(payload)).ToString();
}
catch
{
text = Encoding.UTF8.GetString(payload);
}
break;
case MessageFormat.Json:
try
{
text = JObject.Parse(Encoding.UTF8.GetString(payload)).ToString();
}
catch
{
text = Encoding.UTF8.GetString(payload);
}
break;
}
//if (Service1.SwitchLog)
//{
// LogNet.WriteDebug($"收到主题 {topic}: {text}");
//}
//if (truncateLongMessage && text.Length > 200)
//{
// text = text.Substring(0, 200) + "...";
//}
string timestamp = DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss.fff");
OnMessageReceived?.Invoke(timestamp, topic, text);
}
catch (Exception ex)
{
Logger.WriteError($"处理收到的消息时出错: {ex.Message}");
}
}
public void UnsubscribeMessage(string topic)
{
try
{
JYResult operateResult = mqttClient.UnSubscribeMessage(new string[] { topic });
if (operateResult.IsSuccess)
{
Logger.WriteMqtt($"成功取消订阅主题: {topic}");
}
else
{
Logger.WriteMqtt($"取消订阅主题失败 {topic}: {operateResult.Message}");
}
}
catch (Exception ex)
{
Logger.WriteError($"取消订阅主题时出错 {topic}: {ex.Message}");
}
}
/// <summary>
/// 重连
/// </summary>
/// <returns></returns>
private async Task ReconnectWithRetry()
{
int attempts = 0;
while (!mqttClient.IsConnected && attempts < MaxReconnectAttempts)
{
try
{
await mqttClient.ConnectServerAsync();
Logger.WriteMqtt("重新连接成功.");
return;
}
catch (Exception ex)
{
attempts++;
Logger.WriteMqtt($"重新连接尝试 {attempts} 失败的: {ex.Message}");
await Task.Delay(ReconnectDelay);
}
}
if (!mqttClient.IsConnected)
{
Logger.WriteMqtt("尝试次数达到上限后仍无法重新连接.");
}
}
private async Task ReconnectAsync()
{
Logger.WriteMqtt("正在尝试重新连接...");
mqttIsConnected = false;
int retryCount = 0;
const int maxRetries = 5;
const int retryDelayMs = 5000;
while (!mqttIsConnected && retryCount < maxRetries)
{
try
{
await ConnectMqttAsync();
if (mqttIsConnected)
{
Logger.WriteMqtt("重新连接成功");
return;
}
}
catch (Exception ex)
{
Logger.WriteMqtt($"重新连接尝试 {retryCount + 1} 失败的: {ex.Message}");
}
retryCount++;
await Task.Delay(retryDelayMs);
}
if (!mqttIsConnected)
{
Logger.WriteMqtt("多次尝试后仍无法重新连接");
}
}
public async Task Disconnect()
{
await mqttClient.ConnectCloseAsync();
mqttIsConnected = false;
Logger.WriteMqtt("与 MQTT服务端断开连接!!!");
}
}
}