using System; using System.Net.Sockets; using System.Text; using System.Threading; using System.Threading.Tasks; using Crestron.SimplSharp; using Crestron.SimplSharp.CrestronSockets; using ThreadingTimeout = System.Threading.Timeout; using NetSocketException = System.Net.Sockets.SocketException; namespace PepperDash.Core { /// /// A class to handle basic UDP communications to a remote endpoint /// public class GenericUdpClient : Device, ISocketStatusWithStreamDebugging, IAutoReconnect { private const string SplusKey = "Uninitialized UdpClient"; private readonly object stateLock = new object(); private readonly Timer reconnectTimer; private UdpClient client; private CancellationTokenSource receiveCancellationTokenSource; private bool connectEnabled; private bool connectionRefusedLogged; private SocketStatus clientStatus = SocketStatus.SOCKET_STATUS_NO_CONNECT; /// /// Object to enable stream debugging /// public CommunicationStreamDebugging StreamDebugging { get; private set; } /// /// Fires when data is received from the remote endpoint and returns it as a byte array /// public event EventHandler BytesReceived; /// /// Fires when data is received from the remote endpoint and returns it as text /// public event EventHandler TextReceived; /// /// Fires when the socket status changes /// public event EventHandler ConnectionChange; /// /// Address of remote endpoint /// public string Hostname { get; set; } /// /// Port on remote endpoint /// public int Port { get; set; } /// /// Another S+ helper because large port numbers can be treated as signed ints /// public ushort UPort { get { return Convert.ToUInt16(Port); } set { Port = Convert.ToInt32(value); } } /// /// Defaults to 2000 /// public int BufferSize { get; set; } /// /// True when the local socket is created and associated with the configured remote endpoint /// public bool IsConnected { get { return ClientStatus == SocketStatus.SOCKET_STATUS_CONNECTED; } } /// /// S+ helper for IsConnected /// public ushort UIsConnected { get { return (ushort)(IsConnected ? 1 : 0); } } /// /// The current socket status of the client /// public SocketStatus ClientStatus { get { lock (stateLock) { return clientStatus; } } private set { var shouldFireEvent = false; lock (stateLock) { if (clientStatus != value) { clientStatus = value; shouldFireEvent = true; } } if (shouldFireEvent) ConnectionChange?.Invoke(this, new GenericSocketStatusChageEventArgs(this)); } } /// /// Ushort representation of client status /// public ushort UStatus { get { return (ushort)ClientStatus; } } /// /// Gets or sets the AutoReconnect /// public bool AutoReconnect { get; set; } /// /// S+ helper for AutoReconnect /// public ushort UAutoReconnect { get { return (ushort)(AutoReconnect ? 1 : 0); } set { AutoReconnect = value == 1; } } /// /// Milliseconds to wait before attempting to reconnect. Defaults to 5000 /// public int AutoReconnectIntervalMs { get; set; } /// /// Constructor /// public GenericUdpClient(string key, string address, int port, int bufferSize) : base(key) { StreamDebugging = new CommunicationStreamDebugging(key); CrestronEnvironment.ProgramStatusEventHandler += CrestronEnvironment_ProgramStatusEventHandler; AutoReconnectIntervalMs = 5000; Hostname = address; Port = port; BufferSize = bufferSize; reconnectTimer = new Timer(o => { if (connectEnabled) Connect(); }, null, ThreadingTimeout.Infinite, ThreadingTimeout.Infinite); } /// /// Constructor for S+ /// public GenericUdpClient() : base(SplusKey) { StreamDebugging = new CommunicationStreamDebugging(SplusKey); CrestronEnvironment.ProgramStatusEventHandler += CrestronEnvironment_ProgramStatusEventHandler; AutoReconnectIntervalMs = 5000; BufferSize = 2000; reconnectTimer = new Timer(o => { if (connectEnabled) Connect(); }, null, ThreadingTimeout.Infinite, ThreadingTimeout.Infinite); } /// /// Initialize method /// public void Initialize(string key) { Key = key; } private void CrestronEnvironment_ProgramStatusEventHandler(eProgramStatusEventType programEventType) { if (programEventType == eProgramStatusEventType.Stopping) { Debug.Console(1, this, "Program stopping. Closing connection"); Deactivate(); } } /// /// Deactivate method /// public override bool Deactivate() { Disconnect(); return true; } /// /// Connect method /// public void Connect() { if (string.IsNullOrEmpty(Hostname)) { Debug.Console(1, Debug.ErrorLogLevel.Warning, "GenericUdpClient '{0}': No address set", Key); ClientStatus = SocketStatus.SOCKET_STATUS_NO_CONNECT; return; } if (Port < 1 || Port > 65535) { Debug.Console(1, Debug.ErrorLogLevel.Warning, "GenericUdpClient '{0}': Invalid port", Key); ClientStatus = SocketStatus.SOCKET_STATUS_NO_CONNECT; return; } var hostname = Hostname; var port = Port; var bufferSize = BufferSize; UdpClient newClient = null; CancellationTokenSource newReceiveCancellationTokenSource = null; CancellationToken startReceiveToken = default(CancellationToken); var shouldStartReceive = false; lock (stateLock) { connectEnabled = true; if (client != null) return; } try { newReceiveCancellationTokenSource = new CancellationTokenSource(); newClient = new UdpClient(); newClient.Client.ReceiveBufferSize = bufferSize; newClient.Client.SendBufferSize = bufferSize; newClient.Connect(hostname, port); lock (stateLock) { if (!connectEnabled || client != null) { newClient.Close(); newReceiveCancellationTokenSource.Cancel(); newReceiveCancellationTokenSource.Dispose(); return; } receiveCancellationTokenSource = newReceiveCancellationTokenSource; client = newClient; ClientStatus = SocketStatus.SOCKET_STATUS_CONNECTED; reconnectTimer.Change(ThreadingTimeout.Infinite, ThreadingTimeout.Infinite); startReceiveToken = receiveCancellationTokenSource.Token; shouldStartReceive = true; } if (shouldStartReceive) StartReceive(startReceiveToken); } catch (Exception ex) { Debug.LogMessage(ex, "Error connecting UDP client {0}", this, Key); if (newClient != null) newClient.Close(); if (newReceiveCancellationTokenSource != null) { newReceiveCancellationTokenSource.Cancel(); newReceiveCancellationTokenSource.Dispose(); } lock (stateLock) { if (connectEnabled && client == null) { ClientStatus = SocketStatus.SOCKET_STATUS_NO_CONNECT; StartReconnectTimer(); } } } } /// /// Disconnect method /// public void Disconnect() { lock (stateLock) { connectEnabled = false; reconnectTimer.Change(ThreadingTimeout.Infinite, ThreadingTimeout.Infinite); CleanupClient(); ClientStatus = SocketStatus.SOCKET_STATUS_NO_CONNECT; } } /// /// SendText method /// public void SendText(string text) { this.PrintSentText(text); var bytes = Encoding.GetEncoding(28591).GetBytes(text); SendBytes(bytes); } /// /// SendBytes method /// public void SendBytes(byte[] bytes) { if (bytes == null) return; try { this.PrintSentBytes(bytes); if (!IsConnected || client == null) Connect(); var udpClient = client; if (!IsConnected || udpClient == null) { Debug.Console(1, Debug.ErrorLogLevel.Warning, "GenericUdpClient '{0}': Cannot send bytes because the client is not connected", Key); return; } udpClient.Send(bytes, bytes.Length); } catch (Exception ex) { Debug.LogMessage(ex, "Error sending UDP bytes for {0}", this, Key); HandleDisconnected(); } } private void StartReceive(CancellationToken token) { Task.Run(async () => { while (!token.IsCancellationRequested) { try { var udpClient = client; if (udpClient == null) return; var result = await udpClient.ReceiveAsync().ConfigureAwait(false); var bytes = result.Buffer; if (bytes == null || bytes.Length == 0) continue; connectionRefusedLogged = false; var text = Encoding.GetEncoding(28591).GetString(bytes, 0, bytes.Length); this.PrintReceivedBytes(bytes); this.PrintReceivedText(text); BytesReceived?.Invoke(this, new GenericCommMethodReceiveBytesArgs(bytes)); TextReceived?.Invoke(this, new GenericCommMethodReceiveTextArgs(text)); } catch (ObjectDisposedException) { return; } catch (InvalidOperationException) { return; } catch (NetSocketException ex) { if (ex.SocketErrorCode == SocketError.ConnectionRefused) { if (!connectionRefusedLogged) { Debug.Console(1, Debug.ErrorLogLevel.Warning, "GenericUdpClient '{0}': Remote endpoint refused UDP traffic or is no longer listening", Key); connectionRefusedLogged = true; } HandleDisconnected(); return; } Debug.LogMessage(ex, "UDP receive error for {0}", this, Key); if (AutoReconnect) { HandleDisconnected(); return; } continue; } catch (Exception ex) { Debug.LogMessage(ex, "Unexpected UDP receive error for {0}", this, Key); if (AutoReconnect) { HandleDisconnected(); return; } continue; } } }, token); } private void HandleDisconnected() { lock (stateLock) { CleanupClient(); ClientStatus = SocketStatus.SOCKET_STATUS_NO_CONNECT; StartReconnectTimer(); } } private void StartReconnectTimer() { if (AutoReconnect && connectEnabled) reconnectTimer.Change(AutoReconnectIntervalMs, ThreadingTimeout.Infinite); } private void CleanupClient() { if (receiveCancellationTokenSource != null) { receiveCancellationTokenSource.Cancel(); receiveCancellationTokenSource.Dispose(); receiveCancellationTokenSource = null; } if (client != null) { client.Close(); client = null; } } } }