Files
laatzen/Common/Hardware/Common.Hardware.SIRT/SIRTBroadcaster.Internal.cs
T

317 lines
9.7 KiB
C#

namespace Common.Hardware.SIRT
{
using Common.Hardware.Ports;
using System.Collections.Concurrent;
using System.Threading.Tasks;
using System.Threading;
using System.Collections.Generic;
using System.Linq;
using System;
public partial class SIRTBroadcaster
{
private readonly ConcurrentDictionary<uint, ConcurrentQueue<SIRTTask>> tasks;
private readonly ManualResetEventSlim processingWaitHandle;
private readonly SerialPort serialPort;
private CancellationTokenSource cancellationTokenSource;
private CancellationToken cancellationToken;
private void ActivateSIRT()
{
if (this.TryRun(SIRTTask.ActivateSIRT, x => x.State == SIRTTaskState.Completed))
{
var response = this.TryRun(SIRTTask.IdentifySIRT, x => x.Last);
if (response?.Length == 8)
{
this.HexId = BitConverter.ToString(response);
this.Frequency = this.HexId[10] == '0' ? 868 : 433;
this.Activated?.Invoke();
}
else
{
// this.Logging?.Invoke($"SIRT on {this.COMPort} not identified.");
}
}
else
{
// this.Logging?.Invoke($"SIRT on {this.COMPort} not activated.");
}
}
private ICollection<uint> GetAddressesInPamPool()
{
var addresses = new List<uint>();
var task = SIRTTask.ReadPamPool;
if (this.TryRun(task))
{
var payload = task.Last;
if (payload?.Length == 20)
{
var _address = 0U;
for (var i = 0; i < 20; i += 4)
{
_address = payload.ToUInt32BE(i);
if (_address != uint.MinValue && _address != uint.MaxValue)
{
addresses.Add(_address);
}
}
}
}
return addresses;
}
private void OnStateChanged(SIRTTask task)
{
if (task.State == SIRTTaskState.Completed)
{
if (this.tasks.TryGetValue(task.RequestAddress, out var requestTasks))
{
if (requestTasks.TryDequeue(out var _))
{
this.TaskRemoveEventHandlers(task);
}
}
else if (this.tasks.TryGetValue(task.ResponseAddress, out var responseTasks))
{
if (responseTasks.TryDequeue(out var _))
{
this.TaskRemoveEventHandlers(task);
}
}
}
else if (task.State == SIRTTaskState.Cancelled)
{
this.TaskRemoveEventHandlers(task);
this.OnTaskTimeout(task);
}
}
private void OnTaskTimeout(SIRTTask task)
{
var deleteRequestAddress = SIRTTask.DeleteFromPamPool(task.RequestAddress);
this.TryRun(deleteRequestAddress);
var deleteResponseAddress = SIRTTask.DeleteFromPamPool(task.ResponseAddress);
this.TryRun(deleteResponseAddress);
}
private void OnTaskUpdated(SIRTTask task)
{
if (this.tasks.TryGetValue(task.RequestAddress, out var tasksQueue))
{
foreach (var _task in tasksQueue.AsEnumerable())
{
_task.UpdateFrom(task);
}
}
}
private void SerialPort_Connected()
{
this.cancellationTokenSource = new CancellationTokenSource();
this.cancellationToken = this.cancellationTokenSource.Token;
Task.Factory.StartNew(this.StartReadingMessages, this.cancellationToken);
Task.Factory.StartNew(this.ActivateSIRT, this.cancellationToken);
this.Connected?.Invoke();
}
private void SerialPort_Disconnected()
{
try
{
this.cancellationTokenSource?.Cancel();
}
catch (Exception e)
{
this.SerialPort_Logging(e.ToString());
}
this.Disconnected?.Invoke();
}
private void SerialPort_Logging(string e)
{
//this.Logging?.Invoke(e);
}
private void SIRTConnection_MessageReceived(SIRTMessage message)
{
if (this.tasks.TryGetValue(message.Address, out var tasksQueue))
{
if (tasksQueue.TryPeek(out var task))
{
task.PushResponse(message);
this.Logging?.Invoke(message.ToString());
}
}
}
private void StartProcessingTasksAsync()
{
Task.Factory.StartNew(async state =>
{
var connection = state as SIRTBroadcaster;
var totalTasksCount = 0;
while (!connection.cancellationToken.IsCancellationRequested)
{
totalTasksCount = 0;
foreach (var tasksQueue in connection.tasks.Values.ToArray())
{
totalTasksCount += tasksQueue.Count;
if (tasksQueue.TryPeek(out var task) && task.State == SIRTTaskState.Pending)
{
connection.serialPort.Write(task.Start());
this.Logging?.Invoke(task.ToString());
}
await Task.Delay(50);
}
if (totalTasksCount <= 0)
{
this.processingWaitHandle.Wait();
this.processingWaitHandle.Reset();
}
else
{
await Task.Delay(100);
}
}
}, this, this.cancellationToken);
}
private void StartReadingMessages()
{
var buffer = new List<byte>();
int bufferLength;
int startIX;
int lengthIX;
int dataLength;
int stopIX;
int messageLength;
byte[] message;
foreach (var bytes in this.serialPort.ReadBytes())
{
if (bytes?.Length > 0)
{
buffer.AddRange(bytes);
}
bufferLength = buffer.Count;
while (bufferLength >= 13)
{
startIX = buffer.IndexOf(SIRTMessage.SIRT2PC);
if (startIX < 0)
{
buffer.Clear();
break;
}
lengthIX = startIX + 8;
if (lengthIX >= bufferLength)
{
break;
}
dataLength = buffer[lengthIX];
stopIX = lengthIX + dataLength + 3;
if (stopIX >= bufferLength)
{
break;
}
else if (buffer[stopIX] != SIRTMessage.STOP)
{
buffer.RemoveRange(0, stopIX + 1);
}
else
{
messageLength = stopIX - startIX + 1;
message = new byte[messageLength];
buffer.CopyTo(startIX, message, 0, messageLength);
buffer.RemoveRange(0, startIX + messageLength);
this.MessageReceived?.Invoke(new SIRTMessage(message));
}
bufferLength = buffer.Count;
}
}
}
private void TaskAddEventHandlers(SIRTTask task)
{
task.Cancelled += this.OnStateChanged;
task.Completed += this.OnStateChanged;
task.Timeouted += this.OnTaskTimeout;
task.Updated += this.OnTaskUpdated;
}
private void TaskRemoveEventHandlers(SIRTTask task)
{
task.Cancelled -= this.OnStateChanged;
task.Completed -= this.OnStateChanged;
task.Timeouted -= this.OnTaskTimeout;
task.Updated -= this.OnTaskUpdated;
}
private bool TryRun<TTask>(TTask task) where TTask : SIRTTask
=> this.TryRun(task, x => x.State == SIRTTaskState.Completed);
private TResult TryRun<TTask, TResult>(TTask task, Func<TTask, TResult> result) where TTask : SIRTTask
{
var manualResetEventSlim = new ManualResetEventSlim();
void taskStateChanged(SIRTTask _)
=> manualResetEventSlim.Set();
void messageReceived(SIRTMessage message)
{
if (message.Address == task.RequestAddress)
{
task.PushResponse(message);
}
}
this.MessageReceived += messageReceived;
task.Cancelled += taskStateChanged;
task.Completed += taskStateChanged;
this.serialPort.Write(task.Start());
// this.Logging?.Invoke(task.ToString());
manualResetEventSlim.Wait(task.Timeout);
manualResetEventSlim.Reset();
task.Cancelled += taskStateChanged;
task.Completed += taskStateChanged;
this.MessageReceived -= messageReceived;
return result(task);
}
}
}