EOL progress and refactoring

This commit is contained in:
Stoyan Zlatev
2024-07-17 13:13:41 +02:00
parent 1ab24c37cd
commit 205b00d7da
269 changed files with 4105 additions and 4102 deletions
@@ -1,5 +1,6 @@
namespace Common.Hardware.Interfaces.Ports.SIRT
{
using Common.Hardware.Interfaces.Ports.SIRT.Requests;
using Common.Hardware.Interfaces.Ports.SIRT.Responses;
using Common.Hardware.Interfaces.Ports.SIRT.Tasks;
@@ -43,7 +44,7 @@
public void Run(SIRTTask task)
{
task.AddressChanged += address => this.tasks
.GetOrAdd(task.Address, new Queue<SIRTTask>())
.GetOrAdd(address, new Queue<SIRTTask>())
.Enqueue(task);
this.tasks
@@ -127,47 +128,6 @@
return this.Activated;
}
private void StartTrackingUnavailableDevicesAsync()
{
Task.Factory
.StartNew(state =>
{
if (state is SIRTBroadcaster broadcaster)
{
var awaiter = new ManualResetEventSlim();
while (!broadcaster.cancellationToken.IsCancellationRequested)
{
var devices = broadcaster.endpoints;
var now = DateTime.Now;
foreach (var address in devices.Keys)
{
if (devices.TryGetValue(address, out var timestamp) && (now - timestamp).TotalSeconds > 30)
{
devices.TryRemove(address, out var _);
broadcaster.DeviceUnavailable?.Invoke(address);
}
}
awaiter.Wait(1000, broadcaster.cancellationToken);
awaiter.Reset();
}
}
}, this, this.cancellationToken)
.ContinueWith(task =>
{
if (task.Exception != null && task.AsyncState is SIRTBroadcaster broadcaster)
{
var exception = task.Exception.ToString();
broadcaster.Notify(exception);
broadcaster.Close();
}
}, this.cancellationToken);
}
private void SirtBroadcaster_DataReceived(Byte[] data)
{
var address = default(UInt32);
@@ -205,56 +165,6 @@
, updateValueFactory: (_1, _2) => DateTime.Now);
}
private void StartPamPoolTrackingAsync()
{
Task.Factory
.StartNew(state =>
{
if (state is SIRTBroadcaster broadcaster)
{
var awaiter = new ManualResetEventSlim();
while (!broadcaster.cancellationToken.IsCancellationRequested)
{
var task = new ReadPamPoolSIRTTask();
task.Cancelled += _ => awaiter.Set();
task.Done += _ => awaiter.Set();
task.AddressesParsed += addresses =>
{
broadcaster.PamPoolChanged?.Invoke(addresses);
var length = addresses.Length;
for (Byte i = 0; i < length; i++)
{
var address = addresses[i];
broadcaster.pamPool.AddOrUpdate(i, address, (k, v) => address);
}
Thread.Sleep(1);
awaiter.Set();
};
broadcaster.Run(task);
awaiter.Wait(broadcaster.cancellationToken);
awaiter.Reset();
}
}
}, this, this.cancellationToken)
.ContinueWith(task =>
{
if (task.Exception != null && task.AsyncState is SIRTBroadcaster broadcaster)
{
var exception = task.Exception.ToString();
broadcaster.Notify(exception);
broadcaster.Close();
}
}, this.cancellationToken);
}
private void StartExecutingPendingTasksAsync()
{
Task.Factory
@@ -312,5 +222,104 @@
}, this.cancellationToken);
}
private void StartPamPoolTrackingAsync()
{
Task.Factory
.StartNew(state =>
{
if (state is SIRTBroadcaster broadcaster)
{
var awaiter = new ManualResetEventSlim();
while (!broadcaster.cancellationToken.IsCancellationRequested)
{
var task = new ReadPamPoolSIRTTask();
task.Cancelled += _ => awaiter.Set();
task.Done += _ => awaiter.Set();
task.AddressesParsed += addresses =>
{
broadcaster.PamPoolChanged?.Invoke(addresses);
var length = addresses.Length;
for (Byte i = 0; i < length; i++)
{
var address = addresses[i];
broadcaster.pamPool.AddOrUpdate(i, address, (k, v) => address);
}
Thread.Sleep(1);
awaiter.Set();
};
broadcaster.Run(task);
awaiter.Wait(broadcaster.cancellationToken);
awaiter.Reset();
}
}
}, this, this.cancellationToken)
.ContinueWith(task =>
{
if (task.Exception != null && task.AsyncState is SIRTBroadcaster broadcaster)
{
var exception = task.Exception.ToString();
broadcaster.Notify(exception);
broadcaster.Close();
}
}, this.cancellationToken);
}
private void StartTrackingUnavailableDevicesAsync()
{
Task.Factory
.StartNew(state =>
{
if (state is SIRTBroadcaster broadcaster)
{
var awaiter = new ManualResetEventSlim();
while (!broadcaster.cancellationToken.IsCancellationRequested)
{
var devices = broadcaster.endpoints;
var now = DateTime.Now;
foreach (var address in devices.Keys)
{
if (devices.TryGetValue(address, out var timestamp) && (now - timestamp).TotalSeconds > 10)
{
if (devices.TryRemove(address, out var _))
{
this.Write(new WritePamRequest
{
DelPamAdr = true,
Address = address,
});
broadcaster.DeviceUnavailable?.Invoke(address);
};
}
}
awaiter.Wait(1000, broadcaster.cancellationToken);
awaiter.Reset();
}
}
}, this, this.cancellationToken)
.ContinueWith(task =>
{
if (task.Exception != null && task.AsyncState is SIRTBroadcaster broadcaster)
{
var exception = task.Exception.ToString();
broadcaster.Notify(exception);
broadcaster.Close();
}
}, this.cancellationToken);
}
}
}
@@ -33,38 +33,24 @@
public Boolean IsRunning { get; internal set; }
internal void BeginAddResponse(FromAirResponse response)
{
if (this.IsDone || this.IsCancelled)
{
return;
}
if (this.Timestamp is null)
{
this.Timestamp = DateTime.Now;
}
var done = this.FromAirReceived(response);
this.EndAddResponse(done);
}
internal void BeginAddResponse(AcknOrAnswerResponse response)
{
if (this.IsDone || this.IsCancelled)
if (this.CanReceiveResponses())
{
return;
}
var done = this.AcknOrAnswerReceived(response);
if (this.Timestamp is null)
this.EndAddResponse(done);
}
}
internal void BeginAddResponse(FromAirResponse response)
{
if (this.CanReceiveResponses())
{
this.Timestamp = DateTime.Now;
var done = this.FromAirReceived(response);
this.EndAddResponse(done);
}
var done = this.AcknOrAnswerReceived(response);
this.EndAddResponse(done);
}
internal void NotifyRemoved()
@@ -107,7 +93,19 @@
{
this.Address = address;
this.AddressChanged?.Invoke(address);
this.AddressChanged?.Invoke(this.Address);
}
private Boolean CanReceiveResponses()
{
var canReceiveResponses = !this.IsDone && !this.IsCancelled;
if (canReceiveResponses && this.Timestamp is null)
{
this.Timestamp = DateTime.Now;
}
return canReceiveResponses;
}
private void EndAddResponse(Boolean isDone)
@@ -55,13 +55,5 @@
return parsed;
}
internal void Reset()
{
this.IsDone = false;
this.IsCancelled = false;
this.IsRunning = false;
this.Timestamp = null;
}
}
}
@@ -0,0 +1,17 @@
namespace Common.Hardware.Interfaces.Ports.SIRT.Tasks
{
using Common.Hardware.Interfaces.Ports.SIRT.Requests;
using System;
public class WritePamTask : SIRTTask
{
public WritePamTask(SIRTRequest request, UInt32 timeout) : base(request, timeout)
{
}
public WritePamTask(SIRTRequest request, UInt32 timeout) : base(request, timeout)
{
}
}
}