Upgrade (develop) - GCI - big change
This commit is contained in:
@@ -0,0 +1,314 @@
|
||||
using System;
|
||||
using System.Collections.Concurrent;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace GenesisCordonelInterface.Core.Threading
|
||||
{
|
||||
/*
|
||||
ApiWorker – per-slot sequential execution worker
|
||||
|
||||
This class provides a lightweight background worker that executes actions
|
||||
sequentially on a dedicated thread.
|
||||
|
||||
PRIMARY PURPOSE
|
||||
---------------
|
||||
ApiWorker is designed to safely execute hardware-related operations
|
||||
(e.g. meter communication) without blocking the UI thread and without
|
||||
allowing concurrent access to the same device.
|
||||
|
||||
Each ApiWorker instance typically represents:
|
||||
1 worker = 1 slot = 1 meter = 1 communication channel
|
||||
|
||||
KEY PROPERTIES
|
||||
--------------
|
||||
- Single dedicated background thread
|
||||
- FIFO queue (first-in, first-out)
|
||||
- Sequential execution (NO parallelism inside one worker)
|
||||
- Thread-safe enqueueing
|
||||
- Task-based async interface for callers
|
||||
|
||||
WHY THIS IS IMPORTANT
|
||||
--------------------
|
||||
Hardware communication (serial ports, meters, etc.) is usually NOT thread-safe.
|
||||
If multiple commands are executed in parallel, communication may break or corrupt data.
|
||||
|
||||
ApiWorker guarantees:
|
||||
- operations are executed one-by-one
|
||||
- order is preserved
|
||||
- no race conditions on the device
|
||||
|
||||
HIGH-LEVEL FLOW
|
||||
---------------
|
||||
Caller (UI/API)
|
||||
|
|
||||
v
|
||||
RunAsync(...)
|
||||
|
|
||||
v
|
||||
TaskCompletionSource created
|
||||
|
|
||||
v
|
||||
Action wrapped into queue item
|
||||
|
|
||||
v
|
||||
Added to BlockingCollection queue
|
||||
|
|
||||
v
|
||||
Worker thread consumes queue
|
||||
|
|
||||
v
|
||||
Action executed (blocking HW call)
|
||||
|
|
||||
v
|
||||
Result propagated via TaskCompletionSource
|
||||
|
|
||||
v
|
||||
Caller receives result via await
|
||||
|
||||
GRAPH
|
||||
-----
|
||||
Caller thread (UI)
|
||||
|
|
||||
v
|
||||
RunAsync()
|
||||
|
|
||||
v
|
||||
Queue (BlockingCollection)
|
||||
|
|
||||
v
|
||||
-----------------------------
|
||||
| Worker Thread (background)|
|
||||
| while(queue) |
|
||||
| Execute Action |
|
||||
-----------------------------
|
||||
|
|
||||
v
|
||||
Task result (await)
|
||||
|
||||
THREADING MODEL
|
||||
---------------
|
||||
- Producer/Consumer pattern
|
||||
- Producer: any thread calling RunAsync
|
||||
- Consumer: single worker thread
|
||||
- Synchronization handled by BlockingCollection
|
||||
|
||||
MAIN COMPONENTS
|
||||
---------------
|
||||
1. BlockingCollection<Action> queue
|
||||
- thread-safe queue
|
||||
- stores work items
|
||||
- supports blocking consumption
|
||||
|
||||
2. Dedicated Thread
|
||||
- runs WorkerLoop()
|
||||
- continuously processes queue
|
||||
|
||||
3. TaskCompletionSource<T>
|
||||
- bridges sync execution → async API
|
||||
- allows caller to await result
|
||||
|
||||
METHODS
|
||||
-------
|
||||
|
||||
RunAsync<T>(Func<T>)
|
||||
--------------------
|
||||
- Enqueues a function returning a value
|
||||
- Wraps it into Action
|
||||
- Executes on worker thread
|
||||
- Returns Task<T> to caller
|
||||
|
||||
RunAsync(Action)
|
||||
----------------
|
||||
- Convenience overload for void methods
|
||||
- Internally wraps into Func<object>
|
||||
|
||||
WorkerLoop()
|
||||
------------
|
||||
- Infinite loop consuming queue
|
||||
- Executes actions one-by-one
|
||||
- Stops when queue is completed
|
||||
|
||||
Dispose()
|
||||
---------
|
||||
- Stops accepting new items
|
||||
- Cleans up queue
|
||||
- Does NOT forcibly stop running task
|
||||
|
||||
CANCELLATION
|
||||
------------
|
||||
- CancellationToken is checked BEFORE execution
|
||||
- If cancelled → Task is cancelled
|
||||
- Does NOT interrupt running operation
|
||||
|
||||
IMPORTANT LIMITATIONS
|
||||
--------------------
|
||||
- No parallel execution inside one worker (by design)
|
||||
- Long-running action blocks worker thread
|
||||
- No built-in timeout handling
|
||||
- Dispose does not abort running work
|
||||
|
||||
WHEN TO USE
|
||||
-----------
|
||||
✔ Per-device communication (serial, TCP, HW)
|
||||
✔ Ordered execution required
|
||||
✔ UI must stay responsive
|
||||
|
||||
WHEN NOT TO USE
|
||||
---------------
|
||||
✘ CPU parallel processing (use Task.Run / Parallel)
|
||||
✘ High-throughput parallel workloads
|
||||
✘ Fire-and-forget background tasks
|
||||
|
||||
SUMMARY
|
||||
-------
|
||||
ApiWorker is a simple, robust solution for:
|
||||
"Execute commands sequentially per resource, asynchronously from UI"
|
||||
|
||||
It is a perfect fit for:
|
||||
- hardware interfaces
|
||||
- device drivers
|
||||
- IO-bound serialized workflows
|
||||
*/
|
||||
|
||||
public sealed class ApiWorker : IDisposable
|
||||
{
|
||||
/// <summary>
|
||||
/// Thread-safe FIFO queue holding work items.
|
||||
/// </summary>
|
||||
private readonly BlockingCollection<Action> queue = new BlockingCollection<Action>();
|
||||
|
||||
/// <summary>
|
||||
/// Dedicated worker thread processing the queue.
|
||||
/// </summary>
|
||||
private readonly Thread thread;
|
||||
|
||||
/// <summary>
|
||||
/// Indicates whether this worker has been disposed.
|
||||
/// </summary>
|
||||
private bool disposed;
|
||||
|
||||
public string Name { get; private set; }
|
||||
public int QueueLength { get { return queue.Count; } }
|
||||
public bool IsBusy { get; private set; }
|
||||
public string CurrentOperation { get; private set; }
|
||||
public string LastError { get; private set; }
|
||||
public DateTime LastActivity { get; private set; }
|
||||
|
||||
/// <summary>
|
||||
/// Creates a new ApiWorker with its own background thread.
|
||||
/// </summary>
|
||||
public ApiWorker(string name)
|
||||
{
|
||||
Name = name;
|
||||
LastActivity = DateTime.Now;
|
||||
|
||||
thread = new Thread(WorkerLoop)
|
||||
{
|
||||
IsBackground = true,
|
||||
Name = name
|
||||
};
|
||||
|
||||
thread.Start();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Enqueues a function returning a value for sequential execution.
|
||||
/// </summary>
|
||||
public Task<T> RunAsync<T>(
|
||||
Func<T> action,
|
||||
CancellationToken token = default(CancellationToken),
|
||||
string operationName = null)
|
||||
{
|
||||
if (action == null)
|
||||
throw new ArgumentNullException(nameof(action));
|
||||
|
||||
if (disposed)
|
||||
throw new ObjectDisposedException(nameof(ApiWorker));
|
||||
|
||||
var tcs = new TaskCompletionSource<T>();
|
||||
|
||||
queue.Add(() =>
|
||||
{
|
||||
if (token.IsCancellationRequested)
|
||||
{
|
||||
tcs.TrySetCanceled();
|
||||
return;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
IsBusy = true;
|
||||
CurrentOperation = operationName ?? action.Method.Name;
|
||||
LastActivity = DateTime.Now;
|
||||
LastError = null;
|
||||
|
||||
var result = action();
|
||||
tcs.TrySetResult(result);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
LastError = ex.Message;
|
||||
tcs.TrySetException(ex);
|
||||
}
|
||||
finally
|
||||
{
|
||||
IsBusy = false;
|
||||
CurrentOperation = null;
|
||||
LastActivity = DateTime.Now;
|
||||
}
|
||||
}, token);
|
||||
|
||||
return tcs.Task;
|
||||
}
|
||||
|
||||
public Task RunAsync(
|
||||
Action action,
|
||||
CancellationToken token = default(CancellationToken),
|
||||
string operationName = null)
|
||||
{
|
||||
return RunAsync<object>(() =>
|
||||
{
|
||||
action();
|
||||
return null;
|
||||
}, token, operationName);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Enqueues a void action for sequential execution.
|
||||
/// </summary>
|
||||
public Task RunAsync(Action action, CancellationToken token = default)
|
||||
{
|
||||
return RunAsync<object>(() =>
|
||||
{
|
||||
action();
|
||||
return null;
|
||||
}, token);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Main worker loop processing queued actions.
|
||||
/// </summary>
|
||||
private void WorkerLoop()
|
||||
{
|
||||
foreach (var item in queue.GetConsumingEnumerable())
|
||||
{
|
||||
item();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Stops the worker and releases resources.
|
||||
/// </summary>
|
||||
public void Dispose()
|
||||
{
|
||||
if (disposed)
|
||||
return;
|
||||
|
||||
disposed = true;
|
||||
|
||||
queue.CompleteAdding();
|
||||
queue.Dispose();
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user