You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 

528 line
17 KiB

  1. using System;
  2. using System.Collections.Concurrent;
  3. using System.Text;
  4. using System.Threading.Tasks;
  5. using Windows.Security.Cryptography.Certificates;
  6. using Windows.UI.Core;
  7. using Windows.UI.Xaml;
  8. using MQTTnet.Client;
  9. using MQTTnet.Diagnostics;
  10. using MQTTnet.Implementations;
  11. using MQTTnet.ManagedClient;
  12. using MQTTnet.Protocol;
  13. using MQTTnet.Server;
  14. namespace MQTTnet.TestApp.UniversalWindows
  15. {
  16. public sealed partial class MainPage
  17. {
  18. private readonly ConcurrentQueue<MqttNetLogMessage> _traceMessages = new ConcurrentQueue<MqttNetLogMessage>();
  19. private IMqttClient _mqttClient;
  20. private IMqttServer _mqttServer;
  21. public MainPage()
  22. {
  23. InitializeComponent();
  24. MqttNetGlobalLogger.LogMessagePublished += OnTraceMessagePublished;
  25. }
  26. private async void OnTraceMessagePublished(object sender, MqttNetLogMessagePublishedEventArgs e)
  27. {
  28. _traceMessages.Enqueue(e.TraceMessage);
  29. while (_traceMessages.Count > 100)
  30. {
  31. _traceMessages.TryDequeue(out _);
  32. }
  33. var logText = new StringBuilder();
  34. foreach (var traceMessage in _traceMessages)
  35. {
  36. logText.AppendFormat(
  37. "[{0:yyyy-MM-dd HH:mm:ss.fff}] [{1}] [{2}] [{3}] [{4}]{5}", traceMessage.Timestamp,
  38. traceMessage.Level,
  39. traceMessage.Source,
  40. traceMessage.ThreadId,
  41. traceMessage.Message,
  42. Environment.NewLine);
  43. if (traceMessage.Exception != null)
  44. {
  45. logText.AppendLine(traceMessage.Exception.ToString());
  46. }
  47. }
  48. await Trace.Dispatcher.RunAsync(CoreDispatcherPriority.Low, () =>
  49. {
  50. Trace.Text = logText.ToString();
  51. });
  52. }
  53. private async void Connect(object sender, RoutedEventArgs e)
  54. {
  55. var tlsOptions = new MqttClientTlsOptions
  56. {
  57. UseTls = UseTls.IsChecked == true,
  58. IgnoreCertificateChainErrors = true,
  59. IgnoreCertificateRevocationErrors = true,
  60. AllowUntrustedCertificates = true
  61. };
  62. var options = new MqttClientOptions { ClientId = ClientId.Text };
  63. if (UseTcp.IsChecked == true)
  64. {
  65. options.ChannelOptions = new MqttClientTcpOptions
  66. {
  67. Server = Server.Text,
  68. TlsOptions = tlsOptions
  69. };
  70. }
  71. if (UseWs.IsChecked == true)
  72. {
  73. options.ChannelOptions = new MqttClientWebSocketOptions
  74. {
  75. Uri = Server.Text,
  76. TlsOptions = tlsOptions
  77. };
  78. }
  79. if (options.ChannelOptions == null)
  80. {
  81. throw new InvalidOperationException();
  82. }
  83. options.Credentials = new MqttClientCredentials
  84. {
  85. Username = User.Text,
  86. Password = Password.Text
  87. };
  88. options.CleanSession = CleanSession.IsChecked == true;
  89. try
  90. {
  91. if (_mqttClient != null)
  92. {
  93. await _mqttClient.DisconnectAsync();
  94. _mqttClient.ApplicationMessageReceived -= OnApplicationMessageReceived;
  95. }
  96. var factory = new MqttFactory();
  97. _mqttClient = factory.CreateMqttClient();
  98. _mqttClient.ApplicationMessageReceived += OnApplicationMessageReceived;
  99. await _mqttClient.ConnectAsync(options);
  100. }
  101. catch (Exception exception)
  102. {
  103. Trace.Text += exception + Environment.NewLine;
  104. }
  105. }
  106. private async void OnApplicationMessageReceived(object sender, MqttApplicationMessageReceivedEventArgs eventArgs)
  107. {
  108. var item = $"Timestamp: {DateTime.Now:O} | Topic: {eventArgs.ApplicationMessage.Topic} | Payload: {Encoding.UTF8.GetString(eventArgs.ApplicationMessage.Payload)} | QoS: {eventArgs.ApplicationMessage.QualityOfServiceLevel}";
  109. await Dispatcher.RunAsync(CoreDispatcherPriority.Low, () =>
  110. {
  111. if (AddReceivedMessagesToList.IsChecked == true)
  112. {
  113. ReceivedMessages.Items.Add(item);
  114. }
  115. });
  116. }
  117. private async void Publish(object sender, RoutedEventArgs e)
  118. {
  119. if (_mqttClient == null)
  120. {
  121. return;
  122. }
  123. try
  124. {
  125. var qos = MqttQualityOfServiceLevel.AtMostOnce;
  126. if (QoS1.IsChecked == true)
  127. {
  128. qos = MqttQualityOfServiceLevel.AtLeastOnce;
  129. }
  130. if (QoS2.IsChecked == true)
  131. {
  132. qos = MqttQualityOfServiceLevel.ExactlyOnce;
  133. }
  134. var payload = new byte[0];
  135. if (Text.IsChecked == true)
  136. {
  137. payload = Encoding.UTF8.GetBytes(Payload.Text);
  138. }
  139. if (Base64.IsChecked == true)
  140. {
  141. payload = Convert.FromBase64String(Payload.Text);
  142. }
  143. var message = new MqttApplicationMessageBuilder()
  144. .WithTopic(Topic.Text)
  145. .WithPayload(payload)
  146. .WithQualityOfServiceLevel(qos)
  147. .WithRetainFlag(Retain.IsChecked == true)
  148. .Build();
  149. await _mqttClient.PublishAsync(message);
  150. }
  151. catch (Exception exception)
  152. {
  153. Trace.Text += exception + Environment.NewLine;
  154. }
  155. }
  156. private async void Disconnect(object sender, RoutedEventArgs e)
  157. {
  158. try
  159. {
  160. await _mqttClient.DisconnectAsync();
  161. }
  162. catch (Exception exception)
  163. {
  164. Trace.Text += exception + Environment.NewLine;
  165. }
  166. }
  167. private void Clear(object sender, RoutedEventArgs e)
  168. {
  169. Trace.Text = string.Empty;
  170. }
  171. private async void Subscribe(object sender, RoutedEventArgs e)
  172. {
  173. if (_mqttClient == null)
  174. {
  175. return;
  176. }
  177. try
  178. {
  179. var qos = MqttQualityOfServiceLevel.AtMostOnce;
  180. if (SubscribeQoS1.IsChecked == true)
  181. {
  182. qos = MqttQualityOfServiceLevel.AtLeastOnce;
  183. }
  184. if (SubscribeQoS2.IsChecked == true)
  185. {
  186. qos = MqttQualityOfServiceLevel.ExactlyOnce;
  187. }
  188. await _mqttClient.SubscribeAsync(new TopicFilter(SubscribeTopic.Text, qos));
  189. }
  190. catch (Exception exception)
  191. {
  192. Trace.Text += exception + Environment.NewLine;
  193. }
  194. }
  195. private async void Unsubscribe(object sender, RoutedEventArgs e)
  196. {
  197. if (_mqttClient == null)
  198. {
  199. return;
  200. }
  201. try
  202. {
  203. await _mqttClient.UnsubscribeAsync(SubscribeTopic.Text);
  204. }
  205. catch (Exception exception)
  206. {
  207. Trace.Text += exception + Environment.NewLine;
  208. }
  209. }
  210. // This code is for the Wiki at GitHub!
  211. // ReSharper disable once UnusedMember.Local
  212. private async Task WikiCode()
  213. {
  214. {
  215. // Write all trace messages to the console window.
  216. MqttNetGlobalLogger.LogMessagePublished += (s, e) =>
  217. {
  218. Console.WriteLine($">> [{e.TraceMessage.Timestamp:O}] [{e.TraceMessage.ThreadId}] [{e.TraceMessage.Source}] [{e.TraceMessage.Level}]: {e.TraceMessage.Message}");
  219. if (e.TraceMessage.Exception != null)
  220. {
  221. Console.WriteLine(e.TraceMessage.Exception);
  222. }
  223. };
  224. }
  225. {
  226. // Use a custom identifier for the trace messages.
  227. var clientOptions = new MqttClientOptionsBuilder()
  228. .Build();
  229. }
  230. {
  231. // Create a new MQTT client.
  232. var factory = new MqttFactory();
  233. var mqttClient = factory.CreateMqttClient();
  234. {
  235. // Create TCP based options using the builder.
  236. var options = new MqttClientOptionsBuilder()
  237. .WithClientId("Client1")
  238. .WithTcpServer("broker.hivemq.com")
  239. .WithCredentials("bud", "%spencer%")
  240. .WithTls()
  241. .WithCleanSession()
  242. .Build();
  243. await mqttClient.ConnectAsync(options);
  244. }
  245. {
  246. // Use TCP connection.
  247. var options = new MqttClientOptionsBuilder()
  248. .WithTcpServer("broker.hivemq.com", 1883) // Port is optional
  249. .Build();
  250. }
  251. {
  252. // Use secure TCP connection.
  253. var options = new MqttClientOptionsBuilder()
  254. .WithTcpServer("broker.hivemq.com")
  255. .WithTls()
  256. .Build();
  257. }
  258. {
  259. // Use WebSocket connection.
  260. var options = new MqttClientOptionsBuilder()
  261. .WithWebSocketServer("broker.hivemq.com:8000/mqtt")
  262. .Build();
  263. await mqttClient.ConnectAsync(options);
  264. }
  265. {
  266. // Create TCP based options manually
  267. var options = new MqttClientOptions
  268. {
  269. ClientId = "Client1",
  270. Credentials = new MqttClientCredentials
  271. {
  272. Username = "bud",
  273. Password = "%spencer%"
  274. },
  275. ChannelOptions = new MqttClientTcpOptions
  276. {
  277. Server = "broker.hivemq.org",
  278. TlsOptions = new MqttClientTlsOptions
  279. {
  280. UseTls = true
  281. }
  282. },
  283. };
  284. }
  285. {
  286. // Subscribe to a topic
  287. await mqttClient.SubscribeAsync(new TopicFilterBuilder().WithTopic("my/topic").Build());
  288. // Unsubscribe from a topic
  289. await mqttClient.UnsubscribeAsync("my/topic");
  290. // Publish an application message
  291. var applicationMessage = new MqttApplicationMessageBuilder()
  292. .WithTopic("A/B/C")
  293. .WithPayload("Hello World")
  294. .WithAtLeastOnceQoS()
  295. .Build();
  296. await mqttClient.PublishAsync(applicationMessage);
  297. }
  298. }
  299. // ----------------------------------
  300. {
  301. var options = new MqttServerOptions();
  302. options.ConnectionValidator = c =>
  303. {
  304. if (c.ClientId.Length < 10)
  305. {
  306. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedIdentifierRejected;
  307. return;
  308. }
  309. if (c.Username != "mySecretUser")
  310. {
  311. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedBadUsernameOrPassword;
  312. return;
  313. }
  314. if (c.Password != "mySecretPassword")
  315. {
  316. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedBadUsernameOrPassword;
  317. return;
  318. }
  319. c.ReturnCode = MqttConnectReturnCode.ConnectionAccepted;
  320. };
  321. var factory = new MqttFactory();
  322. var mqttServer = factory.CreateMqttServer();
  323. await mqttServer.StartAsync(options);
  324. Console.WriteLine("Press any key to exit.");
  325. Console.ReadLine();
  326. await mqttServer.StopAsync();
  327. }
  328. // ----------------------------------
  329. // For UWP apps:
  330. MqttTcpChannel.CustomIgnorableServerCertificateErrorsResolver = o =>
  331. {
  332. if (o.Server == "server_with_revoked_cert")
  333. {
  334. return new[] { ChainValidationResult.Revoked };
  335. }
  336. return new ChainValidationResult[0];
  337. };
  338. {
  339. // Start a MQTT server.
  340. var mqttServer = new MqttFactory().CreateMqttServer();
  341. await mqttServer.StartAsync(new MqttServerOptions());
  342. Console.WriteLine("Press any key to exit.");
  343. Console.ReadLine();
  344. await mqttServer.StopAsync();
  345. }
  346. {
  347. // Configure MQTT server.
  348. var optionsBuilder = new MqttServerOptionsBuilder()
  349. .WithConnectionBacklog(100)
  350. .WithDefaultEndpointPort(1884);
  351. var options = new MqttServerOptions
  352. {
  353. };
  354. options.ConnectionValidator = c =>
  355. {
  356. if (c.ClientId != "Highlander")
  357. {
  358. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedIdentifierRejected;
  359. return;
  360. }
  361. c.ReturnCode = MqttConnectReturnCode.ConnectionAccepted;
  362. };
  363. var mqttServer = new MqttFactory().CreateMqttServer();
  364. await mqttServer.StartAsync(optionsBuilder.Build());
  365. }
  366. {
  367. // Setup client validator.
  368. var options = new MqttServerOptions
  369. {
  370. ConnectionValidator = c =>
  371. {
  372. if (c.ClientId.Length < 10)
  373. {
  374. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedIdentifierRejected;
  375. return;
  376. }
  377. if (c.Username != "mySecretUser")
  378. {
  379. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedBadUsernameOrPassword;
  380. return;
  381. }
  382. if (c.Password != "mySecretPassword")
  383. {
  384. c.ReturnCode = MqttConnectReturnCode.ConnectionRefusedBadUsernameOrPassword;
  385. return;
  386. }
  387. c.ReturnCode = MqttConnectReturnCode.ConnectionAccepted;
  388. }
  389. };
  390. }
  391. {
  392. // Create a new MQTT server.
  393. var mqttServer = new MqttFactory().CreateMqttServer();
  394. }
  395. {
  396. // Setup and start a managed MQTT client.
  397. var options = new ManagedMqttClientOptionsBuilder()
  398. .WithAutoReconnectDelay(TimeSpan.FromSeconds(5))
  399. .WithClientOptions(new MqttClientOptionsBuilder()
  400. .WithClientId("Client1")
  401. .WithTcpServer("broker.hivemq.com")
  402. .WithTls().Build())
  403. .Build();
  404. var mqttClient = new MqttFactory().CreateManagedMqttClient();
  405. await mqttClient.SubscribeAsync(new TopicFilterBuilder().WithTopic("my/topic").Build());
  406. await mqttClient.StartAsync(options);
  407. }
  408. }
  409. private async void StartServer(object sender, RoutedEventArgs e)
  410. {
  411. if (_mqttServer != null)
  412. {
  413. return;
  414. }
  415. JsonServerStorage storage = null;
  416. if (ServerPersistRetainedMessages.IsChecked == true)
  417. {
  418. storage = new JsonServerStorage();
  419. if (ServerClearRetainedMessages.IsChecked == true)
  420. {
  421. storage.Clear();
  422. }
  423. }
  424. _mqttServer = new MqttFactory().CreateMqttServer();
  425. var options = new MqttServerOptions();
  426. options.DefaultEndpointOptions.Port = int.Parse(ServerPort.Text);
  427. options.Storage = storage;
  428. await _mqttServer.StartAsync(options);
  429. }
  430. private async void StopServer(object sender, RoutedEventArgs e)
  431. {
  432. if (_mqttServer == null)
  433. {
  434. return;
  435. }
  436. await _mqttServer.StopAsync();
  437. _mqttServer = null;
  438. }
  439. private void ClearReceivedMessages(object sender, RoutedEventArgs e)
  440. {
  441. ReceivedMessages.Items.Clear();
  442. }
  443. }
  444. }