This commit is contained in:
2019-12-04 18:57:18 +01:00
parent 4692422c9a
commit 6263791dff
225 changed files with 33065 additions and 2 deletions
+141
View File
@@ -0,0 +1,141 @@
namespace Swan.Threading
{
using System;
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;
/// <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 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
}
}
+292
View File
@@ -0,0 +1,292 @@
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>
/// 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
? 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(
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();
}
}
}
}
+328
View File
@@ -0,0 +1,328 @@
namespace Swan.Threading
{
using System;
using System.Threading;
using System.Threading.Tasks;
/// <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 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
? 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);
}
}
}
}
+240
View File
@@ -0,0 +1,240 @@
namespace Swan.Threading
{
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;
/// <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>
/// 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));
}
}
}
+151
View File
@@ -0,0 +1,151 @@
namespace Swan.Threading
{
using System;
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;
/// <summary>
/// Represents a class that implements delay logic for thread workers.
/// </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;
}
}
}
}
}