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;
}
}
}
}