Serial port stream added

This commit is contained in:
Stoyan Zlatev
2023-10-18 16:53:08 +02:00
parent 1b59ac55ce
commit 35ba0a62ea
6 changed files with 91 additions and 61 deletions
@@ -3,7 +3,6 @@ using OmniPlus.Connection.Enums;
using OmniPlus.Connection.Events;
using System;
using System.Collections.Concurrent;
using System.IO;
using System.IO.Ports;
@@ -18,8 +17,8 @@ namespace OmniPlus.Connection
private SerialPortState state;
protected bool reading;
protected SerialPortStream stream;
protected SerialPortOptions options;
protected ConcurrentQueue<byte> stream;
protected SerialPortConnection()
=> this.options = new SerialPortOptions();
@@ -71,8 +70,8 @@ namespace OmniPlus.Connection
try
{
this.stream = new ConcurrentQueue<byte>();
this.serialPort = new System.IO.Ports.SerialPort
this.stream = new SerialPortStream();
this.serialPort = new SerialPort
{
PortName = options.PortName,
BaudRate = (int)options.BaudRate,
@@ -149,19 +148,15 @@ namespace OmniPlus.Connection
{
if (this.open && sender is System.IO.Ports.SerialPort port && port.BytesToRead > 0)
{
var bytesLength = port.BytesToRead;
var bytes = new byte[bytesLength];
var bytes = new byte[port.BytesToRead];
var lock_object = new object();
lock (lock_object)
{
port.Read(bytes, 0, bytesLength);
port.Read(bytes, 0, port.BytesToRead);
}
foreach (var item in bytes)
{
this.stream.Enqueue(item);
}
this.stream.Append(bytes);
if (!this.reading)
{
@@ -0,0 +1,66 @@
using System;
namespace OmniPlus.Connection
{
public class SerialPortStream
{
private readonly int minBufferLength;
private byte[] buffer;
private int dataLength;
public SerialPortStream(int minBufferLength = byte.MaxValue)
{
this.minBufferLength = minBufferLength;
this.buffer = new byte[this.minBufferLength];
}
public byte this[int index] => this.buffer[index];
public void Append(byte[] bytes)
{
var bytesLength = bytes?.Length ?? 0;
if (bytesLength <= 0)
{
return;
}
var lock_object = new object();
lock (lock_object)
{
if (this.dataLength + bytesLength > this.buffer.Length)
{
var buffer = new byte[this.buffer.Length * 2];
Array.Copy(this.buffer, 0, buffer, 0, this.dataLength);
}
this.dataLength += bytesLength;
Array.Copy(bytes, 0, this.buffer, this.dataLength, bytesLength);
}
}
public void CopyTo(byte[] destinationArray, int destinationIndex, int sourceIndex, int length)
{
var requiredBufferLength = sourceIndex + length;
if (destinationArray is null || destinationIndex < 0 || sourceIndex < 0 || requiredBufferLength > this.buffer.Length)
{
return;
}
var lock_object = new object();
lock (lock_object)
{
this.dataLength = this.buffer.Length - requiredBufferLength;
Array.Copy(this.buffer, sourceIndex, destinationArray, destinationIndex, length);
Array.Copy(this.buffer, requiredBufferLength, this.buffer, 0, this.dataLength);
}
}
}
}
+15 -47
View File
@@ -44,8 +44,8 @@
#region Temporary
// TODO: remove for real connections
this.stream = new System.Collections.Concurrent.ConcurrentQueue<byte>(
new byte[] { 0xFF, 0x01, 0x25, 0x00, 0x80, 0x9B, 0x33, 0x04, 0x0D, 0x05, 0x08, 0x0C, 0xFF, 0x01, 0x25, 0x00, 0x80, 0x88 });
this.stream = new SerialPortStream();
this.stream.Append(new byte[] { 0xFF, 0x01, 0x25, 0x00, 0x80, 0x9B, 0x33, 0x04, 0x0D, 0x05, 0x08, 0x0C, 0xFF, 0x01, 0x25, 0x00, 0x80, 0x88 });
this.SerialPortConnection_Read();
#endregion
@@ -65,60 +65,28 @@
{
this.reading = true;
var message = default(IrdaMessage);
var message = default(byte[]);
var length = 0;
var start = 0;
while (this.reading)
while (true)
{
if (!this.stream.TryDequeue(out byte start))
if (this.stream[start] == MSG_START)
{
length = this.stream[start + 2] * 2 + 5;
message = new byte[length];
this.stream.CopyTo(message, 0, start, length);
break;
}
if (start != MSG_START)
{
continue;
}
_ = this.stream.TryDequeue(out byte type);
_ = this.stream.TryDequeue(out byte length);
var payloadLength = length * 2;
var payload = new byte[payloadLength];
for (int i = 0; i < payloadLength; i++)
{
if (this.stream.TryDequeue(out byte paybyte))
{
payload[i] = paybyte;
}
else break;
}
message = new IrdaMessage(payload);
message.Start = start;
message.Type = type;
message.DataLength = length;
var crc16 = new byte[2];
for (int i = 0; i < IrdaMessage.CRC_LENGTH; i++)
{
if (this.stream.TryDequeue(out byte crcbyte))
{
crc16[i] = crcbyte;
}
else break;
}
message.CopyCRC16(crc16, 0);
break;
start++;
}
this.reading = false;
this.messageReceived?.Invoke(message);
this.reading = false;
}
}
}
@@ -76,5 +76,8 @@ namespace OmniPlus.Connections.Messages
return payload;
}
public static implicit operator IrdaMessage(byte[] bytes)
=> new IrdaMessage(bytes, true);
}
}
+1
View File
@@ -54,6 +54,7 @@
<Compile Include="Connection\SerialPortMessage.cs" />
<Compile Include="Connection\SerialPortOptions.cs" />
<Compile Include="Connection\Enums\DataBits.cs" />
<Compile Include="Connection\SerialPortStream.cs" />
<Compile Include="Features\DeviceSpecific\BatteryVoltage.cs" />
<Compile Include="Features\DeviceSpecific\BatteryVoltageMeasurementCalibrationFactors.cs" />
<Compile Include="Features\DeviceSpecific\DataLogData.cs" />
-3
View File
@@ -1,8 +1,5 @@
using OmniPlus.Features.DeviceSpecific;
using OmniPlus.Features.Enums;
using OmniPlus.Features.Types;
using System;
namespace OmniPlus
{