namespace Common.Hardware.SIRT { using System; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; public sealed class SIRTStream { private readonly SIRTPort sirtPort; private readonly bool log_RX_TX; private readonly List stream; public SIRTStream(string portName, bool log_RX_TX = true) { this.sirtPort = new SIRTPort(portName, false); this.stream = new List(); this.log_RX_TX = log_RX_TX; } public bool IsOpen => this.sirtPort.IsOpen; public string PortName => this.sirtPort.PortName; public void Close() => this.sirtPort.Close(); public void Open() => this.sirtPort.Open(); public IEnumerable ReadMessages() { void StartReading() { foreach (var bytes in this.sirtPort.ReadBytes()) { this.stream.AddRange(bytes); } } var cancellationTokenSource = new CancellationTokenSource(); var cancellationToken = cancellationTokenSource.Token; var task = Task.Factory.StartNew(StartReading, cancellationToken); var bytesCount = default(int); var startIndex = default(int); var lengthIndex = default(int); var length = default(int); var endIndex = default(int); var messageLength = default(int); var message = default(byte[]); while (this.sirtPort.IsOpen) { bytesCount = this.stream.Count; // until can parse messages from the buffer while (bytesCount > 0) { startIndex = this.stream.IndexOf(SIRTConstants.SIRT_PC, 0); // when no start byte found clear and break if (startIndex < 0) { this.stream.Clear(); break; } lengthIndex = startIndex + SIRTConstants.LEN_IX; // when not enough bytes to the length index just break if (lengthIndex >= bytesCount) { break; } length = this.stream[lengthIndex]; endIndex = lengthIndex + length + 3; // when not enough bytes to the end of the message just break if (endIndex >= bytesCount) { break; } messageLength = endIndex - startIndex + 1; message = new byte[messageLength]; this.stream.CopyTo(startIndex, message, 0, messageLength); this.stream.RemoveRange(0, endIndex + 1); // when message dose not ends with 0x16 clear bytes and break if (message[messageLength - 1] != SIRTConstants.MSG_END) { this.LogMessage($"XX|{BitConverter.ToString(message)}"); break; } if (this.log_RX_TX) { this.LogMessage($"RX|{BitConverter.ToString(message)}"); } // in case there are more messages in the buffer bytesCount = this.stream.Count; yield return message; } } try { cancellationTokenSource.Cancel(); } catch (Exception e) { this.LogMessage($"{e.Message}"); } } public void Write(params byte[] bytes) { this.sirtPort.Write(bytes); if (this.log_RX_TX) { this.LogMessage($"TX|{BitConverter.ToString(bytes)}"); } } private void LogMessage(object message) => SIRTLogger.LogMessage($"{nameof(SIRTStream)}|{this.PortName}|{message}"); } }