|
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240 |
- using System;
- using System.Collections.Concurrent;
- using System.Collections.Generic;
- using System.Linq;
- using System.Net;
- using System.Security.Authentication;
- using System.Security.Cryptography.X509Certificates;
- using System.Threading;
- using System.Threading.Tasks;
- using MQTTnet.Channel;
- using MQTTnet.Client.Options;
- using MQTTnet.Exceptions;
- using MQTTnet.Internal;
- using SuperSocket.ClientEngine;
- using WebSocket4Net;
-
- namespace MQTTnet.Extensions.WebSocket4Net
- {
- public sealed class WebSocket4NetMqttChannel : IMqttChannel
- {
- readonly BlockingCollection<byte> _receiveBuffer = new BlockingCollection<byte>();
-
- readonly IMqttClientOptions _clientOptions;
- readonly MqttClientWebSocketOptions _webSocketOptions;
-
- WebSocket _webSocket;
-
- public WebSocket4NetMqttChannel(IMqttClientOptions clientOptions, MqttClientWebSocketOptions webSocketOptions)
- {
- _clientOptions = clientOptions ?? throw new ArgumentNullException(nameof(clientOptions));
- _webSocketOptions = webSocketOptions ?? throw new ArgumentNullException(nameof(webSocketOptions));
- }
-
- public string Endpoint => _webSocketOptions.Uri;
-
- public bool IsSecureConnection { get; private set; }
-
- public X509Certificate2 ClientCertificate { get; }
-
- public async Task ConnectAsync(CancellationToken cancellationToken)
- {
- var uri = _webSocketOptions.Uri;
- if (!uri.StartsWith("ws://", StringComparison.OrdinalIgnoreCase) && !uri.StartsWith("wss://", StringComparison.OrdinalIgnoreCase))
- {
- if (_webSocketOptions.TlsOptions?.UseTls == false)
- {
- uri = "ws://" + uri;
- }
- else
- {
- uri = "wss://" + uri;
- }
- }
-
- var sslProtocols = _webSocketOptions?.TlsOptions?.SslProtocol ?? SslProtocols.None;
- var subProtocol = _webSocketOptions.SubProtocols.FirstOrDefault() ?? string.Empty;
-
- var cookies = new List<KeyValuePair<string, string>>();
- if (_webSocketOptions.CookieContainer != null)
- {
- throw new NotSupportedException("Cookies are not supported.");
- }
-
- List<KeyValuePair<string, string>> customHeaders = null;
- if (_webSocketOptions.RequestHeaders != null)
- {
- customHeaders = _webSocketOptions.RequestHeaders.Select(i => new KeyValuePair<string, string>(i.Key, i.Value)).ToList();
- }
-
- EndPoint proxy = null;
- if (_webSocketOptions.ProxyOptions != null)
- {
- throw new NotSupportedException("Proxies are not supported.");
- }
-
- // The user agent can be empty always because it is just added to the custom headers as "User-Agent".
- var userAgent = string.Empty;
-
- var origin = string.Empty;
- var webSocketVersion = WebSocketVersion.None;
- var receiveBufferSize = 0;
-
- var certificates = new X509CertificateCollection();
- if (_webSocketOptions?.TlsOptions?.Certificates != null)
- {
- foreach (var certificate in _webSocketOptions.TlsOptions.Certificates)
- {
- #if WINDOWS_UWP
- certificates.Add(new X509Certificate(certificate));
- #else
- certificates.Add(certificate);
- #endif
-
- }
- }
-
- _webSocket = new WebSocket(uri, subProtocol, cookies, customHeaders, userAgent, origin, webSocketVersion, proxy, sslProtocols, receiveBufferSize)
- {
- NoDelay = true,
- Security =
- {
- AllowUnstrustedCertificate = _webSocketOptions?.TlsOptions?.AllowUntrustedCertificates == true,
- AllowCertificateChainErrors = _webSocketOptions?.TlsOptions?.IgnoreCertificateChainErrors == true,
- Certificates = certificates
- }
- };
-
- await ConnectInternalAsync(cancellationToken).ConfigureAwait(false);
-
- IsSecureConnection = uri.StartsWith("wss://", StringComparison.OrdinalIgnoreCase);
- }
-
- public Task DisconnectAsync(CancellationToken cancellationToken)
- {
- if (_webSocket != null && _webSocket.State == WebSocketState.Open)
- {
- _webSocket.Close();
- }
-
- return Task.FromResult(0);
- }
-
- public Task<int> ReadAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
- {
- var readBytes = 0;
- while (count > 0 && !cancellationToken.IsCancellationRequested)
- {
- if (!_receiveBuffer.TryTake(out var @byte))
- {
- if (readBytes == 0)
- {
- // Block until at least one byte was received.
- @byte = _receiveBuffer.Take(cancellationToken);
- }
- else
- {
- return Task.FromResult(readBytes);
- }
- }
-
- buffer[offset] = @byte;
- offset++;
- count--;
- readBytes++;
- }
-
- return Task.FromResult(readBytes);
- }
-
- public Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
- {
- _webSocket.Send(buffer, offset, count);
- return Task.FromResult(0);
- }
-
- public void Dispose()
- {
- if (_webSocket == null)
- {
- return;
- }
-
- _webSocket.DataReceived -= OnDataReceived;
- _webSocket.Error -= OnError;
- _webSocket.Dispose();
- _webSocket = null;
- }
-
- static void OnError(object sender, ErrorEventArgs e)
- {
- System.Diagnostics.Debug.Write(e.Exception.ToString());
- }
-
- void OnDataReceived(object sender, DataReceivedEventArgs e)
- {
- foreach (var @byte in e.Data)
- {
- _receiveBuffer.Add(@byte);
- }
- }
-
- async Task ConnectInternalAsync(CancellationToken cancellationToken)
- {
- _webSocket.Error += OnError;
- _webSocket.DataReceived += OnDataReceived;
-
- var taskCompletionSource = new TaskCompletionSource<Exception>();
-
- void ErrorHandler(object sender, ErrorEventArgs e)
- {
- taskCompletionSource.TrySetResult(e.Exception);
- }
-
- void SuccessHandler(object sender, EventArgs e)
- {
- taskCompletionSource.TrySetResult(null);
- }
-
- try
- {
- _webSocket.Opened += SuccessHandler;
- _webSocket.Error += ErrorHandler;
-
- #pragma warning disable AsyncFixer02 // Long-running or blocking operations inside an async method
- _webSocket.Open();
- #pragma warning restore AsyncFixer02 // Long-running or blocking operations inside an async method
-
- var exception = await MqttTaskTimeout.WaitAsync(c =>
- {
- c.Register(() => taskCompletionSource.TrySetCanceled());
- return taskCompletionSource.Task;
- }, _clientOptions.CommunicationTimeout, cancellationToken).ConfigureAwait(false);
-
- if (exception != null)
- {
- if (exception is AuthenticationException authenticationException)
- {
- throw new MqttCommunicationException(authenticationException.InnerException);
- }
-
- if (exception is OperationCanceledException)
- {
- throw new MqttCommunicationTimedOutException();
- }
-
- throw new MqttCommunicationException(exception);
- }
- }
- catch (OperationCanceledException)
- {
- throw new MqttCommunicationTimedOutException();
- }
- finally
- {
- _webSocket.Opened -= SuccessHandler;
- _webSocket.Error -= ErrorHandler;
- }
- }
- }
- }
|