diff --git a/src/CanKit.Pro.Actor/IProtocolActor.cs b/src/CanKit.Pro.Actor/IProtocolActor.cs index 31a30a9..c3e54d5 100644 --- a/src/CanKit.Pro.Actor/IProtocolActor.cs +++ b/src/CanKit.Pro.Actor/IProtocolActor.cs @@ -1,4 +1,5 @@ using System; +using System.Threading; using System.Threading.Tasks; namespace CanKit.Pro.Actor @@ -30,13 +31,26 @@ public interface IProtocolActor : IDisposable /// is surfaced through the returned task's fault, not through /// — the caller is already positioned to observe /// it by awaiting. + /// + /// withdraws for as long as it + /// has not started: a token cancelled before the call, or while the item waits in the + /// mailbox, cancels the returned task at once and the work never runs. Work that has + /// already started is never interrupted — it runs to completion, and cancelling then has no + /// effect on the task's outcome, so the mailbox's single-writer discipline (FR-RAW-021) is + /// not weakened by a half-finished work item. + /// /// - Task PostAsync(Action work); + /// The work to run on the mailbox loop. + /// Withdraws the work while it is still queued. + Task PostAsync(Action work, CancellationToken cancellationToken = default); /// - /// Same as but returns 's result. + /// Same as but returns + /// 's result. /// - Task PostAsync(Func work); + /// The work to run on the mailbox loop. + /// Withdraws the work while it is still queued. + Task PostAsync(Func work, CancellationToken cancellationToken = default); /// /// Schedules to run on the actor's mailbox loop once @@ -52,8 +66,8 @@ public interface IProtocolActor : IDisposable /// Raised whenever a posted work item (via ) or a scheduled callback /// (via ) throws — the actor's single, defined channel for /// background exceptions (FR-RAW-023). The mailbox loop keeps running afterward; one - /// failing item never stops the actor. Never raised for / - /// failures, which surface through their own returned + /// failing item never stops the actor. Never raised for / + /// failures, which surface through their own returned /// task instead. /// event EventHandler BackgroundExceptionOccurred; diff --git a/src/CanKit.Pro.Actor/ProtocolActor.cs b/src/CanKit.Pro.Actor/ProtocolActor.cs index 2c7418c..26e2466 100644 --- a/src/CanKit.Pro.Actor/ProtocolActor.cs +++ b/src/CanKit.Pro.Actor/ProtocolActor.cs @@ -187,7 +187,7 @@ public sealed class ProtocolActor : IProtocolActor /// /// /// Lets a public sync API safely detect that it is already on the actor loop and run the - /// requested work inline instead of routing it through and + /// requested work inline instead of routing it through and /// synchronously waiting on the returned task — which would deadlock the loop against /// itself. External callers still take the marshal-through-mailbox path exactly as /// before. @@ -302,50 +302,130 @@ public void Post(Action work) } /// - public Task PostAsync(Action work) + public Task PostAsync(Action work, CancellationToken cancellationToken = default) { if (work is null) throw new ArgumentNullException(nameof(work)); - var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - PostInternal( - () => + if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken); + var call = new WithdrawableCall(cancellationToken); + return PostWithdrawable(call, () => + { + // Disposed on every path out, including the one where the token withdrew the item. + using (call) { + if (!call.TryStart()) return; try { work(); - tcs.TrySetResult(null); + call.Completion.TrySetResult(null); } catch (Exception ex) { - tcs.TrySetException(ex); + call.Completion.TrySetException(ex); } - }, - // If the marshal itself fails (SynchronizationContext mode, Send throws before - // ever invoking the wrapped work above), the wrapper's own try/catch never runs, - // so nothing would otherwise complete this task -- it would hang forever even - // though PostAsync failures are documented to surface via the returned task. - onDispatchFailure: ex => tcs.TrySetException(ex)); - return tcs.Task; + } + }); } /// - public Task PostAsync(Func work) + public Task PostAsync(Func work, CancellationToken cancellationToken = default) { if (work is null) throw new ArgumentNullException(nameof(work)); - var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - PostInternal( - () => + if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken); + var call = new WithdrawableCall(cancellationToken); + return PostWithdrawable(call, () => + { + using (call) { + if (!call.TryStart()) return; try { - tcs.TrySetResult(work()); + call.Completion.TrySetResult(work()); } catch (Exception ex) { - tcs.TrySetException(ex); + call.Completion.TrySetException(ex); } - }, - onDispatchFailure: ex => tcs.TrySetException(ex)); - return tcs.Task; + } + }); + } + + private Task PostWithdrawable(WithdrawableCall call, Action work) + { + try + { + PostInternal( + work, + // If the marshal itself fails (SynchronizationContext mode, Send throws before + // ever invoking the wrapped work), the wrapper's own try/catch never runs, so + // nothing would otherwise complete this task -- it would hang forever even + // though PostAsync failures are documented to surface via the returned task. + onDispatchFailure: ex => + { + call.Fail(ex); + call.Dispose(); + }); + } + catch + { + // Refused before anything was enqueued (the actor was disposed): no wrapper will + // ever run to release the registration on the caller's token, and no task is + // returned to hold it either. A long-lived token would keep this call alive. + call.Dispose(); + throw; + } + + return call.Completion.Task; + } + + /// + /// One PostAsync call that its token can still withdraw. The mailbox item and the + /// token's callback race for a single transition out of "queued": whichever wins decides + /// whether the work runs (the loop won, the token is too late) or is skipped and its task + /// cancelled (the token won). Exactly one of the two, so work that has started is never + /// reported as cancelled, and work that was reported as cancelled never runs. + /// + private sealed class WithdrawableCall : IDisposable + { + private const int Queued = 0, Started = 1, Withdrawn = 2; + + private readonly CancellationToken _token; + private readonly CancellationTokenRegistration _registration; + private int _state; + + internal WithdrawableCall(CancellationToken token) + { + _token = token; + Completion = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + // A token that can never be cancelled needs no registration and allocates none. + _registration = token.CanBeCanceled ? token.Register(Withdraw) : default; + } + + internal TaskCompletionSource Completion { get; } + + /// Claims the work for the loop; false if the token withdrew it first. + internal bool TryStart() => Interlocked.CompareExchange(ref _state, Started, Queued) == Queued; + + /// + /// Faults the task because the marshal failed before the work could start -- unless the + /// token withdrew the item first, in which case the task is (or is about to be) + /// cancelled and a fault must not overtake that. Decided by the same transition + /// and race for, not by whichever of + /// and + /// gets there first. + /// + internal void Fail(Exception exception) + { + if (Interlocked.CompareExchange(ref _state, Started, Queued) != Withdrawn) + Completion.TrySetException(exception); + } + + private void Withdraw() + { + if (Interlocked.CompareExchange(ref _state, Withdrawn, Queued) == Queued) + Completion.TrySetCanceled(_token); + } + + public void Dispose() => _registration.Dispose(); } private void PostInternal(Action work, Action? onDispatchFailure) diff --git a/src/CanKit.Pro.Actor/README.md b/src/CanKit.Pro.Actor/README.md index c635de8..f883b47 100644 --- a/src/CanKit.Pro.Actor/README.md +++ b/src/CanKit.Pro.Actor/README.md @@ -51,6 +51,10 @@ using var timeout = actor.Schedule(TimeSpan.FromMilliseconds(150), () => channel unrelated caller thread, never lost as an unobserved task exception. `PostAsync` failures surface through the returned task instead, since the caller is already positioned to observe them by awaiting. +- **`PostAsync` can be withdrawn while it is still queued.** A `CancellationToken` cancelled before + the call, or while the item waits in the mailbox, cancels the returned task at once and the work + never runs. Work that has already started is never interrupted: it runs through and the task + reports its result, so a half-finished work item never breaks the single-writer discipline. - **Configurable execution context** (FR-RAW-024): `ActorExecutionMode.DedicatedThread` (default) pins the loop to one real `Thread` for its entire lifetime — demonstrably the same thread for every callback. `ActorExecutionMode.ThreadPool` is cheaper for many short-lived instances but diff --git a/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.Actor.approved.txt b/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.Actor.approved.txt index 2ae5c6e..a559ce6 100644 --- a/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.Actor.approved.txt +++ b/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.Actor.approved.txt @@ -10,8 +10,8 @@ namespace CanKit.Pro.Actor { event System.EventHandler BackgroundExceptionOccurred; void Post(System.Action work); - System.Threading.Tasks.Task PostAsync(System.Action work); - System.Threading.Tasks.Task PostAsync(System.Func work); + System.Threading.Tasks.Task PostAsync(System.Action work, System.Threading.CancellationToken cancellationToken = default); + System.Threading.Tasks.Task PostAsync(System.Func work, System.Threading.CancellationToken cancellationToken = default); System.IDisposable Schedule(System.TimeSpan delay, System.Action callback); } public sealed class ProtocolActor : CanKit.Pro.Actor.IProtocolActor, System.IDisposable @@ -21,8 +21,8 @@ namespace CanKit.Pro.Actor public event System.EventHandler? BackgroundExceptionOccurred; public void Dispose() { } public void Post(System.Action work) { } - public System.Threading.Tasks.Task PostAsync(System.Action work) { } - public System.Threading.Tasks.Task PostAsync(System.Func work) { } + public System.Threading.Tasks.Task PostAsync(System.Action work, System.Threading.CancellationToken cancellationToken = default) { } + public System.Threading.Tasks.Task PostAsync(System.Func work, System.Threading.CancellationToken cancellationToken = default) { } public System.IDisposable Schedule(System.TimeSpan delay, System.Action callback) { } } } \ No newline at end of file diff --git a/tests/CanKit.Pro.Tests/Infrastructure/VirtualClock.cs b/tests/CanKit.Pro.Tests/Infrastructure/VirtualClock.cs index 4529953..0e0eed9 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/VirtualClock.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/VirtualClock.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Threading; using System.Threading.Tasks; using CanKit.Pro.Actor; @@ -23,7 +24,7 @@ namespace CanKit.Pro.Tests.Infrastructure; /// /// Why two round-trips. 's loop runs /// wait → DrainMailbox → DrainPendingTimerInserts → FireDueTimers. A single -/// completes during the drain, so awaiting it can +/// completes during the drain, so awaiting it can /// resume while that same iteration is still inside FireDueTimers. The second one is /// drained in the next iteration, which the loop reaches only after the first /// iteration's timer callbacks have returned — so awaiting it is a proof that they did, resting diff --git a/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs b/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs index 249e659..6c7045e 100644 --- a/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs @@ -330,10 +330,10 @@ public int RunQueuedWork() public IDisposable Schedule(TimeSpan delay, Action callback) => new NeverDue(); - public Task PostAsync(Action work) + public Task PostAsync(Action work, CancellationToken cancellationToken = default) => throw new NotSupportedException("The monitor only uses Post and Schedule; an ask would need a real loop."); - public Task PostAsync(Func work) + public Task PostAsync(Func work, CancellationToken cancellationToken = default) => throw new NotSupportedException("The monitor only uses Post and Schedule; an ask would need a real loop."); #pragma warning disable CS0067 // Nothing is run on the caller's behalf here, so this never fires. diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/HeartbeatProducerTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/HeartbeatProducerTests.cs index 73c83b8..251ee72 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/HeartbeatProducerTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/HeartbeatProducerTests.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Threading; using System.Threading.Tasks; using AwesomeAssertions; using CanKit.Pro.Actor; @@ -136,8 +137,8 @@ private sealed class RecordingActor : IProtocolActor public event EventHandler? BackgroundExceptionOccurred { add { } remove { } } public void Post(Action work) => work(); - public Task PostAsync(Action work) { work(); return Task.CompletedTask; } - public Task PostAsync(Func work) => Task.FromResult(work()); + public Task PostAsync(Action work, CancellationToken cancellationToken = default) { work(); return Task.CompletedTask; } + public Task PostAsync(Func work, CancellationToken cancellationToken = default) => Task.FromResult(work()); public IDisposable Schedule(TimeSpan delay, Action callback) { diff --git a/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs b/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs index 767236e..777a343 100644 --- a/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs @@ -279,9 +279,9 @@ public IDisposable Schedule(TimeSpan delay, Action callback) public void Post(Action work) => throw new NotSupportedException(); - public Task PostAsync(Action work) => throw new NotSupportedException(); + public Task PostAsync(Action work, CancellationToken cancellationToken = default) => throw new NotSupportedException(); - public Task PostAsync(Func work) => throw new NotSupportedException(); + public Task PostAsync(Func work, CancellationToken cancellationToken = default) => throw new NotSupportedException(); #pragma warning disable CS0067 public event EventHandler? BackgroundExceptionOccurred; diff --git a/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs b/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs index bdb71be..279effa 100644 --- a/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs @@ -82,6 +82,95 @@ public async Task Concurrent_Callers_From_Real_Threads_Produce_Exact_Count_No_Da counter.Should().Be(n); } + [Fact] + public async Task PostAsync_With_An_Already_Cancelled_Token_Is_Cancelled_And_Never_Runs() + { + using var actor = new ProtocolActor(); + using var cts = new CancellationTokenSource(); + cts.Cancel(); + var ran = false; + + var plain = actor.PostAsync(() => ran = true, cts.Token); + var generic = actor.PostAsync(() => ran = true, cts.Token); + + (await Record.ExceptionAsync(() => plain)).Should().BeAssignableTo(); + (await Record.ExceptionAsync(() => generic)).Should().BeAssignableTo(); + await actor.PostAsync(() => 0); + ran.Should().BeFalse("a token cancelled before the call withdraws the work"); + } + + [Fact] + public async Task PostAsync_Cancelled_While_Queued_Releases_The_Caller_At_Once_And_Skips_The_Work() + { + using var actor = new ProtocolActor(ActorExecutionMode.DedicatedThread); + using var release = new ManualResetEventSlim(); + using var started = new ManualResetEventSlim(); + actor.Post(() => { started.Set(); release.Wait(); }); + started.Wait(TimeSpan.FromSeconds(10)).Should().BeTrue(); + + using var cts = new CancellationTokenSource(); + var ran = 0; + var plain = actor.PostAsync(() => Interlocked.Increment(ref ran), cts.Token); + var generic = actor.PostAsync(() => Interlocked.Increment(ref ran), cts.Token); + + cts.Cancel(); + + // Both tasks are released while the loop is still blocked: cancelling does not wait for + // the item's turn in the mailbox. + (await Record.ExceptionAsync(() => WithinTenSeconds(plain))).Should().BeAssignableTo(); + (await Record.ExceptionAsync(() => WithinTenSeconds(generic))).Should().BeAssignableTo(); + + release.Set(); + await actor.PostAsync(() => 0); + ran.Should().Be(0, "an item withdrawn while queued is skipped when its turn comes"); + } + + [Fact] + public async Task PostAsync_Cancelled_After_The_Work_Started_Does_Not_Interrupt_It() + { + using var actor = new ProtocolActor(ActorExecutionMode.DedicatedThread); + using var release = new ManualResetEventSlim(); + using var started = new ManualResetEventSlim(); + using var cts = new CancellationTokenSource(); + + var task = actor.PostAsync(() => + { + started.Set(); + release.Wait(); + return 42; + }, cts.Token); + started.Wait(TimeSpan.FromSeconds(10)).Should().BeTrue(); + + // Past the point where the token can withdraw the item: the work runs through and the + // task reports its result instead of a cancellation. + cts.Cancel(); + release.Set(); + + (await task).Should().Be(42); + (await actor.PostAsync(() => 7)).Should().Be(7); + } + + // A cancelled-while-queued task that is never released would otherwise hang the run instead of + // failing it. Task.WaitAsync does not exist on every target framework this suite runs on. + private static async Task WithinTenSeconds(Task task) + { + var finished = await Task.WhenAny(task, Task.Delay(TimeSpan.FromSeconds(10))); + if (!ReferenceEquals(finished, task)) throw new TimeoutException("the task was not released"); + await task; + } + + [Fact] + public async Task PostAsync_With_A_Live_Token_Behaves_Like_Without_One() + { + using var actor = new ProtocolActor(); + using var cts = new CancellationTokenSource(); + + (await actor.PostAsync(() => 5, cts.Token)).Should().Be(5); + await actor.PostAsync(() => { }, cts.Token); + Func act = () => actor.PostAsync(() => throw new InvalidOperationException("boom"), cts.Token); + await act.Should().ThrowAsync().WithMessage("boom"); + } + [Fact] public async Task PostAsync_Propagates_Exception_Through_The_Returned_Task_Without_Raising_BackgroundEvent() { @@ -219,6 +308,35 @@ public async Task PostAsync_Fails_Its_Task_Instead_Of_Hanging_When_Synchronizati await act.Should().ThrowAsync(); } + [Fact] + public void PostAsync_On_A_Disposed_Actor_Throws_Synchronously_Also_With_A_Live_Token() + { + var actor = new ProtocolActor(); + actor.Dispose(); + using var cts = new CancellationTokenSource(); + + // Not a returned task that never completes: the refusal is the call's own exception, + // exactly as without a token. + Action plain = () => { _ = actor.PostAsync(() => { }, cts.Token); }; + Action generic = () => { _ = actor.PostAsync(() => 0, cts.Token); }; + plain.Should().Throw(); + generic.Should().Throw(); + } + + [Fact] + public async Task PostAsync_With_A_Live_Token_Faults_Its_Task_When_SynchronizationContext_Send_Throws() + { + using var actor = new ProtocolActor(ActorExecutionMode.SynchronizationContext, new AlwaysThrowingSynchronizationContext()); + using var cts = new CancellationTokenSource(); + + Func plain = () => actor.PostAsync(() => { }, cts.Token); + Func generic = () => actor.PostAsync(() => 0, cts.Token); + + // A fault, not a hang and not a cancellation: the token was never cancelled. + await plain.Should().ThrowAsync(); + await generic.Should().ThrowAsync(); + } + [Fact] public async Task Dispose_Rejects_New_Work_With_ObjectDisposedException() {