using HslCommunication; using HslCommunication.Core.Net; using HslCommunication.Enthernet; using MQTTnet; using MQTTnet.Client.Connecting; using MQTTnet.Client.Disconnecting; using MQTTnet.Client.Options; using MQTTnet.Client.Receiving; using MQTTnet.Extensions.ManagedClient; using MQTTnet.Formatter; using System; using System.Collections.Generic; using System.Drawing; using System.Linq; using System.Net; using System.Text; using System.Text.RegularExpressions; using System.Threading; using System.Threading.Tasks; /// /// /// namespace MesWork { /// /// /// public partial class Temp { } /// /// /// public partial class MesWorkForm { /// /// /// public static NetComplexClient complexClient = null; /// /// /// static private IManagedMqttClient mqttClient = null; /// /// /// string topic = ""; /// /// /// static bool OnLine_Mqtt = false; /// /// /// /// /// public async Task StartAsync(string uri) { var array = uri.Split(new char[] { ':', ',' }); string mqttUser = array[2]; string mqttPassword = ""; var options = new MqttClientOptions { ClientId = $"{mqttUser}_{Guid.NewGuid().ToString()}", ProtocolVersion = MqttProtocolVersion.V311, ChannelOptions = new MqttClientTcpOptions { Server = array[0], Port = Convert.ToInt32(array[1]) }, CleanSession = true, KeepAlivePeriod = TimeSpan.FromSeconds(15), Credentials = new MqttClientCredentials { Username = mqttUser, Password = Encoding.UTF8.GetBytes(mqttPassword) } }; mqttClient = new MqttFactory().CreateManagedMqttClient(); mqttClient.ConnectedHandler = new MqttClientConnectedHandlerDelegate(OnSubscriberConnected); mqttClient.DisconnectedHandler = new MqttClientDisconnectedHandlerDelegate(OnSubscriberDisconnected); mqttClient.ApplicationMessageReceivedHandler = new MqttApplicationMessageReceivedHandlerDelegate(OnSubscriberMessageReceived); await mqttClient.StartAsync(new ManagedMqttClientOptions { ClientOptions = options }); this.topic = array[2]; await mqttClient.SubscribeAsync(new MqttTopicFilter { Topic = topic }); } /// /// /// /// private void ClientConnection(string MesServer_IpPortName) { string serverIP = MesServer_IpPortName.Split(':')[0]; string portString = MesServer_IpPortName.Split(':')[1].Split(',')[0]; string opNameSocket = MesServer_IpPortName.Split(',')[1]; if (!IPAddress.TryParse(serverIP, out IPAddress address)) { return; } if (!int.TryParse(portString, out int port)) { return; } try { // 连接 connect complexClient = new NetComplexClient(); complexClient.ClientAlias = opNameSocket; complexClient.EndPointServer = new IPEndPoint(address, port); //complexClient.Token = new Guid(textBox3.Text); complexClient.AcceptString += ComplexClient_AcceptString; complexClient.MessageAlerts += ComplexClient_MessageAlerts; complexClient.ClientStart(); } catch (Exception ex) { HslCommunication.BasicFramework.SoftBasic.ShowExceptionMessage(ex); } } /// /// /// /// /// /// private static async Task Publish(string publish_topic, string publish_msg) { if (mqttClient != null) { await mqttClient.PublishAsync(new MqttApplicationMessageBuilder().WithTopic(publish_topic).WithPayload(publish_msg).WithAtMostOnceQoS().Build()); } } /// /// /// /// private async Task DisConnectService() { if (mqttClient != null) { await mqttClient.UnsubscribeAsync(new string[] { topic }); await mqttClient.StopAsync(); } } /// /// /// /// private void ComplexClient_MessageAlerts(string text) { //if (InvokeRequired) //{ // Invoke(new Action(ComplexClient_MessageAlerts), text); // return; //} } /// /// 连接完成 /// /// private void OnSubscriberConnected(MqttClientConnectedEventArgs x) { OnLine_Mqtt = true; // textBox1_robotData.Text += "连接完成" + "\r\n"; var rec = x.AuthenticateResult.ResultCode.ToString(); // label_State.Text = rec; // textBox1_robotData.Text += rec + "\r\n"; } /// /// 断开连接 /// /// private void OnSubscriberDisconnected(MqttClientDisconnectedEventArgs x) { OnLine_Mqtt = false; try { // textBox1_robotData.Text += "断开连接" + "\r\n"; var rec = x.ClientWasConnected.ToString(); // textBox1_robotData.Text += rec + "\r\n"; } catch (Exception err) { // textBox1_robotData.Text += $"Message:{err.Message}\r\n\r\nStackTrace:{err.StackTrace}" + "\r\n"; } } /// /// 接收信息 /// /// private void OnSubscriberMessageReceived(MqttApplicationMessageReceivedEventArgs x) { var recStr = x.ApplicationMessage.ConvertPayloadToString(); ThreadPool.QueueUserWorkItem(LoadClientInfo, recStr); } /// /// 接收消息 /// /// /// /// private void ComplexClient_AcceptString(AppSession arg1, NetHandle arg2, string arg3) { ThreadPool.QueueUserWorkItem(LoadClientInfo, (string)arg3); } public static void SendMsg(string str) { try { switch (ConnectType) { case 1: complexClient.Send(1, str); break; case 2: //var s = $"{MqttTargetTopic}\\?"; //var sendString = Regex.Split(str, s, RegexOptions.IgnoreCase); //TOBS??SCADAMain?DT|C_B_JointMoveTo| var topic = $"{str.Split('?')[2]}"; var content = $"{str.Split('?')[3]}"; Publish(topic, content); break; default: break; } } catch (Exception err) { } } public static void SendMsg(string opName, string topic, string content) { try { switch (ConnectType) { case 1: complexClient.Send(1, content); break; case 2: Publish(topic + opName, content); break; default: break; } } catch (Exception err) { } } private void ConnectService() { switch (ConnectType) { case 1: ClientConnection(IpPortSub); break; case 2: var _task = StartAsync(IpPortSub); break; default: break; } } } }