Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
164 changes: 164 additions & 0 deletions IntervalAction.Test/IntervalActionTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@
{
Interlocked.Increment(ref executions);
// Simulate a long running task.
Thread.Sleep(500);

Check warning on line 49 in IntervalAction.Test/IntervalActionTests.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Do not use 'Thread.Sleep()' in a test.

Check warning on line 49 in IntervalAction.Test/IntervalActionTests.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Do not use 'Thread.Sleep()' in a test.

Check warning on line 49 in IntervalAction.Test/IntervalActionTests.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Do not use 'Thread.Sleep()' in a test.

Check warning on line 49 in IntervalAction.Test/IntervalActionTests.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Do not use 'Thread.Sleep()' in a test.
},
IntervalType = IntervalType.FromLastStart
};
Expand Down Expand Up @@ -322,6 +322,170 @@
StringAssert.Contains(exception.StackTrace, nameof(ThrowFromNamedMethod));
}

[TestMethod]
public async Task AsyncActionRunsNeverOverlap()
{
// Arrange: each run awaits far longer than the interval between runs
int running = 0;
int maxRunning = 0;
int started = 0;
IntervalActionOptions options = new()
{
PollingInterval = TimeSpan.FromMilliseconds(10),
ActionInterval = TimeSpan.FromMilliseconds(20),
AsyncAction = async _ =>
{
int nowRunning = Interlocked.Increment(ref running);
Interlocked.Increment(ref started);
InterlockedMax(ref maxRunning, nowRunning);
await Task.Delay(200, CancellationToken.None).ConfigureAwait(false);
Interlocked.Decrement(ref running);
},
IntervalType = IntervalType.FromLastCompletion
};

// Act
IntervalAction intervalAction = IntervalAction.Start(options);
await Task.Delay(1000).ConfigureAwait(false);

Check warning on line 349 in IntervalAction.Test/IntervalActionTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Consider using the overload that accepts a CancellationToken and pass 'TestContext.CancellationToken'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_IntervalAction&issues=AaEL4nquoRaNIpvVE-WF&open=AaEL4nquoRaNIpvVE-WF&pullRequest=78
intervalAction.Stop();

// Assert: a run lasts until its task completes, so the next one waits for it
Assert.IsGreaterThanOrEqualTo(2, Volatile.Read(ref started), "Expected the async action to run more than once.");
Assert.AreEqual(1, Volatile.Read(ref maxRunning), "Async action runs should never overlap.");

intervalAction.RethrowExceptions();
}

[TestMethod]
public async Task AsyncActionFromLastCompletionMeasuresFromWhenTheWorkFinished()
{
// Arrange: record when each run starts and when its awaited work finishes
TimeSpan actionInterval = TimeSpan.FromMilliseconds(100);
List<(long Start, long Finish)> runs = [];
Lock runsLock = new();
IntervalActionOptions options = new()
{
PollingInterval = TimeSpan.FromMilliseconds(10),
ActionInterval = actionInterval,
AsyncAction = async _ =>
{
long start = System.Diagnostics.Stopwatch.GetTimestamp();
await Task.Delay(150, CancellationToken.None).ConfigureAwait(false);
long finish = System.Diagnostics.Stopwatch.GetTimestamp();
lock (runsLock)
{
runs.Add((start, finish));
}
},
IntervalType = IntervalType.FromLastCompletion
};

// Act
IntervalAction intervalAction = IntervalAction.Start(options);
await Task.Delay(1200).ConfigureAwait(false);

Check warning on line 385 in IntervalAction.Test/IntervalActionTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Consider using the overload that accepts a CancellationToken and pass 'TestContext.CancellationToken'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_IntervalAction&issues=AaEL4nquoRaNIpvVE-WI&open=AaEL4nquoRaNIpvVE-WI&pullRequest=78
intervalAction.Stop();
intervalAction.RethrowExceptions();

// Assert: each run starts at least an interval after the previous one's work finished
(long Start, long Finish)[] snapshot;
lock (runsLock)
{
snapshot = [.. runs];
}

Assert.IsGreaterThanOrEqualTo(2, snapshot.Length, "Expected at least two completed runs.");
for (int i = 1; i < snapshot.Length; i++)
{
TimeSpan gap = TimeSpan.FromTicks((long)((snapshot[i].Start - snapshot[i - 1].Finish) * ((double)TimeSpan.TicksPerSecond / System.Diagnostics.Stopwatch.Frequency)));
Assert.IsGreaterThanOrEqualTo(actionInterval, gap, $"Run {i} started {gap.TotalMilliseconds:F0} ms after the previous run finished.");
}
}

[TestMethod]
public async Task AsyncActionExceptionAfterAwaitIsRethrown()
{
// Arrange
string exceptionMessage = "Thrown after an await";
IntervalActionOptions options = new()
{
PollingInterval = TimeSpan.FromMilliseconds(10),
ActionInterval = TimeSpan.Zero,
AsyncAction = async _ =>
{
await Task.Delay(20, CancellationToken.None).ConfigureAwait(false);
throw new InvalidOperationException(exceptionMessage);
},
IntervalType = IntervalType.FromLastStart
};

// Act
IntervalAction intervalAction = IntervalAction.Start(options);
await WaitForPollingToFaultAsync(intervalAction).ConfigureAwait(false);

// Assert: the exception reaches RethrowExceptions instead of the thread pool
InvalidOperationException exception = Assert.ThrowsExactly<InvalidOperationException>(intervalAction.RethrowExceptions);
Assert.AreEqual(exceptionMessage, exception.Message);
intervalAction.Stop();
}

[TestMethod]
public async Task StopCancelsTheAsyncActionToken()
{
// Arrange: the action waits on its token for far longer than the test
using ManualResetEventSlim started = new();
IntervalActionOptions options = new()
{
PollingInterval = TimeSpan.FromMilliseconds(10),
ActionInterval = TimeSpan.FromHours(1),
AsyncAction = async cancellationToken =>
{
started.Set();
await Task.Delay(TimeSpan.FromHours(1), cancellationToken).ConfigureAwait(false);
},
IntervalType = IntervalType.FromLastCompletion
};

IntervalAction intervalAction = IntervalAction.Start(options);
Assert.IsTrue(started.Wait(TimeSpan.FromSeconds(10)), "The action should start.");

Check warning on line 449 in IntervalAction.Test/IntervalActionTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Consider using the overload that accepts a CancellationToken and pass 'TestContext.CancellationToken'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_IntervalAction&issues=AaEL4nquoRaNIpvVE-WG&open=AaEL4nquoRaNIpvVE-WG&pullRequest=78
Task? actionTask = intervalAction.ActionTask;
Assert.IsNotNull(actionTask, "The action should still be running.");

// Act
intervalAction.Stop();

// Assert: the run ends promptly, and giving up on cancellation is not reported as a failure
Assert.AreSame(actionTask, await Task.WhenAny(actionTask, Task.Delay(TimeSpan.FromSeconds(10))).ConfigureAwait(false), "Stop should cancel the running async action.");

Check warning on line 457 in IntervalAction.Test/IntervalActionTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Consider using the overload that accepts a CancellationToken and pass 'TestContext.CancellationToken'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_IntervalAction&issues=AaEL4nquoRaNIpvVE-WH&open=AaEL4nquoRaNIpvVE-WH&pullRequest=78
Assert.IsFalse(actionTask.IsFaulted, "A run that ends because Stop cancelled it should not fault.");
intervalAction.RethrowExceptions();
}

[TestMethod]
public void StartThrowsWhenBothActionAndAsyncActionAreSet()
{
IntervalActionOptions options = new()
{
Action = () => { },
AsyncAction = _ => Task.CompletedTask,
};

_ = Assert.ThrowsExactly<ArgumentException>(() => IntervalAction.Start(options));
}

private static void InterlockedMax(ref int target, int value)
{
int current = Volatile.Read(ref target);
while (value > current)
{
int seen = Interlocked.CompareExchange(ref target, value, current);
if (seen == current)
{
return;
}

current = seen;
}
}

[TestMethod]
public async Task RethrowExceptionsReportsAnActionThatThrowsAfterStop()
{
Expand Down
112 changes: 105 additions & 7 deletions IntervalAction/IntervalAction.cs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,22 @@ public class IntervalAction
/// </summary>
private Action Action { get; init; } = null!;

/// <summary>
/// Gets the asynchronous action to be executed at each interval in place of <see cref="Action"/>,
/// or <see langword="null"/> if a synchronous action was given.
/// </summary>
private Func<CancellationToken, Task>? AsyncAction { get; init; }

/// <summary>
/// Gets or sets the cancellation source handed to the running <see cref="AsyncAction"/>, which
/// <see cref="Stop"/> cancels, or <see langword="null"/> when no asynchronous run is in flight.
/// </summary>
/// <remarks>
/// Only read, cancelled or replaced while holding <see cref="Lock"/>. A run takes it out of this
/// property before disposing it, so a cancel never meets a disposed source.
/// </remarks>
private CancellationTokenSource? ActionCancellation { get; set; }

/// <summary>
/// Gets the interval at which the action should be executed.
/// </summary>
Expand Down Expand Up @@ -119,15 +135,31 @@ private IntervalAction() { }
/// </summary>
/// <param name="intervalActionOptions">The options for configuring the interval action.</param>
/// <returns>A new instance of <see cref="IntervalAction"/>.</returns>
/// <exception cref="ArgumentNullException">Thrown if <paramref name="intervalActionOptions"/> or its <see cref="IntervalActionOptions.Action"/> is null.</exception>
/// <exception cref="ArgumentNullException">
/// Thrown if <paramref name="intervalActionOptions"/> is null, or if its <see cref="IntervalActionOptions.Action"/>
/// is null and no <see cref="IntervalActionOptions.AsyncAction"/> is given.
/// </exception>
/// <exception cref="ArgumentException">
/// Thrown if both <see cref="IntervalActionOptions.Action"/> and <see cref="IntervalActionOptions.AsyncAction"/> are set.
/// </exception>
/// <exception cref="ArgumentOutOfRangeException">
/// Thrown if <see cref="IntervalActionOptions.PollingInterval"/> is zero or negative, which includes
/// <see cref="Timeout.InfiniteTimeSpan"/>.
/// </exception>
public static IntervalAction Start(IntervalActionOptions intervalActionOptions)
{
Ensure.NotNull(intervalActionOptions);
Ensure.NotNull(intervalActionOptions.Action);

if (intervalActionOptions.AsyncAction is null)
{
Ensure.NotNull(intervalActionOptions.Action);
}
else if (intervalActionOptions.Action is not null && !ReferenceEquals(intervalActionOptions.Action, IntervalActionOptions.NoAction))
{
throw new ArgumentException(
$"Set either {nameof(IntervalActionOptions.Action)} or {nameof(IntervalActionOptions.AsyncAction)}, not both.",
nameof(intervalActionOptions));
}

// Rejected here rather than left to Task.Delay in the polling loop: a negative interval would
// fault the loop after the first run, an infinite one would leave Restart() and Stop() waiting
Expand All @@ -143,7 +175,8 @@ public static IntervalAction Start(IntervalActionOptions intervalActionOptions)
IntervalAction intervalAction = new()
{
PollingInterval = intervalActionOptions.PollingInterval,
Action = intervalActionOptions.Action,
Action = intervalActionOptions.Action ?? IntervalActionOptions.NoAction,
AsyncAction = intervalActionOptions.AsyncAction,
ActionInterval = intervalActionOptions.ActionInterval,
IntervalType = intervalActionOptions.IntervalType
};
Expand All @@ -162,17 +195,30 @@ public static IntervalAction Start(IntervalActionOptions intervalActionOptions)
/// </remarks>
public void Stop()
{
CancellationTokenSource? actionCancellation;

lock (Lock)
{
StopGeneration++;
ShouldPoll = false;
actionCancellation = ActionCancellation;
}

// Cancelled outside the lock, because cancelling runs the action's own token callbacks
try
{
actionCancellation?.Cancel();
}
catch (ObjectDisposedException)
{
// The run finished and disposed its source after the lock was released: nothing to cancel
}
}

/// <summary>
/// Restarts the polling of the action.
/// </summary>
public void Restart() => RestartAsync().Wait();
public void Restart() => RestartAsync().Wait(CancellationToken.None);

/// <summary>
/// Asynchronously restarts the polling of the action.
Expand Down Expand Up @@ -242,14 +288,14 @@ private async Task RestartCoreAsync(Task previousRestart, long stopGeneration)
while (shouldPoll)
{
TryRun();
await Task.Delay(PollingInterval).ConfigureAwait(false);
await Task.Delay(PollingInterval, CancellationToken.None).ConfigureAwait(false);

lock (Lock)
{
shouldPoll = ShouldPoll;
}
}
});
}, CancellationToken.None);
}

/// <summary>
Expand Down Expand Up @@ -277,6 +323,14 @@ internal bool TryRun()

if (ActionInterval >= TimeSpan.Zero && ActionTask is null && HasIntervalElapsed())
{
if (AsyncAction is { } asyncAction)
{
CancellationTokenSource cancellation = new();
ActionCancellation = cancellation;
ActionTask = Task.Run(() => RunAsyncAction(asyncAction, cancellation), CancellationToken.None);
return true;
}

ActionTask = Task.Run(() =>
{
if (IntervalType == IntervalType.FromLastStart)
Expand All @@ -290,7 +344,7 @@ internal bool TryRun()
{
RecordRun();
}
});
}, CancellationToken.None);

return true;
}
Expand All @@ -299,6 +353,50 @@ internal bool TryRun()
}
}

/// <summary>
/// Runs <see cref="AsyncAction"/> once, so that the returned task completes when its work does.
/// </summary>
/// <param name="asyncAction">The action to run.</param>
/// <param name="cancellation">The source whose token the action receives, and which this run disposes.</param>
/// <returns>A task that completes, or faults, when the action does.</returns>
private async Task RunAsyncAction(Func<CancellationToken, Task> asyncAction, CancellationTokenSource cancellation)
{
try
{
if (IntervalType == IntervalType.FromLastStart)
{
RecordRun();
}

try
{
await asyncAction(cancellation.Token).ConfigureAwait(false);
}
catch (OperationCanceledException) when (cancellation.IsCancellationRequested)
{
// Stop() asked the action to give up, and it did: that is how a run ends after a stop,
// not a failure to report
}

if (IntervalType == IntervalType.FromLastCompletion)
{
RecordRun();
}
}
finally
{
lock (Lock)
{
if (ActionCancellation == cancellation)
{
ActionCancellation = null;
}
}

cancellation.Dispose();
}
}

/// <summary>
/// Reports whether <see cref="ActionInterval"/> has passed since the last run, measured on a
/// monotonic clock. Callers hold <see cref="Lock"/>.
Expand Down
Loading
Loading