Codingstyle nullable
This commit is contained in:
+133
-138
@@ -1,141 +1,136 @@
|
||||
namespace Swan.Threading
|
||||
{
|
||||
using System;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
using System;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Swan.Threading {
|
||||
/// <summary>
|
||||
/// Represents logic providing several delay mechanisms.
|
||||
/// </summary>
|
||||
/// <example>
|
||||
/// The following example shows how to implement delay mechanisms.
|
||||
/// <code>
|
||||
/// using Swan.Threading;
|
||||
///
|
||||
/// public class Example
|
||||
/// {
|
||||
/// public static void Main()
|
||||
/// {
|
||||
/// // using the ThreadSleep strategy
|
||||
/// using (var delay = new DelayProvider(DelayProvider.DelayStrategy.ThreadSleep))
|
||||
/// {
|
||||
/// // retrieve how much time was delayed
|
||||
/// var time = delay.WaitOne();
|
||||
/// }
|
||||
/// }
|
||||
/// }
|
||||
/// </code>
|
||||
/// </example>
|
||||
public sealed class DelayProvider : IDisposable {
|
||||
private readonly Object _syncRoot = new Object();
|
||||
private readonly Stopwatch _delayStopwatch = new Stopwatch();
|
||||
|
||||
private Boolean _isDisposed;
|
||||
private IWaitEvent _delayEvent;
|
||||
|
||||
/// <summary>
|
||||
/// Represents logic providing several delay mechanisms.
|
||||
/// Initializes a new instance of the <see cref="DelayProvider"/> class.
|
||||
/// </summary>
|
||||
/// <example>
|
||||
/// The following example shows how to implement delay mechanisms.
|
||||
/// <code>
|
||||
/// using Swan.Threading;
|
||||
///
|
||||
/// public class Example
|
||||
/// {
|
||||
/// public static void Main()
|
||||
/// {
|
||||
/// // using the ThreadSleep strategy
|
||||
/// using (var delay = new DelayProvider(DelayProvider.DelayStrategy.ThreadSleep))
|
||||
/// {
|
||||
/// // retrieve how much time was delayed
|
||||
/// var time = delay.WaitOne();
|
||||
/// }
|
||||
/// }
|
||||
/// }
|
||||
/// </code>
|
||||
/// </example>
|
||||
public sealed class DelayProvider : IDisposable
|
||||
{
|
||||
private readonly object _syncRoot = new object();
|
||||
private readonly Stopwatch _delayStopwatch = new Stopwatch();
|
||||
|
||||
private bool _isDisposed;
|
||||
private IWaitEvent _delayEvent;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="DelayProvider"/> class.
|
||||
/// </summary>
|
||||
/// <param name="strategy">The strategy.</param>
|
||||
public DelayProvider(DelayStrategy strategy = DelayStrategy.TaskDelay)
|
||||
{
|
||||
Strategy = strategy;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Enumerates the different ways of providing delays.
|
||||
/// </summary>
|
||||
public enum DelayStrategy
|
||||
{
|
||||
/// <summary>
|
||||
/// Using the Thread.Sleep(15) mechanism.
|
||||
/// </summary>
|
||||
ThreadSleep,
|
||||
|
||||
/// <summary>
|
||||
/// Using the Task.Delay(1).Wait mechanism.
|
||||
/// </summary>
|
||||
TaskDelay,
|
||||
|
||||
/// <summary>
|
||||
/// Using a wait event that completes in a background ThreadPool thread.
|
||||
/// </summary>
|
||||
ThreadPool,
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the selected delay strategy.
|
||||
/// </summary>
|
||||
public DelayStrategy Strategy { get; }
|
||||
|
||||
/// <summary>
|
||||
/// Creates the smallest possible, synchronous delay based on the selected strategy.
|
||||
/// </summary>
|
||||
/// <returns>The elapsed time of the delay.</returns>
|
||||
public TimeSpan WaitOne()
|
||||
{
|
||||
lock (_syncRoot)
|
||||
{
|
||||
if (_isDisposed) return TimeSpan.Zero;
|
||||
|
||||
_delayStopwatch.Restart();
|
||||
|
||||
switch (Strategy)
|
||||
{
|
||||
case DelayStrategy.ThreadSleep:
|
||||
DelaySleep();
|
||||
break;
|
||||
case DelayStrategy.TaskDelay:
|
||||
DelayTask();
|
||||
break;
|
||||
case DelayStrategy.ThreadPool:
|
||||
DelayThreadPool();
|
||||
break;
|
||||
}
|
||||
|
||||
return _delayStopwatch.Elapsed;
|
||||
}
|
||||
}
|
||||
|
||||
#region Dispose Pattern
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Dispose()
|
||||
{
|
||||
lock (_syncRoot)
|
||||
{
|
||||
if (_isDisposed) return;
|
||||
_isDisposed = true;
|
||||
|
||||
_delayEvent?.Dispose();
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region Private Delay Mechanisms
|
||||
|
||||
private static void DelaySleep() => Thread.Sleep(15);
|
||||
|
||||
private static void DelayTask() => Task.Delay(1).Wait();
|
||||
|
||||
private void DelayThreadPool()
|
||||
{
|
||||
if (_delayEvent == null)
|
||||
_delayEvent = WaitEventFactory.Create(isCompleted: true, useSlim: true);
|
||||
|
||||
_delayEvent.Begin();
|
||||
ThreadPool.QueueUserWorkItem(s =>
|
||||
{
|
||||
DelaySleep();
|
||||
_delayEvent.Complete();
|
||||
});
|
||||
|
||||
_delayEvent.Wait();
|
||||
}
|
||||
|
||||
#endregion
|
||||
}
|
||||
/// <param name="strategy">The strategy.</param>
|
||||
public DelayProvider(DelayStrategy strategy = DelayStrategy.TaskDelay) => this.Strategy = strategy;
|
||||
|
||||
/// <summary>
|
||||
/// Enumerates the different ways of providing delays.
|
||||
/// </summary>
|
||||
public enum DelayStrategy {
|
||||
/// <summary>
|
||||
/// Using the Thread.Sleep(15) mechanism.
|
||||
/// </summary>
|
||||
ThreadSleep,
|
||||
|
||||
/// <summary>
|
||||
/// Using the Task.Delay(1).Wait mechanism.
|
||||
/// </summary>
|
||||
TaskDelay,
|
||||
|
||||
/// <summary>
|
||||
/// Using a wait event that completes in a background ThreadPool thread.
|
||||
/// </summary>
|
||||
ThreadPool,
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the selected delay strategy.
|
||||
/// </summary>
|
||||
public DelayStrategy Strategy {
|
||||
get;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Creates the smallest possible, synchronous delay based on the selected strategy.
|
||||
/// </summary>
|
||||
/// <returns>The elapsed time of the delay.</returns>
|
||||
public TimeSpan WaitOne() {
|
||||
lock(this._syncRoot) {
|
||||
if(this._isDisposed) {
|
||||
return TimeSpan.Zero;
|
||||
}
|
||||
|
||||
this._delayStopwatch.Restart();
|
||||
|
||||
switch(this.Strategy) {
|
||||
case DelayStrategy.ThreadSleep:
|
||||
DelaySleep();
|
||||
break;
|
||||
case DelayStrategy.TaskDelay:
|
||||
DelayTask();
|
||||
break;
|
||||
case DelayStrategy.ThreadPool:
|
||||
this.DelayThreadPool();
|
||||
break;
|
||||
}
|
||||
|
||||
return this._delayStopwatch.Elapsed;
|
||||
}
|
||||
}
|
||||
|
||||
#region Dispose Pattern
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Dispose() {
|
||||
lock(this._syncRoot) {
|
||||
if(this._isDisposed) {
|
||||
return;
|
||||
}
|
||||
|
||||
this._isDisposed = true;
|
||||
|
||||
this._delayEvent?.Dispose();
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region Private Delay Mechanisms
|
||||
|
||||
private static void DelaySleep() => Thread.Sleep(15);
|
||||
|
||||
private static void DelayTask() => Task.Delay(1).Wait();
|
||||
|
||||
private void DelayThreadPool() {
|
||||
if(this._delayEvent == null) {
|
||||
this._delayEvent = WaitEventFactory.Create(isCompleted: true, useSlim: true);
|
||||
}
|
||||
|
||||
this._delayEvent.Begin();
|
||||
_ = ThreadPool.QueueUserWorkItem(s => {
|
||||
DelaySleep();
|
||||
this._delayEvent.Complete();
|
||||
});
|
||||
|
||||
this._delayEvent.Wait();
|
||||
}
|
||||
|
||||
#endregion
|
||||
}
|
||||
}
|
||||
+245
-286
@@ -1,292 +1,251 @@
|
||||
namespace Swan.Threading
|
||||
{
|
||||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Swan.Threading {
|
||||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
/// <summary>
|
||||
/// Provides a base implementation for application workers
|
||||
/// that perform continuous, long-running tasks. This class
|
||||
/// provides the ability to perform fine-grained control on these tasks.
|
||||
/// </summary>
|
||||
/// <seealso cref="IWorker" />
|
||||
public abstract class ThreadWorkerBase : WorkerBase {
|
||||
private readonly Object _syncLock = new Object();
|
||||
private readonly Thread _thread;
|
||||
|
||||
/// <summary>
|
||||
/// Provides a base implementation for application workers
|
||||
/// that perform continuous, long-running tasks. This class
|
||||
/// provides the ability to perform fine-grained control on these tasks.
|
||||
/// Initializes a new instance of the <see cref="ThreadWorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <seealso cref="IWorker" />
|
||||
public abstract class ThreadWorkerBase : WorkerBase
|
||||
{
|
||||
private readonly object _syncLock = new object();
|
||||
private readonly Thread _thread;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="ThreadWorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="priority">The thread priority.</param>
|
||||
/// <param name="period">The interval of cycle execution.</param>
|
||||
/// <param name="delayProvider">The cycle delay provide implementation.</param>
|
||||
protected ThreadWorkerBase(string name, ThreadPriority priority, TimeSpan period, IWorkerDelayProvider delayProvider)
|
||||
: base(name, period)
|
||||
{
|
||||
DelayProvider = delayProvider;
|
||||
_thread = new Thread(RunWorkerLoop)
|
||||
{
|
||||
IsBackground = true,
|
||||
Priority = priority,
|
||||
Name = name,
|
||||
};
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="ThreadWorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="period">The execution interval.</param>
|
||||
protected ThreadWorkerBase(string name, TimeSpan period)
|
||||
: this(name, ThreadPriority.Normal, period, WorkerDelayProvider.Default)
|
||||
{
|
||||
// placeholder
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Provides an implementation on a cycle delay provider.
|
||||
/// </summary>
|
||||
protected IWorkerDelayProvider DelayProvider { get; }
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StartAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Paused || WorkerState == WorkerState.Waiting)
|
||||
return ResumeAsync();
|
||||
|
||||
if (WorkerState != WorkerState.Created)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
if (IsStopRequested)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
var task = QueueStateChange(StateChangeRequest.Start);
|
||||
_thread.Start();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> PauseAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState != WorkerState.Running && WorkerState != WorkerState.Waiting)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
return IsStopRequested ? Task.FromResult(WorkerState) : QueueStateChange(StateChangeRequest.Pause);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> ResumeAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Created)
|
||||
return StartAsync();
|
||||
|
||||
if (WorkerState != WorkerState.Paused && WorkerState != WorkerState.Waiting)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
return IsStopRequested ? Task.FromResult(WorkerState) : QueueStateChange(StateChangeRequest.Resume);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StopAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Stopped || WorkerState == WorkerState.Created)
|
||||
{
|
||||
WorkerState = WorkerState.Stopped;
|
||||
return Task.FromResult(WorkerState);
|
||||
}
|
||||
|
||||
return QueueStateChange(StateChangeRequest.Stop);
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Suspends execution queues a new new cycle for execution. The delay is given in
|
||||
/// milliseconds. When overridden in a derived class the wait handle will be set
|
||||
/// whenever an interrupt is received.
|
||||
/// </summary>
|
||||
/// <param name="wantedDelay">The remaining delay to wait for in the cycle.</param>
|
||||
/// <param name="delayTask">Contains a reference to a task with the scheduled period delay.</param>
|
||||
/// <param name="token">The cancellation token to cancel waiting.</param>
|
||||
protected virtual void ExecuteCycleDelay(int wantedDelay, Task delayTask, CancellationToken token) =>
|
||||
DelayProvider?.ExecuteCycleDelay(wantedDelay, delayTask, token);
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override void OnDisposing()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if ((_thread.ThreadState & ThreadState.Unstarted) != ThreadState.Unstarted)
|
||||
_thread.Join();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Implements worker control, execution and delay logic in a loop.
|
||||
/// </summary>
|
||||
private void RunWorkerLoop()
|
||||
{
|
||||
while (WorkerState != WorkerState.Stopped && !IsDisposing && !IsDisposed)
|
||||
{
|
||||
CycleStopwatch.Restart();
|
||||
var interruptToken = CycleCancellation.Token;
|
||||
var period = Period.TotalMilliseconds >= int.MaxValue ? -1 : Convert.ToInt32(Math.Floor(Period.TotalMilliseconds));
|
||||
var delayTask = Task.Delay(period, interruptToken);
|
||||
var initialWorkerState = WorkerState;
|
||||
|
||||
// Lock the cycle and capture relevant state valid for this cycle
|
||||
CycleCompletedEvent.Reset();
|
||||
|
||||
// Process the tasks that are awaiting
|
||||
if (ProcessStateChangeRequests())
|
||||
continue;
|
||||
|
||||
try
|
||||
{
|
||||
if (initialWorkerState == WorkerState.Waiting &&
|
||||
!interruptToken.IsCancellationRequested)
|
||||
{
|
||||
// Mark the state as Running
|
||||
WorkerState = WorkerState.Running;
|
||||
|
||||
// Call the execution logic
|
||||
ExecuteCycleLogic(interruptToken);
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
OnCycleException(ex);
|
||||
}
|
||||
finally
|
||||
{
|
||||
// Update the state
|
||||
WorkerState = initialWorkerState == WorkerState.Paused
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="priority">The thread priority.</param>
|
||||
/// <param name="period">The interval of cycle execution.</param>
|
||||
/// <param name="delayProvider">The cycle delay provide implementation.</param>
|
||||
protected ThreadWorkerBase(String name, ThreadPriority priority, TimeSpan period, IWorkerDelayProvider delayProvider) : base(name, period) {
|
||||
this.DelayProvider = delayProvider;
|
||||
this._thread = new Thread(this.RunWorkerLoop) {
|
||||
IsBackground = true,
|
||||
Priority = priority,
|
||||
Name = name,
|
||||
};
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="ThreadWorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="period">The execution interval.</param>
|
||||
protected ThreadWorkerBase(String name, TimeSpan period) : this(name, ThreadPriority.Normal, period, WorkerDelayProvider.Default) {
|
||||
// placeholder
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Provides an implementation on a cycle delay provider.
|
||||
/// </summary>
|
||||
protected IWorkerDelayProvider DelayProvider {
|
||||
get;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StartAsync() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState == WorkerState.Paused || this.WorkerState == WorkerState.Waiting) {
|
||||
return this.ResumeAsync();
|
||||
}
|
||||
|
||||
if(this.WorkerState != WorkerState.Created) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
if(this.IsStopRequested) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
Task<WorkerState> task = this.QueueStateChange(StateChangeRequest.Start);
|
||||
this._thread.Start();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> PauseAsync() {
|
||||
lock(this._syncLock) {
|
||||
return this.WorkerState != WorkerState.Running && this.WorkerState != WorkerState.Waiting ? Task.FromResult(this.WorkerState) : this.IsStopRequested ? Task.FromResult(this.WorkerState) : this.QueueStateChange(StateChangeRequest.Pause);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> ResumeAsync() {
|
||||
lock(this._syncLock) {
|
||||
return this.WorkerState == WorkerState.Created ? this.StartAsync() : this.WorkerState != WorkerState.Paused && this.WorkerState != WorkerState.Waiting ? Task.FromResult(this.WorkerState) : this.IsStopRequested ? Task.FromResult(this.WorkerState) : this.QueueStateChange(StateChangeRequest.Resume);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StopAsync() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState == WorkerState.Stopped || this.WorkerState == WorkerState.Created) {
|
||||
this.WorkerState = WorkerState.Stopped;
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
return this.QueueStateChange(StateChangeRequest.Stop);
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Suspends execution queues a new new cycle for execution. The delay is given in
|
||||
/// milliseconds. When overridden in a derived class the wait handle will be set
|
||||
/// whenever an interrupt is received.
|
||||
/// </summary>
|
||||
/// <param name="wantedDelay">The remaining delay to wait for in the cycle.</param>
|
||||
/// <param name="delayTask">Contains a reference to a task with the scheduled period delay.</param>
|
||||
/// <param name="token">The cancellation token to cancel waiting.</param>
|
||||
protected virtual void ExecuteCycleDelay(Int32 wantedDelay, Task delayTask, CancellationToken token) =>
|
||||
this.DelayProvider?.ExecuteCycleDelay(wantedDelay, delayTask, token);
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override void OnDisposing() {
|
||||
lock(this._syncLock) {
|
||||
if((this._thread.ThreadState & ThreadState.Unstarted) != ThreadState.Unstarted) {
|
||||
this._thread.Join();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Implements worker control, execution and delay logic in a loop.
|
||||
/// </summary>
|
||||
private void RunWorkerLoop() {
|
||||
while(this.WorkerState != WorkerState.Stopped && !this.IsDisposing && !this.IsDisposed) {
|
||||
this.CycleStopwatch.Restart();
|
||||
CancellationToken interruptToken = this.CycleCancellation.Token;
|
||||
Int32 period = this.Period.TotalMilliseconds >= Int32.MaxValue ? -1 : Convert.ToInt32(Math.Floor(this.Period.TotalMilliseconds));
|
||||
Task delayTask = Task.Delay(period, interruptToken);
|
||||
WorkerState initialWorkerState = this.WorkerState;
|
||||
|
||||
// Lock the cycle and capture relevant state valid for this cycle
|
||||
this.CycleCompletedEvent.Reset();
|
||||
|
||||
// Process the tasks that are awaiting
|
||||
if(this.ProcessStateChangeRequests()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
try {
|
||||
if(initialWorkerState == WorkerState.Waiting &&
|
||||
!interruptToken.IsCancellationRequested) {
|
||||
// Mark the state as Running
|
||||
this.WorkerState = WorkerState.Running;
|
||||
|
||||
// Call the execution logic
|
||||
this.ExecuteCycleLogic(interruptToken);
|
||||
}
|
||||
} catch(Exception ex) {
|
||||
this.OnCycleException(ex);
|
||||
} finally {
|
||||
// Update the state
|
||||
this.WorkerState = initialWorkerState == WorkerState.Paused
|
||||
? WorkerState.Paused
|
||||
: WorkerState.Waiting;
|
||||
|
||||
// Signal the cycle has been completed so new cycles can be executed
|
||||
CycleCompletedEvent.Set();
|
||||
|
||||
if (!interruptToken.IsCancellationRequested)
|
||||
{
|
||||
var cycleDelay = ComputeCycleDelay(initialWorkerState);
|
||||
if (cycleDelay == Timeout.Infinite)
|
||||
delayTask = Task.Delay(Timeout.Infinite, interruptToken);
|
||||
|
||||
ExecuteCycleDelay(
|
||||
: WorkerState.Waiting;
|
||||
|
||||
// Signal the cycle has been completed so new cycles can be executed
|
||||
this.CycleCompletedEvent.Set();
|
||||
|
||||
if(!interruptToken.IsCancellationRequested) {
|
||||
Int32 cycleDelay = this.ComputeCycleDelay(initialWorkerState);
|
||||
if(cycleDelay == Timeout.Infinite) {
|
||||
delayTask = Task.Delay(Timeout.Infinite, interruptToken);
|
||||
}
|
||||
|
||||
this.ExecuteCycleDelay(
|
||||
cycleDelay,
|
||||
delayTask,
|
||||
CycleCancellation.Token);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ClearStateChangeRequests();
|
||||
WorkerState = WorkerState.Stopped;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Queues a transition in worker state for processing. Returns a task that can be awaited
|
||||
/// when the operation completes.
|
||||
/// </summary>
|
||||
/// <param name="request">The request.</param>
|
||||
/// <returns>The awaitable task.</returns>
|
||||
private Task<WorkerState> QueueStateChange(StateChangeRequest request)
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (StateChangeTask != null)
|
||||
return StateChangeTask;
|
||||
|
||||
var waitingTask = new Task<WorkerState>(() =>
|
||||
{
|
||||
StateChangedEvent.Wait();
|
||||
lock (_syncLock)
|
||||
{
|
||||
StateChangeTask = null;
|
||||
return WorkerState;
|
||||
}
|
||||
});
|
||||
|
||||
StateChangeTask = waitingTask;
|
||||
StateChangedEvent.Reset();
|
||||
StateChangeRequests[request] = true;
|
||||
waitingTask.Start();
|
||||
CycleCancellation.Cancel();
|
||||
|
||||
return waitingTask;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Processes the state change request by checking pending events and scheduling
|
||||
/// cycle execution accordingly. The <see cref="WorkerState"/> is also updated.
|
||||
/// </summary>
|
||||
/// <returns>Returns <c>true</c> if the execution should be terminated. <c>false</c> otherwise.</returns>
|
||||
private bool ProcessStateChangeRequests()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
var hasRequest = false;
|
||||
var currentState = WorkerState;
|
||||
|
||||
// Update the state in the given priority
|
||||
if (StateChangeRequests[StateChangeRequest.Stop] || IsDisposing || IsDisposed)
|
||||
{
|
||||
hasRequest = true;
|
||||
WorkerState = WorkerState.Stopped;
|
||||
}
|
||||
else if (StateChangeRequests[StateChangeRequest.Pause])
|
||||
{
|
||||
hasRequest = true;
|
||||
WorkerState = WorkerState.Paused;
|
||||
}
|
||||
else if (StateChangeRequests[StateChangeRequest.Start] || StateChangeRequests[StateChangeRequest.Resume])
|
||||
{
|
||||
hasRequest = true;
|
||||
WorkerState = WorkerState.Waiting;
|
||||
}
|
||||
|
||||
// Signals all state changes to continue
|
||||
// as a command has been handled.
|
||||
if (hasRequest)
|
||||
{
|
||||
ClearStateChangeRequests();
|
||||
OnStateChangeProcessed(currentState, WorkerState);
|
||||
}
|
||||
|
||||
return hasRequest;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Signals all state change requests to set.
|
||||
/// </summary>
|
||||
private void ClearStateChangeRequests()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
// Mark all events as completed
|
||||
StateChangeRequests[StateChangeRequest.Start] = false;
|
||||
StateChangeRequests[StateChangeRequest.Pause] = false;
|
||||
StateChangeRequests[StateChangeRequest.Resume] = false;
|
||||
StateChangeRequests[StateChangeRequest.Stop] = false;
|
||||
|
||||
StateChangedEvent.Set();
|
||||
CycleCompletedEvent.Set();
|
||||
}
|
||||
}
|
||||
}
|
||||
this.CycleCancellation.Token);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
this.ClearStateChangeRequests();
|
||||
this.WorkerState = WorkerState.Stopped;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Queues a transition in worker state for processing. Returns a task that can be awaited
|
||||
/// when the operation completes.
|
||||
/// </summary>
|
||||
/// <param name="request">The request.</param>
|
||||
/// <returns>The awaitable task.</returns>
|
||||
private Task<WorkerState> QueueStateChange(StateChangeRequest request) {
|
||||
lock(this._syncLock) {
|
||||
if(this.StateChangeTask != null) {
|
||||
return this.StateChangeTask;
|
||||
}
|
||||
|
||||
Task<WorkerState> waitingTask = new Task<WorkerState>(() => {
|
||||
this.StateChangedEvent.Wait();
|
||||
lock(this._syncLock) {
|
||||
this.StateChangeTask = null;
|
||||
return this.WorkerState;
|
||||
}
|
||||
});
|
||||
|
||||
this.StateChangeTask = waitingTask;
|
||||
this.StateChangedEvent.Reset();
|
||||
this.StateChangeRequests[request] = true;
|
||||
waitingTask.Start();
|
||||
this.CycleCancellation.Cancel();
|
||||
|
||||
return waitingTask;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Processes the state change request by checking pending events and scheduling
|
||||
/// cycle execution accordingly. The <see cref="WorkerState"/> is also updated.
|
||||
/// </summary>
|
||||
/// <returns>Returns <c>true</c> if the execution should be terminated. <c>false</c> otherwise.</returns>
|
||||
private Boolean ProcessStateChangeRequests() {
|
||||
lock(this._syncLock) {
|
||||
Boolean hasRequest = false;
|
||||
WorkerState currentState = this.WorkerState;
|
||||
|
||||
// Update the state in the given priority
|
||||
if(this.StateChangeRequests[StateChangeRequest.Stop] || this.IsDisposing || this.IsDisposed) {
|
||||
hasRequest = true;
|
||||
this.WorkerState = WorkerState.Stopped;
|
||||
} else if(this.StateChangeRequests[StateChangeRequest.Pause]) {
|
||||
hasRequest = true;
|
||||
this.WorkerState = WorkerState.Paused;
|
||||
} else if(this.StateChangeRequests[StateChangeRequest.Start] || this.StateChangeRequests[StateChangeRequest.Resume]) {
|
||||
hasRequest = true;
|
||||
this.WorkerState = WorkerState.Waiting;
|
||||
}
|
||||
|
||||
// Signals all state changes to continue
|
||||
// as a command has been handled.
|
||||
if(hasRequest) {
|
||||
this.ClearStateChangeRequests();
|
||||
this.OnStateChangeProcessed(currentState, this.WorkerState);
|
||||
}
|
||||
|
||||
return hasRequest;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Signals all state change requests to set.
|
||||
/// </summary>
|
||||
private void ClearStateChangeRequests() {
|
||||
lock(this._syncLock) {
|
||||
// Mark all events as completed
|
||||
this.StateChangeRequests[StateChangeRequest.Start] = false;
|
||||
this.StateChangeRequests[StateChangeRequest.Pause] = false;
|
||||
this.StateChangeRequests[StateChangeRequest.Resume] = false;
|
||||
this.StateChangeRequests[StateChangeRequest.Stop] = false;
|
||||
|
||||
this.StateChangedEvent.Set();
|
||||
this.CycleCompletedEvent.Set();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+296
-324
@@ -1,328 +1,300 @@
|
||||
namespace Swan.Threading
|
||||
{
|
||||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
using System;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Swan.Threading {
|
||||
/// <summary>
|
||||
/// Provides a base implementation for application workers.
|
||||
/// </summary>
|
||||
/// <seealso cref="IWorker" />
|
||||
public abstract class TimerWorkerBase : WorkerBase {
|
||||
private readonly Object _syncLock = new Object();
|
||||
private readonly Timer _timer;
|
||||
private Boolean _isTimerAlive = true;
|
||||
|
||||
/// <summary>
|
||||
/// Provides a base implementation for application workers.
|
||||
/// Initializes a new instance of the <see cref="TimerWorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <seealso cref="IWorker" />
|
||||
public abstract class TimerWorkerBase : WorkerBase
|
||||
{
|
||||
private readonly object _syncLock = new object();
|
||||
private readonly Timer _timer;
|
||||
private bool _isTimerAlive = true;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="TimerWorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="period">The execution interval.</param>
|
||||
protected TimerWorkerBase(string name, TimeSpan period)
|
||||
: base(name, period)
|
||||
{
|
||||
// Instantiate the timer that will be used to schedule cycles
|
||||
_timer = new Timer(
|
||||
ExecuteTimerCallback,
|
||||
this,
|
||||
Timeout.Infinite,
|
||||
Timeout.Infinite);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StartAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Paused || WorkerState == WorkerState.Waiting)
|
||||
return ResumeAsync();
|
||||
|
||||
if (WorkerState != WorkerState.Created)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
if (IsStopRequested)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
var task = QueueStateChange(StateChangeRequest.Start);
|
||||
Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> PauseAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState != WorkerState.Running && WorkerState != WorkerState.Waiting)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
if (IsStopRequested)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
var task = QueueStateChange(StateChangeRequest.Pause);
|
||||
Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> ResumeAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Created)
|
||||
return StartAsync();
|
||||
|
||||
if (WorkerState != WorkerState.Paused && WorkerState != WorkerState.Waiting)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
if (IsStopRequested)
|
||||
return Task.FromResult(WorkerState);
|
||||
|
||||
var task = QueueStateChange(StateChangeRequest.Resume);
|
||||
Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StopAsync()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Stopped || WorkerState == WorkerState.Created)
|
||||
{
|
||||
WorkerState = WorkerState.Stopped;
|
||||
return Task.FromResult(WorkerState);
|
||||
}
|
||||
|
||||
var task = QueueStateChange(StateChangeRequest.Stop);
|
||||
Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Schedules a new cycle for execution. The delay is given in
|
||||
/// milliseconds. Passing a delay of 0 means a new cycle should be executed
|
||||
/// immediately.
|
||||
/// </summary>
|
||||
/// <param name="delay">The delay.</param>
|
||||
protected void ScheduleCycle(int delay)
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (!_isTimerAlive) return;
|
||||
_timer.Change(delay, Timeout.Infinite);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override void Dispose(bool disposing)
|
||||
{
|
||||
base.Dispose(disposing);
|
||||
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (!_isTimerAlive) return;
|
||||
_isTimerAlive = false;
|
||||
_timer.Dispose();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Cancels the current token and schedules a new cycle immediately.
|
||||
/// </summary>
|
||||
private void Interrupt()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (WorkerState == WorkerState.Stopped)
|
||||
return;
|
||||
|
||||
CycleCancellation.Cancel();
|
||||
ScheduleCycle(0);
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Executes the worker cycle control logic.
|
||||
/// This includes processing state change requests,
|
||||
/// the execution of use cycle code,
|
||||
/// and the scheduling of new cycles.
|
||||
/// </summary>
|
||||
private void ExecuteWorkerCycle()
|
||||
{
|
||||
CycleStopwatch.Restart();
|
||||
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (IsDisposing || IsDisposed)
|
||||
{
|
||||
WorkerState = WorkerState.Stopped;
|
||||
|
||||
// Cancel any awaiters
|
||||
try { StateChangedEvent.Set(); }
|
||||
catch { /* Ignore */ }
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
// Prevent running another instance of the cycle
|
||||
if (CycleCompletedEvent.IsSet == false) return;
|
||||
|
||||
// Lock the cycle and capture relevant state valid for this cycle
|
||||
CycleCompletedEvent.Reset();
|
||||
}
|
||||
|
||||
var interruptToken = CycleCancellation.Token;
|
||||
var initialWorkerState = WorkerState;
|
||||
|
||||
// Process the tasks that are awaiting
|
||||
if (ProcessStateChangeRequests())
|
||||
return;
|
||||
|
||||
try
|
||||
{
|
||||
if (initialWorkerState == WorkerState.Waiting &&
|
||||
!interruptToken.IsCancellationRequested)
|
||||
{
|
||||
// Mark the state as Running
|
||||
WorkerState = WorkerState.Running;
|
||||
|
||||
// Call the execution logic
|
||||
ExecuteCycleLogic(interruptToken);
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
OnCycleException(ex);
|
||||
}
|
||||
finally
|
||||
{
|
||||
// Update the state
|
||||
WorkerState = initialWorkerState == WorkerState.Paused
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="period">The execution interval.</param>
|
||||
protected TimerWorkerBase(String name, TimeSpan period) : base(name, period) =>
|
||||
// Instantiate the timer that will be used to schedule cycles
|
||||
this._timer = new Timer(this.ExecuteTimerCallback, this, Timeout.Infinite, Timeout.Infinite);
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StartAsync() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState == WorkerState.Paused || this.WorkerState == WorkerState.Waiting) {
|
||||
return this.ResumeAsync();
|
||||
}
|
||||
|
||||
if(this.WorkerState != WorkerState.Created) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
if(this.IsStopRequested) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
Task<WorkerState> task = this.QueueStateChange(StateChangeRequest.Start);
|
||||
this.Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> PauseAsync() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState != WorkerState.Running && this.WorkerState != WorkerState.Waiting) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
if(this.IsStopRequested) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
Task<WorkerState> task = this.QueueStateChange(StateChangeRequest.Pause);
|
||||
this.Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> ResumeAsync() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState == WorkerState.Created) {
|
||||
return this.StartAsync();
|
||||
}
|
||||
|
||||
if(this.WorkerState != WorkerState.Paused && this.WorkerState != WorkerState.Waiting) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
if(this.IsStopRequested) {
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
Task<WorkerState> task = this.QueueStateChange(StateChangeRequest.Resume);
|
||||
this.Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public override Task<WorkerState> StopAsync() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState == WorkerState.Stopped || this.WorkerState == WorkerState.Created) {
|
||||
this.WorkerState = WorkerState.Stopped;
|
||||
return Task.FromResult(this.WorkerState);
|
||||
}
|
||||
|
||||
Task<WorkerState> task = this.QueueStateChange(StateChangeRequest.Stop);
|
||||
this.Interrupt();
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Schedules a new cycle for execution. The delay is given in
|
||||
/// milliseconds. Passing a delay of 0 means a new cycle should be executed
|
||||
/// immediately.
|
||||
/// </summary>
|
||||
/// <param name="delay">The delay.</param>
|
||||
protected void ScheduleCycle(Int32 delay) {
|
||||
lock(this._syncLock) {
|
||||
if(!this._isTimerAlive) {
|
||||
return;
|
||||
}
|
||||
|
||||
_ = this._timer.Change(delay, Timeout.Infinite);
|
||||
}
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override void Dispose(Boolean disposing) {
|
||||
base.Dispose(disposing);
|
||||
|
||||
lock(this._syncLock) {
|
||||
if(!this._isTimerAlive) {
|
||||
return;
|
||||
}
|
||||
|
||||
this._isTimerAlive = false;
|
||||
this._timer.Dispose();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Cancels the current token and schedules a new cycle immediately.
|
||||
/// </summary>
|
||||
private void Interrupt() {
|
||||
lock(this._syncLock) {
|
||||
if(this.WorkerState == WorkerState.Stopped) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.CycleCancellation.Cancel();
|
||||
this.ScheduleCycle(0);
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Executes the worker cycle control logic.
|
||||
/// This includes processing state change requests,
|
||||
/// the execution of use cycle code,
|
||||
/// and the scheduling of new cycles.
|
||||
/// </summary>
|
||||
private void ExecuteWorkerCycle() {
|
||||
this.CycleStopwatch.Restart();
|
||||
|
||||
lock(this._syncLock) {
|
||||
if(this.IsDisposing || this.IsDisposed) {
|
||||
this.WorkerState = WorkerState.Stopped;
|
||||
|
||||
// Cancel any awaiters
|
||||
try {
|
||||
this.StateChangedEvent.Set();
|
||||
} catch { /* Ignore */ }
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
// Prevent running another instance of the cycle
|
||||
if(this.CycleCompletedEvent.IsSet == false) {
|
||||
return;
|
||||
}
|
||||
|
||||
// Lock the cycle and capture relevant state valid for this cycle
|
||||
this.CycleCompletedEvent.Reset();
|
||||
}
|
||||
|
||||
CancellationToken interruptToken = this.CycleCancellation.Token;
|
||||
WorkerState initialWorkerState = this.WorkerState;
|
||||
|
||||
// Process the tasks that are awaiting
|
||||
if(this.ProcessStateChangeRequests()) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
if(initialWorkerState == WorkerState.Waiting &&
|
||||
!interruptToken.IsCancellationRequested) {
|
||||
// Mark the state as Running
|
||||
this.WorkerState = WorkerState.Running;
|
||||
|
||||
// Call the execution logic
|
||||
this.ExecuteCycleLogic(interruptToken);
|
||||
}
|
||||
} catch(Exception ex) {
|
||||
this.OnCycleException(ex);
|
||||
} finally {
|
||||
// Update the state
|
||||
this.WorkerState = initialWorkerState == WorkerState.Paused
|
||||
? WorkerState.Paused
|
||||
: WorkerState.Waiting;
|
||||
|
||||
lock (_syncLock)
|
||||
{
|
||||
// Signal the cycle has been completed so new cycles can be executed
|
||||
CycleCompletedEvent.Set();
|
||||
|
||||
// Schedule a new cycle
|
||||
ScheduleCycle(!interruptToken.IsCancellationRequested
|
||||
? ComputeCycleDelay(initialWorkerState)
|
||||
: 0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Represents the callback that is executed when the <see cref="_timer"/> ticks.
|
||||
/// </summary>
|
||||
/// <param name="state">The state -- this contains the worker.</param>
|
||||
private void ExecuteTimerCallback(object state) => ExecuteWorkerCycle();
|
||||
|
||||
/// <summary>
|
||||
/// Queues a transition in worker state for processing. Returns a task that can be awaited
|
||||
/// when the operation completes.
|
||||
/// </summary>
|
||||
/// <param name="request">The request.</param>
|
||||
/// <returns>The awaitable task.</returns>
|
||||
private Task<WorkerState> QueueStateChange(StateChangeRequest request)
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (StateChangeTask != null)
|
||||
return StateChangeTask;
|
||||
|
||||
var waitingTask = new Task<WorkerState>(() =>
|
||||
{
|
||||
StateChangedEvent.Wait();
|
||||
lock (_syncLock)
|
||||
{
|
||||
StateChangeTask = null;
|
||||
return WorkerState;
|
||||
}
|
||||
});
|
||||
|
||||
StateChangeTask = waitingTask;
|
||||
StateChangedEvent.Reset();
|
||||
StateChangeRequests[request] = true;
|
||||
waitingTask.Start();
|
||||
CycleCancellation.Cancel();
|
||||
|
||||
return waitingTask;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Processes the state change queue by checking pending events and scheduling
|
||||
/// cycle execution accordingly. The <see cref="WorkerState"/> is also updated.
|
||||
/// </summary>
|
||||
/// <returns>Returns <c>true</c> if the execution should be terminated. <c>false</c> otherwise.</returns>
|
||||
private bool ProcessStateChangeRequests()
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
var currentState = WorkerState;
|
||||
var hasRequest = false;
|
||||
var schedule = 0;
|
||||
|
||||
// Update the state according to request priority
|
||||
if (StateChangeRequests[StateChangeRequest.Stop] || IsDisposing || IsDisposed)
|
||||
{
|
||||
hasRequest = true;
|
||||
WorkerState = WorkerState.Stopped;
|
||||
schedule = StateChangeRequests[StateChangeRequest.Stop] ? Timeout.Infinite : 0;
|
||||
}
|
||||
else if (StateChangeRequests[StateChangeRequest.Pause])
|
||||
{
|
||||
hasRequest = true;
|
||||
WorkerState = WorkerState.Paused;
|
||||
schedule = Timeout.Infinite;
|
||||
}
|
||||
else if (StateChangeRequests[StateChangeRequest.Start] || StateChangeRequests[StateChangeRequest.Resume])
|
||||
{
|
||||
hasRequest = true;
|
||||
WorkerState = WorkerState.Waiting;
|
||||
}
|
||||
|
||||
// Signals all state changes to continue
|
||||
// as a command has been handled.
|
||||
if (hasRequest)
|
||||
{
|
||||
ClearStateChangeRequests(schedule, currentState, WorkerState);
|
||||
}
|
||||
|
||||
return hasRequest;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Signals all state change requests to set.
|
||||
/// </summary>
|
||||
/// <param name="schedule">The cycle schedule.</param>
|
||||
/// <param name="oldState">The previous worker state.</param>
|
||||
/// <param name="newState">The new worker state.</param>
|
||||
private void ClearStateChangeRequests(int schedule, WorkerState oldState, WorkerState newState)
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
// Mark all events as completed
|
||||
StateChangeRequests[StateChangeRequest.Start] = false;
|
||||
StateChangeRequests[StateChangeRequest.Pause] = false;
|
||||
StateChangeRequests[StateChangeRequest.Resume] = false;
|
||||
StateChangeRequests[StateChangeRequest.Stop] = false;
|
||||
|
||||
StateChangedEvent.Set();
|
||||
CycleCompletedEvent.Set();
|
||||
OnStateChangeProcessed(oldState, newState);
|
||||
ScheduleCycle(schedule);
|
||||
}
|
||||
}
|
||||
}
|
||||
: WorkerState.Waiting;
|
||||
|
||||
lock(this._syncLock) {
|
||||
// Signal the cycle has been completed so new cycles can be executed
|
||||
this.CycleCompletedEvent.Set();
|
||||
|
||||
// Schedule a new cycle
|
||||
this.ScheduleCycle(!interruptToken.IsCancellationRequested
|
||||
? this.ComputeCycleDelay(initialWorkerState)
|
||||
: 0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Represents the callback that is executed when the <see cref="_timer"/> ticks.
|
||||
/// </summary>
|
||||
/// <param name="state">The state -- this contains the worker.</param>
|
||||
private void ExecuteTimerCallback(Object state) => this.ExecuteWorkerCycle();
|
||||
|
||||
/// <summary>
|
||||
/// Queues a transition in worker state for processing. Returns a task that can be awaited
|
||||
/// when the operation completes.
|
||||
/// </summary>
|
||||
/// <param name="request">The request.</param>
|
||||
/// <returns>The awaitable task.</returns>
|
||||
private Task<WorkerState> QueueStateChange(StateChangeRequest request) {
|
||||
lock(this._syncLock) {
|
||||
if(this.StateChangeTask != null) {
|
||||
return this.StateChangeTask;
|
||||
}
|
||||
|
||||
Task<WorkerState> waitingTask = new Task<WorkerState>(() => {
|
||||
this.StateChangedEvent.Wait();
|
||||
lock(this._syncLock) {
|
||||
this.StateChangeTask = null;
|
||||
return this.WorkerState;
|
||||
}
|
||||
});
|
||||
|
||||
this.StateChangeTask = waitingTask;
|
||||
this.StateChangedEvent.Reset();
|
||||
this.StateChangeRequests[request] = true;
|
||||
waitingTask.Start();
|
||||
this.CycleCancellation.Cancel();
|
||||
|
||||
return waitingTask;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Processes the state change queue by checking pending events and scheduling
|
||||
/// cycle execution accordingly. The <see cref="WorkerState"/> is also updated.
|
||||
/// </summary>
|
||||
/// <returns>Returns <c>true</c> if the execution should be terminated. <c>false</c> otherwise.</returns>
|
||||
private Boolean ProcessStateChangeRequests() {
|
||||
lock(this._syncLock) {
|
||||
WorkerState currentState = this.WorkerState;
|
||||
Boolean hasRequest = false;
|
||||
Int32 schedule = 0;
|
||||
|
||||
// Update the state according to request priority
|
||||
if(this.StateChangeRequests[StateChangeRequest.Stop] || this.IsDisposing || this.IsDisposed) {
|
||||
hasRequest = true;
|
||||
this.WorkerState = WorkerState.Stopped;
|
||||
schedule = this.StateChangeRequests[StateChangeRequest.Stop] ? Timeout.Infinite : 0;
|
||||
} else if(this.StateChangeRequests[StateChangeRequest.Pause]) {
|
||||
hasRequest = true;
|
||||
this.WorkerState = WorkerState.Paused;
|
||||
schedule = Timeout.Infinite;
|
||||
} else if(this.StateChangeRequests[StateChangeRequest.Start] || this.StateChangeRequests[StateChangeRequest.Resume]) {
|
||||
hasRequest = true;
|
||||
this.WorkerState = WorkerState.Waiting;
|
||||
}
|
||||
|
||||
// Signals all state changes to continue
|
||||
// as a command has been handled.
|
||||
if(hasRequest) {
|
||||
this.ClearStateChangeRequests(schedule, currentState, this.WorkerState);
|
||||
}
|
||||
|
||||
return hasRequest;
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Signals all state change requests to set.
|
||||
/// </summary>
|
||||
/// <param name="schedule">The cycle schedule.</param>
|
||||
/// <param name="oldState">The previous worker state.</param>
|
||||
/// <param name="newState">The new worker state.</param>
|
||||
private void ClearStateChangeRequests(Int32 schedule, WorkerState oldState, WorkerState newState) {
|
||||
lock(this._syncLock) {
|
||||
// Mark all events as completed
|
||||
this.StateChangeRequests[StateChangeRequest.Start] = false;
|
||||
this.StateChangeRequests[StateChangeRequest.Pause] = false;
|
||||
this.StateChangeRequests[StateChangeRequest.Resume] = false;
|
||||
this.StateChangeRequests[StateChangeRequest.Stop] = false;
|
||||
|
||||
this.StateChangedEvent.Set();
|
||||
this.CycleCompletedEvent.Set();
|
||||
this.OnStateChangeProcessed(oldState, newState);
|
||||
this.ScheduleCycle(schedule);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+230
-237
@@ -1,240 +1,233 @@
|
||||
namespace Swan.Threading
|
||||
{
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
#nullable enable
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Swan.Threading {
|
||||
/// <summary>
|
||||
/// Provides base infrastructure for Timer and Thread workers.
|
||||
/// </summary>
|
||||
/// <seealso cref="IWorker" />
|
||||
public abstract class WorkerBase : IWorker, IDisposable {
|
||||
// Since these are API property backers, we use interlocked to read from them
|
||||
// to avoid deadlocked reads
|
||||
private readonly Object _syncLock = new Object();
|
||||
|
||||
private readonly AtomicBoolean _isDisposed = new AtomicBoolean();
|
||||
private readonly AtomicBoolean _isDisposing = new AtomicBoolean();
|
||||
private readonly AtomicEnum<WorkerState> _workerState = new AtomicEnum<WorkerState>(WorkerState.Created);
|
||||
private readonly AtomicTimeSpan _timeSpan;
|
||||
|
||||
/// <summary>
|
||||
/// Provides base infrastructure for Timer and Thread workers.
|
||||
/// Initializes a new instance of the <see cref="WorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <seealso cref="IWorker" />
|
||||
public abstract class WorkerBase : IWorker, IDisposable
|
||||
{
|
||||
// Since these are API property backers, we use interlocked to read from them
|
||||
// to avoid deadlocked reads
|
||||
private readonly object _syncLock = new object();
|
||||
|
||||
private readonly AtomicBoolean _isDisposed = new AtomicBoolean();
|
||||
private readonly AtomicBoolean _isDisposing = new AtomicBoolean();
|
||||
private readonly AtomicEnum<WorkerState> _workerState = new AtomicEnum<WorkerState>(WorkerState.Created);
|
||||
private readonly AtomicTimeSpan _timeSpan;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="WorkerBase"/> class.
|
||||
/// </summary>
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="period">The execution interval.</param>
|
||||
protected WorkerBase(string name, TimeSpan period)
|
||||
{
|
||||
Name = name;
|
||||
_timeSpan = new AtomicTimeSpan(period);
|
||||
|
||||
StateChangeRequests = new Dictionary<StateChangeRequest, bool>(5)
|
||||
{
|
||||
[StateChangeRequest.Start] = false,
|
||||
[StateChangeRequest.Pause] = false,
|
||||
[StateChangeRequest.Resume] = false,
|
||||
[StateChangeRequest.Stop] = false,
|
||||
};
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Enumerates all the different state change requests.
|
||||
/// </summary>
|
||||
protected enum StateChangeRequest
|
||||
{
|
||||
/// <summary>
|
||||
/// No state change request.
|
||||
/// </summary>
|
||||
None,
|
||||
|
||||
/// <summary>
|
||||
/// Start state change request
|
||||
/// </summary>
|
||||
Start,
|
||||
|
||||
/// <summary>
|
||||
/// Pause state change request
|
||||
/// </summary>
|
||||
Pause,
|
||||
|
||||
/// <summary>
|
||||
/// Resume state change request
|
||||
/// </summary>
|
||||
Resume,
|
||||
|
||||
/// <summary>
|
||||
/// Stop state change request
|
||||
/// </summary>
|
||||
Stop,
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public string Name { get; }
|
||||
|
||||
/// <inheritdoc />
|
||||
public TimeSpan Period
|
||||
{
|
||||
get => _timeSpan.Value;
|
||||
set => _timeSpan.Value = value;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public WorkerState WorkerState
|
||||
{
|
||||
get => _workerState.Value;
|
||||
protected set => _workerState.Value = value;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public bool IsDisposed
|
||||
{
|
||||
get => _isDisposed.Value;
|
||||
protected set => _isDisposed.Value = value;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public bool IsDisposing
|
||||
{
|
||||
get => _isDisposing.Value;
|
||||
protected set => _isDisposing.Value = value;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the default period of 15 milliseconds which is the default precision for timers.
|
||||
/// </summary>
|
||||
protected static TimeSpan DefaultPeriod { get; } = TimeSpan.FromMilliseconds(15);
|
||||
|
||||
/// <summary>
|
||||
/// Gets a value indicating whether stop has been requested.
|
||||
/// This is useful to prevent more requests from being issued.
|
||||
/// </summary>
|
||||
protected bool IsStopRequested => StateChangeRequests[StateChangeRequest.Stop];
|
||||
|
||||
/// <summary>
|
||||
/// Gets the cycle stopwatch.
|
||||
/// </summary>
|
||||
protected Stopwatch CycleStopwatch { get; } = new Stopwatch();
|
||||
|
||||
/// <summary>
|
||||
/// Gets the state change requests.
|
||||
/// </summary>
|
||||
protected Dictionary<StateChangeRequest, bool> StateChangeRequests { get; }
|
||||
|
||||
/// <summary>
|
||||
/// Gets the cycle completed event.
|
||||
/// </summary>
|
||||
protected ManualResetEventSlim CycleCompletedEvent { get; } = new ManualResetEventSlim(true);
|
||||
|
||||
/// <summary>
|
||||
/// Gets the state changed event.
|
||||
/// </summary>
|
||||
protected ManualResetEventSlim StateChangedEvent { get; } = new ManualResetEventSlim(true);
|
||||
|
||||
/// <summary>
|
||||
/// Gets the cycle logic cancellation owner.
|
||||
/// </summary>
|
||||
protected CancellationTokenOwner CycleCancellation { get; } = new CancellationTokenOwner();
|
||||
|
||||
/// <summary>
|
||||
/// Gets or sets the state change task.
|
||||
/// </summary>
|
||||
protected Task<WorkerState>? StateChangeTask { get; set; }
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> StartAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> PauseAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> ResumeAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> StopAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Dispose()
|
||||
{
|
||||
Dispose(true);
|
||||
GC.SuppressFinalize(this);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Releases unmanaged and - optionally - managed resources.
|
||||
/// </summary>
|
||||
/// <param name="disposing"><c>true</c> to release both managed and unmanaged resources; <c>false</c> to release only unmanaged resources.</param>
|
||||
protected virtual void Dispose(bool disposing)
|
||||
{
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (IsDisposed || IsDisposing) return;
|
||||
IsDisposing = true;
|
||||
}
|
||||
|
||||
// This also ensures the state change queue gets cleared
|
||||
StopAsync().Wait();
|
||||
StateChangedEvent.Set();
|
||||
CycleCompletedEvent.Set();
|
||||
|
||||
OnDisposing();
|
||||
|
||||
CycleStopwatch.Stop();
|
||||
StateChangedEvent.Dispose();
|
||||
CycleCompletedEvent.Dispose();
|
||||
CycleCancellation.Dispose();
|
||||
|
||||
IsDisposed = true;
|
||||
IsDisposing = false;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Handles the cycle logic exceptions.
|
||||
/// </summary>
|
||||
/// <param name="ex">The exception that was thrown.</param>
|
||||
protected abstract void OnCycleException(Exception ex);
|
||||
|
||||
/// <summary>
|
||||
/// Represents the user defined logic to be executed on a single worker cycle.
|
||||
/// Check the cancellation token continuously if you need responsive interrupts.
|
||||
/// </summary>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
protected abstract void ExecuteCycleLogic(CancellationToken cancellationToken);
|
||||
|
||||
/// <summary>
|
||||
/// This method is called automatically when <see cref="Dispose()"/> is called.
|
||||
/// Makes sure you release all resources within this call.
|
||||
/// </summary>
|
||||
protected abstract void OnDisposing();
|
||||
|
||||
/// <summary>
|
||||
/// Called when a state change request is processed.
|
||||
/// </summary>
|
||||
/// <param name="previousState">The state before the change.</param>
|
||||
/// <param name="newState">The new state.</param>
|
||||
protected virtual void OnStateChangeProcessed(WorkerState previousState, WorkerState newState)
|
||||
{
|
||||
// placeholder
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Computes the cycle delay.
|
||||
/// </summary>
|
||||
/// <param name="initialWorkerState">Initial state of the worker.</param>
|
||||
/// <returns>The number of milliseconds to delay for.</returns>
|
||||
protected int ComputeCycleDelay(WorkerState initialWorkerState)
|
||||
{
|
||||
var elapsedMillis = CycleStopwatch.ElapsedMilliseconds;
|
||||
var period = Period;
|
||||
var periodMillis = period.TotalMilliseconds;
|
||||
var delayMillis = periodMillis - elapsedMillis;
|
||||
|
||||
if (initialWorkerState == WorkerState.Paused || period == TimeSpan.MaxValue || delayMillis >= int.MaxValue)
|
||||
return Timeout.Infinite;
|
||||
|
||||
return elapsedMillis >= periodMillis ? 0 : Convert.ToInt32(Math.Floor(delayMillis));
|
||||
}
|
||||
}
|
||||
/// <param name="name">The name.</param>
|
||||
/// <param name="period">The execution interval.</param>
|
||||
protected WorkerBase(String name, TimeSpan period) {
|
||||
this.Name = name;
|
||||
this._timeSpan = new AtomicTimeSpan(period);
|
||||
|
||||
this.StateChangeRequests = new Dictionary<StateChangeRequest, Boolean>(5) {
|
||||
[StateChangeRequest.Start] = false,
|
||||
[StateChangeRequest.Pause] = false,
|
||||
[StateChangeRequest.Resume] = false,
|
||||
[StateChangeRequest.Stop] = false,
|
||||
};
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Enumerates all the different state change requests.
|
||||
/// </summary>
|
||||
protected enum StateChangeRequest {
|
||||
/// <summary>
|
||||
/// No state change request.
|
||||
/// </summary>
|
||||
None,
|
||||
|
||||
/// <summary>
|
||||
/// Start state change request
|
||||
/// </summary>
|
||||
Start,
|
||||
|
||||
/// <summary>
|
||||
/// Pause state change request
|
||||
/// </summary>
|
||||
Pause,
|
||||
|
||||
/// <summary>
|
||||
/// Resume state change request
|
||||
/// </summary>
|
||||
Resume,
|
||||
|
||||
/// <summary>
|
||||
/// Stop state change request
|
||||
/// </summary>
|
||||
Stop,
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public String Name {
|
||||
get;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public TimeSpan Period {
|
||||
get => this._timeSpan.Value;
|
||||
set => this._timeSpan.Value = value;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public WorkerState WorkerState {
|
||||
get => this._workerState.Value;
|
||||
protected set => this._workerState.Value = value;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Boolean IsDisposed {
|
||||
get => this._isDisposed.Value;
|
||||
protected set => this._isDisposed.Value = value;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Boolean IsDisposing {
|
||||
get => this._isDisposing.Value;
|
||||
protected set => this._isDisposing.Value = value;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the default period of 15 milliseconds which is the default precision for timers.
|
||||
/// </summary>
|
||||
protected static TimeSpan DefaultPeriod { get; } = TimeSpan.FromMilliseconds(15);
|
||||
|
||||
/// <summary>
|
||||
/// Gets a value indicating whether stop has been requested.
|
||||
/// This is useful to prevent more requests from being issued.
|
||||
/// </summary>
|
||||
protected Boolean IsStopRequested => this.StateChangeRequests[StateChangeRequest.Stop];
|
||||
|
||||
/// <summary>
|
||||
/// Gets the cycle stopwatch.
|
||||
/// </summary>
|
||||
protected Stopwatch CycleStopwatch { get; } = new Stopwatch();
|
||||
|
||||
/// <summary>
|
||||
/// Gets the state change requests.
|
||||
/// </summary>
|
||||
protected Dictionary<StateChangeRequest, Boolean> StateChangeRequests {
|
||||
get;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the cycle completed event.
|
||||
/// </summary>
|
||||
protected ManualResetEventSlim CycleCompletedEvent { get; } = new ManualResetEventSlim(true);
|
||||
|
||||
/// <summary>
|
||||
/// Gets the state changed event.
|
||||
/// </summary>
|
||||
protected ManualResetEventSlim StateChangedEvent { get; } = new ManualResetEventSlim(true);
|
||||
|
||||
/// <summary>
|
||||
/// Gets the cycle logic cancellation owner.
|
||||
/// </summary>
|
||||
protected CancellationTokenOwner CycleCancellation { get; } = new CancellationTokenOwner();
|
||||
|
||||
/// <summary>
|
||||
/// Gets or sets the state change task.
|
||||
/// </summary>
|
||||
protected Task<WorkerState>? StateChangeTask {
|
||||
get; set;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> StartAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> PauseAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> ResumeAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public abstract Task<WorkerState> StopAsync();
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Dispose() {
|
||||
this.Dispose(true);
|
||||
GC.SuppressFinalize(this);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Releases unmanaged and - optionally - managed resources.
|
||||
/// </summary>
|
||||
/// <param name="disposing"><c>true</c> to release both managed and unmanaged resources; <c>false</c> to release only unmanaged resources.</param>
|
||||
protected virtual void Dispose(Boolean disposing) {
|
||||
lock(this._syncLock) {
|
||||
if(this.IsDisposed || this.IsDisposing) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.IsDisposing = true;
|
||||
}
|
||||
|
||||
// This also ensures the state change queue gets cleared
|
||||
this.StopAsync().Wait();
|
||||
this.StateChangedEvent.Set();
|
||||
this.CycleCompletedEvent.Set();
|
||||
|
||||
this.OnDisposing();
|
||||
|
||||
this.CycleStopwatch.Stop();
|
||||
this.StateChangedEvent.Dispose();
|
||||
this.CycleCompletedEvent.Dispose();
|
||||
this.CycleCancellation.Dispose();
|
||||
|
||||
this.IsDisposed = true;
|
||||
this.IsDisposing = false;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Handles the cycle logic exceptions.
|
||||
/// </summary>
|
||||
/// <param name="ex">The exception that was thrown.</param>
|
||||
protected abstract void OnCycleException(Exception ex);
|
||||
|
||||
/// <summary>
|
||||
/// Represents the user defined logic to be executed on a single worker cycle.
|
||||
/// Check the cancellation token continuously if you need responsive interrupts.
|
||||
/// </summary>
|
||||
/// <param name="cancellationToken">The cancellation token.</param>
|
||||
protected abstract void ExecuteCycleLogic(CancellationToken cancellationToken);
|
||||
|
||||
/// <summary>
|
||||
/// This method is called automatically when <see cref="Dispose()"/> is called.
|
||||
/// Makes sure you release all resources within this call.
|
||||
/// </summary>
|
||||
protected abstract void OnDisposing();
|
||||
|
||||
/// <summary>
|
||||
/// Called when a state change request is processed.
|
||||
/// </summary>
|
||||
/// <param name="previousState">The state before the change.</param>
|
||||
/// <param name="newState">The new state.</param>
|
||||
protected virtual void OnStateChangeProcessed(WorkerState previousState, WorkerState newState) {
|
||||
// placeholder
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Computes the cycle delay.
|
||||
/// </summary>
|
||||
/// <param name="initialWorkerState">Initial state of the worker.</param>
|
||||
/// <returns>The number of milliseconds to delay for.</returns>
|
||||
protected Int32 ComputeCycleDelay(WorkerState initialWorkerState) {
|
||||
Int64 elapsedMillis = this.CycleStopwatch.ElapsedMilliseconds;
|
||||
TimeSpan period = this.Period;
|
||||
Double periodMillis = period.TotalMilliseconds;
|
||||
Double delayMillis = periodMillis - elapsedMillis;
|
||||
|
||||
return initialWorkerState == WorkerState.Paused || period == TimeSpan.MaxValue || delayMillis >= Int32.MaxValue ? Timeout.Infinite : elapsedMillis >= periodMillis ? 0 : Convert.ToInt32(Math.Floor(delayMillis));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,151 +1,146 @@
|
||||
namespace Swan.Threading
|
||||
{
|
||||
using System;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
using System;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Swan.Threading {
|
||||
/// <summary>
|
||||
/// Represents a class that implements delay logic for thread workers.
|
||||
/// </summary>
|
||||
public static class WorkerDelayProvider {
|
||||
/// <summary>
|
||||
/// Represents a class that implements delay logic for thread workers.
|
||||
/// Gets the default delay provider.
|
||||
/// </summary>
|
||||
public static class WorkerDelayProvider
|
||||
{
|
||||
/// <summary>
|
||||
/// Gets the default delay provider.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider Default => TokenTimeout;
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which simply waits on the task and cancels on
|
||||
/// the cancellation token.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider Token => new TokenCancellableDelay();
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which waits on the task and cancels on both,
|
||||
/// the cancellation token and a wanted delay timeout.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider TokenTimeout => new TokenTimeoutCancellableDelay();
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which uses short sleep intervals of 5ms.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider TokenSleep => new TokenSleepDelay();
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which uses short delay intervals of 5ms and
|
||||
/// a wait on the delay task in the final loop.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider SteppedToken => new SteppedTokenDelay();
|
||||
|
||||
private class TokenCancellableDelay : IWorkerDelayProvider
|
||||
{
|
||||
public void ExecuteCycleDelay(int wantedDelay, Task delayTask, CancellationToken token)
|
||||
{
|
||||
if (wantedDelay == 0 || wantedDelay < -1)
|
||||
return;
|
||||
|
||||
// for wanted delays of less than 30ms it is not worth
|
||||
// passing a timeout or a token as it only adds unnecessary
|
||||
// overhead.
|
||||
if (wantedDelay <= 30)
|
||||
{
|
||||
try { delayTask.Wait(token); }
|
||||
catch { /* ignore */ }
|
||||
return;
|
||||
}
|
||||
|
||||
// only wait on the cancellation token
|
||||
// or until the task completes normally
|
||||
try { delayTask.Wait(token); }
|
||||
catch { /* ignore */ }
|
||||
}
|
||||
}
|
||||
|
||||
private class TokenTimeoutCancellableDelay : IWorkerDelayProvider
|
||||
{
|
||||
public void ExecuteCycleDelay(int wantedDelay, Task delayTask, CancellationToken token)
|
||||
{
|
||||
if (wantedDelay == 0 || wantedDelay < -1)
|
||||
return;
|
||||
|
||||
// for wanted delays of less than 30ms it is not worth
|
||||
// passing a timeout or a token as it only adds unnecessary
|
||||
// overhead.
|
||||
if (wantedDelay <= 30)
|
||||
{
|
||||
try { delayTask.Wait(token); }
|
||||
catch { /* ignore */ }
|
||||
return;
|
||||
}
|
||||
|
||||
try { delayTask.Wait(wantedDelay, token); }
|
||||
catch { /* ignore */ }
|
||||
}
|
||||
}
|
||||
|
||||
private class TokenSleepDelay : IWorkerDelayProvider
|
||||
{
|
||||
private readonly Stopwatch _elapsedWait = new Stopwatch();
|
||||
|
||||
public void ExecuteCycleDelay(int wantedDelay, Task delayTask, CancellationToken token)
|
||||
{
|
||||
_elapsedWait.Restart();
|
||||
|
||||
if (wantedDelay == 0 || wantedDelay < -1)
|
||||
return;
|
||||
|
||||
while (!token.IsCancellationRequested)
|
||||
{
|
||||
Thread.Sleep(5);
|
||||
|
||||
if (wantedDelay != Timeout.Infinite && _elapsedWait.ElapsedMilliseconds >= wantedDelay)
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private class SteppedTokenDelay : IWorkerDelayProvider
|
||||
{
|
||||
private const int StepMilliseconds = 15;
|
||||
private readonly Stopwatch _elapsedWait = new Stopwatch();
|
||||
|
||||
public void ExecuteCycleDelay(int wantedDelay, Task delayTask, CancellationToken token)
|
||||
{
|
||||
_elapsedWait.Restart();
|
||||
|
||||
if (wantedDelay == 0 || wantedDelay < -1)
|
||||
return;
|
||||
|
||||
if (wantedDelay == Timeout.Infinite)
|
||||
{
|
||||
try { delayTask.Wait(wantedDelay, token); }
|
||||
catch { /* Ignore cancelled tasks */ }
|
||||
return;
|
||||
}
|
||||
|
||||
while (!token.IsCancellationRequested)
|
||||
{
|
||||
var remainingWaitTime = wantedDelay - Convert.ToInt32(_elapsedWait.ElapsedMilliseconds);
|
||||
|
||||
// Exit for no remaining wait time
|
||||
if (remainingWaitTime <= 0)
|
||||
break;
|
||||
|
||||
if (remainingWaitTime >= StepMilliseconds)
|
||||
{
|
||||
Task.Delay(StepMilliseconds, token).Wait(token);
|
||||
}
|
||||
else
|
||||
{
|
||||
try { delayTask.Wait(remainingWaitTime); }
|
||||
catch { /* ignore cancellation of task exception */ }
|
||||
}
|
||||
|
||||
if (_elapsedWait.ElapsedMilliseconds >= wantedDelay)
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
public static IWorkerDelayProvider Default => TokenTimeout;
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which simply waits on the task and cancels on
|
||||
/// the cancellation token.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider Token => new TokenCancellableDelay();
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which waits on the task and cancels on both,
|
||||
/// the cancellation token and a wanted delay timeout.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider TokenTimeout => new TokenTimeoutCancellableDelay();
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which uses short sleep intervals of 5ms.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider TokenSleep => new TokenSleepDelay();
|
||||
|
||||
/// <summary>
|
||||
/// Provides a delay implementation which uses short delay intervals of 5ms and
|
||||
/// a wait on the delay task in the final loop.
|
||||
/// </summary>
|
||||
public static IWorkerDelayProvider SteppedToken => new SteppedTokenDelay();
|
||||
|
||||
private class TokenCancellableDelay : IWorkerDelayProvider {
|
||||
public void ExecuteCycleDelay(Int32 wantedDelay, Task delayTask, CancellationToken token) {
|
||||
if(wantedDelay == 0 || wantedDelay < -1) {
|
||||
return;
|
||||
}
|
||||
|
||||
// for wanted delays of less than 30ms it is not worth
|
||||
// passing a timeout or a token as it only adds unnecessary
|
||||
// overhead.
|
||||
if(wantedDelay <= 30) {
|
||||
try {
|
||||
delayTask.Wait(token);
|
||||
} catch { /* ignore */ }
|
||||
return;
|
||||
}
|
||||
|
||||
// only wait on the cancellation token
|
||||
// or until the task completes normally
|
||||
try {
|
||||
delayTask.Wait(token);
|
||||
} catch { /* ignore */ }
|
||||
}
|
||||
}
|
||||
|
||||
private class TokenTimeoutCancellableDelay : IWorkerDelayProvider {
|
||||
public void ExecuteCycleDelay(Int32 wantedDelay, Task delayTask, CancellationToken token) {
|
||||
if(wantedDelay == 0 || wantedDelay < -1) {
|
||||
return;
|
||||
}
|
||||
|
||||
// for wanted delays of less than 30ms it is not worth
|
||||
// passing a timeout or a token as it only adds unnecessary
|
||||
// overhead.
|
||||
if(wantedDelay <= 30) {
|
||||
try {
|
||||
delayTask.Wait(token);
|
||||
} catch { /* ignore */ }
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
_ = delayTask.Wait(wantedDelay, token);
|
||||
} catch { /* ignore */ }
|
||||
}
|
||||
}
|
||||
|
||||
private class TokenSleepDelay : IWorkerDelayProvider {
|
||||
private readonly Stopwatch _elapsedWait = new Stopwatch();
|
||||
|
||||
public void ExecuteCycleDelay(Int32 wantedDelay, Task delayTask, CancellationToken token) {
|
||||
this._elapsedWait.Restart();
|
||||
|
||||
if(wantedDelay == 0 || wantedDelay < -1) {
|
||||
return;
|
||||
}
|
||||
|
||||
while(!token.IsCancellationRequested) {
|
||||
Thread.Sleep(5);
|
||||
|
||||
if(wantedDelay != Timeout.Infinite && this._elapsedWait.ElapsedMilliseconds >= wantedDelay) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private class SteppedTokenDelay : IWorkerDelayProvider {
|
||||
private const Int32 StepMilliseconds = 15;
|
||||
private readonly Stopwatch _elapsedWait = new Stopwatch();
|
||||
|
||||
public void ExecuteCycleDelay(Int32 wantedDelay, Task delayTask, CancellationToken token) {
|
||||
this._elapsedWait.Restart();
|
||||
|
||||
if(wantedDelay == 0 || wantedDelay < -1) {
|
||||
return;
|
||||
}
|
||||
|
||||
if(wantedDelay == Timeout.Infinite) {
|
||||
try {
|
||||
_ = delayTask.Wait(wantedDelay, token);
|
||||
} catch { /* Ignore cancelled tasks */ }
|
||||
return;
|
||||
}
|
||||
|
||||
while(!token.IsCancellationRequested) {
|
||||
Int32 remainingWaitTime = wantedDelay - Convert.ToInt32(this._elapsedWait.ElapsedMilliseconds);
|
||||
|
||||
// Exit for no remaining wait time
|
||||
if(remainingWaitTime <= 0) {
|
||||
break;
|
||||
}
|
||||
|
||||
if(remainingWaitTime >= StepMilliseconds) {
|
||||
Task.Delay(StepMilliseconds, token).Wait(token);
|
||||
} else {
|
||||
try {
|
||||
_ = delayTask.Wait(remainingWaitTime);
|
||||
} catch { /* ignore cancellation of task exception */ }
|
||||
}
|
||||
|
||||
if(this._elapsedWait.ElapsedMilliseconds >= wantedDelay) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user