WingMan – Blame information for rev 36
?pathlinks?
Rev | Author | Line No. | Line |
---|---|---|---|
10 | office | 1 | using System; |
12 | office | 2 | using System.IO; |
34 | office | 3 | using System.IO.Compression; |
10 | office | 4 | using System.Net; |
5 | using System.Threading; |
||
6 | using System.Threading.Tasks; |
||
34 | office | 7 | using LZ4; |
10 | office | 8 | using MQTTnet; |
9 | using MQTTnet.Client; |
||
10 | using MQTTnet.Extensions.ManagedClient; |
||
11 | using MQTTnet.Protocol; |
||
12 | using MQTTnet.Server; |
||
13 | using WingMan.Utilities; |
||
14 | using MqttClientConnectedEventArgs = MQTTnet.Client.MqttClientConnectedEventArgs; |
||
15 | using MqttClientDisconnectedEventArgs = MQTTnet.Client.MqttClientDisconnectedEventArgs; |
||
16 | |||
17 | namespace WingMan.Communication |
||
18 | { |
||
19 | public class MqttCommunication : IDisposable |
||
20 | { |
||
21 | public delegate void ClientAuthenticationFailed(object sender, MqttAuthenticationFailureEventArgs e); |
||
22 | |||
23 | public delegate void ClientConnected(object sender, MqttClientConnectedEventArgs e); |
||
24 | |||
25 | public delegate void ClientConnectionFailed(object sender, MqttManagedProcessFailedEventArgs e); |
||
26 | |||
27 | public delegate void ClientDisconnected(object sender, MqttClientDisconnectedEventArgs e); |
||
28 | |||
29 | public delegate void ClientSubscribed(object sender, MqttClientSubscribedTopicEventArgs e); |
||
30 | |||
31 | public delegate void ClientUnsubscribed(object sender, MqttClientUnsubscribedTopicEventArgs e); |
||
32 | |||
35 | office | 33 | public delegate void MessageReceived(object sender, MqttCommunicationMessageReceivedEventArgs e); |
10 | office | 34 | |
35 | public delegate void ServerAuthenticationFailed(object sender, MqttAuthenticationFailureEventArgs e); |
||
36 | |||
37 | public delegate void ServerClientConnected(object sender, MQTTnet.Server.MqttClientConnectedEventArgs e); |
||
38 | |||
39 | public delegate void ServerClientDisconnected(object sender, MQTTnet.Server.MqttClientDisconnectedEventArgs e); |
||
40 | |||
41 | public delegate void ServerStarted(object sender, EventArgs e); |
||
42 | |||
43 | public delegate void ServerStopped(object sender, EventArgs e); |
||
44 | |||
45 | public MqttCommunication(TaskScheduler taskScheduler, CancellationToken cancellationToken) |
||
46 | { |
||
47 | TaskScheduler = taskScheduler; |
||
48 | CancellationToken = cancellationToken; |
||
49 | |||
50 | Client = new MqttFactory().CreateManagedMqttClient(); |
||
51 | Server = new MqttFactory().CreateMqttServer(); |
||
52 | } |
||
53 | |||
54 | private TaskScheduler TaskScheduler { get; } |
||
55 | |||
56 | private IManagedMqttClient Client { get; } |
||
57 | |||
58 | private IMqttServer Server { get; } |
||
59 | |||
60 | public bool Running { get; set; } |
||
61 | |||
62 | public string Nick { get; set; } |
||
63 | |||
64 | private IPAddress IpAddress { get; set; } |
||
65 | |||
36 | office | 66 | public int Port { get; set; } |
10 | office | 67 | |
68 | private string Password { get; set; } |
||
69 | |||
70 | private CancellationToken CancellationToken { get; } |
||
71 | |||
72 | public MqttCommunicationType Type { get; set; } |
||
73 | |||
74 | public async void Dispose() |
||
75 | { |
||
76 | await Stop(); |
||
77 | } |
||
78 | |||
79 | public event ClientAuthenticationFailed OnClientAuthenticationFailed; |
||
80 | |||
81 | public event ServerAuthenticationFailed OnServerAuthenticationFailed; |
||
82 | |||
83 | public event MessageReceived OnMessageReceived; |
||
84 | |||
85 | public event ClientConnected OnClientConnected; |
||
86 | |||
87 | public event ClientDisconnected OnClientDisconnected; |
||
88 | |||
89 | public event ClientConnectionFailed OnClientConnectionFailed; |
||
90 | |||
91 | public event ClientUnsubscribed OnClientUnsubscribed; |
||
92 | |||
93 | public event ClientSubscribed OnClientSubscribed; |
||
94 | |||
95 | public event ServerClientDisconnected OnServerClientDisconnected; |
||
96 | |||
97 | public event ServerClientConnected OnServerClientConnected; |
||
98 | |||
99 | public event ServerStarted OnServerStarted; |
||
100 | |||
101 | public event ServerStopped OnServerStopped; |
||
102 | |||
103 | public async Task<bool> Start(MqttCommunicationType type, IPAddress ipAddress, int port, string nick, |
||
104 | string password) |
||
105 | { |
||
106 | Type = type; |
||
107 | IpAddress = ipAddress; |
||
108 | Port = port; |
||
109 | Nick = nick; |
||
110 | Password = password; |
||
111 | |||
112 | switch (type) |
||
113 | { |
||
114 | case MqttCommunicationType.Client: |
||
115 | return await StartClient(); |
||
116 | case MqttCommunicationType.Server: |
||
117 | return await StartServer(); |
||
118 | } |
||
119 | |||
120 | return false; |
||
121 | } |
||
122 | |||
123 | private async Task<bool> StartClient() |
||
124 | { |
||
125 | var clientOptions = new MqttClientOptionsBuilder() |
||
126 | .WithTcpServer(IpAddress.ToString(), Port); |
||
127 | |||
128 | // Setup and start a managed MQTT client. |
||
129 | var options = new ManagedMqttClientOptionsBuilder() |
||
130 | .WithClientOptions(clientOptions.Build()) |
||
131 | .Build(); |
||
132 | |||
133 | BindClientHandlers(); |
||
134 | |||
135 | await Client.SubscribeAsync( |
||
136 | new TopicFilterBuilder() |
||
137 | .WithTopic("lobby") |
||
138 | .Build() |
||
139 | ); |
||
140 | |||
141 | await Client.SubscribeAsync( |
||
142 | new TopicFilterBuilder() |
||
143 | .WithTopic("exchange") |
||
144 | .Build() |
||
145 | ); |
||
146 | |||
147 | await Client.SubscribeAsync( |
||
148 | new TopicFilterBuilder() |
||
149 | .WithTopic("execute") |
||
150 | .Build() |
||
151 | ); |
||
152 | |||
153 | await Client.StartAsync(options); |
||
154 | |||
155 | Running = true; |
||
156 | |||
157 | return Running; |
||
158 | } |
||
159 | |||
160 | private async Task StopClient() |
||
161 | { |
||
36 | office | 162 | await Client.StopAsync(); |
163 | |||
10 | office | 164 | UnbindClientHandlers(); |
165 | } |
||
166 | |||
167 | public void BindClientHandlers() |
||
168 | { |
||
169 | Client.Connected += ClientOnConnected; |
||
170 | Client.Disconnected += ClientOnDisconnected; |
||
171 | Client.ConnectingFailed += ClientOnConnectingFailed; |
||
172 | Client.ApplicationMessageReceived += ClientOnApplicationMessageReceived; |
||
173 | } |
||
174 | |||
175 | public void UnbindClientHandlers() |
||
176 | { |
||
177 | Client.Connected -= ClientOnConnected; |
||
178 | Client.Disconnected -= ClientOnDisconnected; |
||
179 | Client.ConnectingFailed -= ClientOnConnectingFailed; |
||
180 | Client.ApplicationMessageReceived -= ClientOnApplicationMessageReceived; |
||
181 | } |
||
182 | |||
183 | private async void ClientOnApplicationMessageReceived(object sender, MqttApplicationMessageReceivedEventArgs e) |
||
184 | { |
||
185 | try |
||
186 | { |
||
34 | office | 187 | using (var inputStream = new MemoryStream(e.ApplicationMessage.Payload)) |
188 | { |
||
36 | office | 189 | using (var decryptedStream = await Aes.Decrypt(inputStream, Password)) |
34 | office | 190 | { |
191 | using (var lz4Decompress = new LZ4Stream(decryptedStream, CompressionMode.Decompress)) |
||
192 | { |
||
35 | office | 193 | var outpuStream = new MemoryStream(); |
194 | await lz4Decompress.CopyToAsync(outpuStream); |
||
10 | office | 195 | |
35 | office | 196 | outpuStream.Position = 0L; |
34 | office | 197 | |
35 | office | 198 | await Task.Delay(0, CancellationToken).ContinueWith( |
199 | _ => OnMessageReceived?.Invoke(sender, |
||
200 | new MqttCommunicationMessageReceivedEventArgs(e.ApplicationMessage.Topic, |
||
201 | outpuStream)), |
||
202 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
34 | office | 203 | } |
204 | } |
||
205 | } |
||
10 | office | 206 | } |
207 | catch (Exception ex) |
||
208 | { |
||
209 | await Task.Delay(0, CancellationToken).ContinueWith( |
||
210 | _ => OnClientAuthenticationFailed?.Invoke(sender, new MqttAuthenticationFailureEventArgs(e, ex)), |
||
211 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
212 | } |
||
213 | } |
||
214 | |||
215 | private async void ClientOnConnectingFailed(object sender, MqttManagedProcessFailedEventArgs e) |
||
216 | { |
||
217 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnClientConnectionFailed?.Invoke(sender, e), |
||
218 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
219 | } |
||
220 | |||
221 | private async void ClientOnDisconnected(object sender, MqttClientDisconnectedEventArgs e) |
||
222 | { |
||
223 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnClientDisconnected?.Invoke(sender, e), |
||
224 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
225 | } |
||
226 | |||
227 | private async void ClientOnConnected(object sender, MqttClientConnectedEventArgs e) |
||
228 | { |
||
229 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnClientConnected?.Invoke(sender, e), |
||
230 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
231 | } |
||
232 | |||
233 | private async Task<bool> StartServer() |
||
234 | { |
||
235 | var optionsBuilder = new MqttServerOptionsBuilder() |
||
236 | .WithDefaultEndpointBoundIPAddress(IpAddress) |
||
237 | .WithDefaultEndpointPort(Port); |
||
238 | |||
239 | BindServerHandlers(); |
||
240 | |||
241 | try |
||
242 | { |
||
243 | await Server.StartAsync(optionsBuilder.Build()); |
||
244 | |||
245 | Running = true; |
||
246 | } |
||
247 | catch (Exception) |
||
248 | { |
||
249 | Running = false; |
||
250 | } |
||
251 | |||
252 | return Running; |
||
253 | } |
||
254 | |||
255 | private async Task StopServer() |
||
256 | { |
||
36 | office | 257 | await Server.StopAsync(); |
258 | |||
10 | office | 259 | UnbindServerHandlers(); |
260 | } |
||
261 | |||
262 | private void BindServerHandlers() |
||
263 | { |
||
264 | Server.Started += ServerOnStarted; |
||
265 | Server.Stopped += ServerOnStopped; |
||
266 | Server.ClientConnected += ServerOnClientConnected; |
||
267 | Server.ClientDisconnected += ServerOnClientDisconnected; |
||
268 | Server.ClientSubscribedTopic += ServerOnClientSubscribedTopic; |
||
269 | Server.ClientUnsubscribedTopic += ServerOnClientUnsubscribedTopic; |
||
270 | Server.ApplicationMessageReceived += ServerOnApplicationMessageReceived; |
||
271 | } |
||
272 | |||
273 | private void UnbindServerHandlers() |
||
274 | { |
||
275 | Server.Started -= ServerOnStarted; |
||
276 | Server.Stopped -= ServerOnStopped; |
||
277 | Server.ClientConnected -= ServerOnClientConnected; |
||
278 | Server.ClientDisconnected -= ServerOnClientDisconnected; |
||
279 | Server.ClientSubscribedTopic -= ServerOnClientSubscribedTopic; |
||
280 | Server.ClientUnsubscribedTopic -= ServerOnClientUnsubscribedTopic; |
||
281 | Server.ApplicationMessageReceived -= ServerOnApplicationMessageReceived; |
||
282 | } |
||
283 | |||
284 | private async void ServerOnApplicationMessageReceived(object sender, MqttApplicationMessageReceivedEventArgs e) |
||
285 | { |
||
286 | try |
||
287 | { |
||
34 | office | 288 | using (var inputStream = new MemoryStream(e.ApplicationMessage.Payload)) |
289 | { |
||
36 | office | 290 | using (var decryptedStream = await Aes.Decrypt(inputStream, Password)) |
34 | office | 291 | { |
292 | using (var lz4Decompress = new LZ4Stream(decryptedStream, CompressionMode.Decompress)) |
||
293 | { |
||
35 | office | 294 | var outpuStream = new MemoryStream(); |
295 | await lz4Decompress.CopyToAsync(outpuStream); |
||
10 | office | 296 | |
35 | office | 297 | outpuStream.Position = 0L; |
34 | office | 298 | |
35 | office | 299 | await Task.Delay(0, CancellationToken).ContinueWith( |
300 | _ => OnMessageReceived?.Invoke(sender, |
||
301 | new MqttCommunicationMessageReceivedEventArgs(e.ApplicationMessage.Topic, |
||
302 | outpuStream)), |
||
303 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
34 | office | 304 | } |
305 | } |
||
306 | } |
||
10 | office | 307 | } |
308 | catch (Exception ex) |
||
309 | { |
||
310 | foreach (var clientSessionStatus in await Server.GetClientSessionsStatusAsync()) |
||
311 | { |
||
312 | if (!string.Equals(clientSessionStatus.ClientId, e.ClientId, StringComparison.Ordinal)) |
||
313 | continue; |
||
314 | |||
315 | await clientSessionStatus.DisconnectAsync(); |
||
316 | } |
||
317 | |||
318 | await Task.Delay(0, CancellationToken).ContinueWith( |
||
319 | _ => OnServerAuthenticationFailed?.Invoke(sender, new MqttAuthenticationFailureEventArgs(e, ex)), |
||
320 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
321 | } |
||
322 | } |
||
323 | |||
324 | private async void ServerOnClientUnsubscribedTopic(object sender, MqttClientUnsubscribedTopicEventArgs e) |
||
325 | { |
||
326 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnClientUnsubscribed?.Invoke(sender, e), |
||
327 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
328 | } |
||
329 | |||
330 | private async void ServerOnClientSubscribedTopic(object sender, MqttClientSubscribedTopicEventArgs e) |
||
331 | { |
||
332 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnClientSubscribed?.Invoke(sender, e), |
||
333 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
334 | } |
||
335 | |||
336 | private async void ServerOnClientDisconnected(object sender, MQTTnet.Server.MqttClientDisconnectedEventArgs e) |
||
337 | { |
||
338 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnServerClientDisconnected?.Invoke(sender, e), |
||
339 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
340 | } |
||
341 | |||
342 | private async void ServerOnClientConnected(object sender, MQTTnet.Server.MqttClientConnectedEventArgs e) |
||
343 | { |
||
344 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnServerClientConnected?.Invoke(sender, e), |
||
345 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
346 | } |
||
347 | |||
12 | office | 348 | private async void ServerOnStopped(object sender, EventArgs e) |
10 | office | 349 | { |
12 | office | 350 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnServerStopped?.Invoke(sender, e), |
351 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
10 | office | 352 | } |
353 | |||
12 | office | 354 | private async void ServerOnStarted(object sender, EventArgs e) |
10 | office | 355 | { |
12 | office | 356 | await Task.Delay(0, CancellationToken).ContinueWith(_ => OnServerStarted?.Invoke(sender, e), |
357 | CancellationToken, TaskContinuationOptions.None, TaskScheduler); |
||
10 | office | 358 | } |
359 | |||
360 | public async Task Stop() |
||
361 | { |
||
362 | switch (Type) |
||
363 | { |
||
364 | case MqttCommunicationType.Server: |
||
365 | await StopServer(); |
||
366 | break; |
||
367 | case MqttCommunicationType.Client: |
||
368 | await StopClient(); |
||
369 | break; |
||
370 | } |
||
371 | |||
372 | Running = false; |
||
373 | } |
||
374 | |||
375 | public async Task Broadcast(string topic, byte[] payload) |
||
376 | { |
||
34 | office | 377 | using (var compressStream = new MemoryStream()) |
10 | office | 378 | { |
34 | office | 379 | using (var lz4Stream = new LZ4Stream(compressStream, CompressionMode.Compress)) |
12 | office | 380 | { |
34 | office | 381 | await lz4Stream.WriteAsync(payload, 0, payload.Length); |
382 | await lz4Stream.FlushAsync(); |
||
383 | |||
384 | compressStream.Position = 0L; |
||
385 | |||
36 | office | 386 | using (var outputStream = await Aes.Encrypt(compressStream, Password)) |
34 | office | 387 | { |
388 | var data = outputStream.ToArray(); |
||
389 | switch (Type) |
||
12 | office | 390 | { |
34 | office | 391 | case MqttCommunicationType.Client: |
392 | await Client.PublishAsync(new[] |
||
393 | { |
||
394 | new MqttApplicationMessage |
||
395 | { |
||
396 | Topic = topic, |
||
397 | Payload = data, |
||
398 | QualityOfServiceLevel = MqttQualityOfServiceLevel.AtMostOnce |
||
399 | } |
||
400 | }); |
||
401 | break; |
||
402 | case MqttCommunicationType.Server: |
||
403 | await Server.PublishAsync(new[] |
||
404 | { |
||
405 | new MqttApplicationMessage |
||
406 | { |
||
407 | Topic = topic, |
||
408 | Payload = data, |
||
409 | QualityOfServiceLevel = MqttQualityOfServiceLevel.AtMostOnce |
||
410 | } |
||
411 | }); |
||
412 | break; |
||
413 | } |
||
414 | } |
||
12 | office | 415 | } |
10 | office | 416 | } |
417 | } |
||
418 | } |
||
419 | } |