From 6516b28c93aac25e234ce29a8c6bf656bbe9e39b Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 29 Sep 2026 20:20:26 +0000 Subject: [PATCH 1/4] feat(actor)!: let IProtocolActor.PostAsync be withdrawn with a CancellationToken PostAsync(Action) and PostAsync(Func) had no way to give up on work that was still waiting in the mailbox. Both now take an optional CancellationToken: a token cancelled before the call, or while the item is queued, cancels the returned task at once and the work never runs. Work that has started is never interrupted and its task reports the result, so the single-writer discipline (FR-RAW-021) never sees a half-finished item. The queued-to-started transition is one atomic step that the loop and the token race for, so work that started is never reported as cancelled and work that was reported as cancelled never runs. The ADR window before 1.3.0 is the last point where an interface member can be changed without a major version, so the signatures are replaced rather than overloaded next to the old ones. BREAKING CHANGE: IProtocolActor.PostAsync and PostAsync gained a CancellationToken parameter. Callers compile unchanged; implementers of IProtocolActor and code compiled against 1.2.x must be rebuilt. Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01UWRpkQzKkNDYz3WiNgWvWU --- src/CanKit.Pro.Actor/IProtocolActor.cs | 24 +++-- src/CanKit.Pro.Actor/ProtocolActor.cs | 83 ++++++++++++++--- src/CanKit.Pro.Actor/README.md | 4 + .../CanKit.Pro.Actor.approved.txt | 8 +- .../Infrastructure/VirtualClock.cs | 3 +- .../TestCases/BusStateMonitorTests.cs | 4 +- .../CANopen/HeartbeatProducerTests.cs | 5 +- .../TestCases/DeadlineTests.cs | 4 +- .../TestCases/ProtocolActorTests.cs | 89 +++++++++++++++++++ 9 files changed, 195 insertions(+), 29 deletions(-) diff --git a/src/CanKit.Pro.Actor/IProtocolActor.cs b/src/CanKit.Pro.Actor/IProtocolActor.cs index 31a30a93..c3e54d5c 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 2c7418cc..48573591 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,107 @@ 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); + if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken); + var call = new WithdrawableCall(cancellationToken); PostInternal( () => { + if (!call.TryStart()) return; try { work(); - tcs.TrySetResult(null); + call.Completion.TrySetResult(null); } catch (Exception ex) { - tcs.TrySetException(ex); + call.Completion.TrySetException(ex); + } + finally + { + call.Dispose(); } }, // 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; + onDispatchFailure: ex => + { + call.Completion.TrySetException(ex); + call.Dispose(); + }); + return call.Completion.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); + if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken); + var call = new WithdrawableCall(cancellationToken); PostInternal( () => { + if (!call.TryStart()) return; try { - tcs.TrySetResult(work()); + call.Completion.TrySetResult(work()); } catch (Exception ex) { - tcs.TrySetException(ex); + call.Completion.TrySetException(ex); + } + finally + { + call.Dispose(); } }, - onDispatchFailure: ex => tcs.TrySetException(ex)); - return tcs.Task; + onDispatchFailure: ex => + { + call.Completion.TrySetException(ex); + call.Dispose(); + }); + 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 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. + if (token.CanBeCanceled) _registration = token.Register(Withdraw); + } + + 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; + + 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 c635de84..f883b47f 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 2ae5c6ed..a559ce6d 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 45299538..0e0eed90 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 249e659c..6c7045e6 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 73c83b8e..251ee722 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 767236ec..777a3434 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 bdb71be7..2f189ca9 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() { From dd1a4ab9af4afb44d1b6d92cf929002200f441c8 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 29 Sep 2026 20:30:14 +0000 Subject: [PATCH 2/4] fix(actor): release the token registration when PostAsync is refused PostAsync registers with the caller's token before it enqueues. When the actor is disposed first, PostInternal throws before anything is enqueued, so no mailbox wrapper ever runs to dispose that registration and no task is returned to hold it. A long-lived token kept every such call alive until it was cancelled or disposed. The enqueue is now wrapped once, in a helper both overloads share, and disposes the registration when it throws. The dispatch-failure handling moves into the same helper so it exists once instead of twice. Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01UWRpkQzKkNDYz3WiNgWvWU --- src/CanKit.Pro.Actor/ProtocolActor.cs | 94 +++++++++++++++------------ 1 file changed, 52 insertions(+), 42 deletions(-) diff --git a/src/CanKit.Pro.Actor/ProtocolActor.cs b/src/CanKit.Pro.Actor/ProtocolActor.cs index 48573591..4dd98ab9 100644 --- a/src/CanKit.Pro.Actor/ProtocolActor.cs +++ b/src/CanKit.Pro.Actor/ProtocolActor.cs @@ -307,34 +307,23 @@ public Task PostAsync(Action work, CancellationToken cancellationToken = default if (work is null) throw new ArgumentNullException(nameof(work)); if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken); var call = new WithdrawableCall(cancellationToken); - PostInternal( - () => + return PostWithdrawable(call, () => + { + if (!call.TryStart()) return; + try { - if (!call.TryStart()) return; - try - { - work(); - call.Completion.TrySetResult(null); - } - catch (Exception ex) - { - call.Completion.TrySetException(ex); - } - finally - { - call.Dispose(); - } - }, - // 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 => + work(); + call.Completion.TrySetResult(null); + } + catch (Exception ex) { call.Completion.TrySetException(ex); + } + finally + { call.Dispose(); - }); - return call.Completion.Task; + } + }); } /// @@ -343,28 +332,49 @@ public Task PostAsync(Func work, CancellationToken cancellationToken = if (work is null) throw new ArgumentNullException(nameof(work)); if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken); var call = new WithdrawableCall(cancellationToken); - PostInternal( - () => + return PostWithdrawable(call, () => + { + if (!call.TryStart()) return; + try { - if (!call.TryStart()) return; - try - { - call.Completion.TrySetResult(work()); - } - catch (Exception ex) - { - call.Completion.TrySetException(ex); - } - finally - { - call.Dispose(); - } - }, - onDispatchFailure: ex => + call.Completion.TrySetResult(work()); + } + catch (Exception ex) { call.Completion.TrySetException(ex); + } + finally + { call.Dispose(); - }); + } + }); + } + + 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.Completion.TrySetException(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; } From f9639e7a9eea9d8837899b739523a27ca68451c0 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 29 Sep 2026 20:34:53 +0000 Subject: [PATCH 3/4] fix(actor): dispose a withdrawn PostAsync call and tidy its ownership When the token won the queued-to-started race, the mailbox wrapper returned before the try/finally that disposes the call, so the registration on the token was never released on that path. The wrapper now owns the call with a using statement, which covers the withdrawn path as well as the completed and faulted ones. The registration field becomes readonly, assigned once in the constructor, which is what the analyser asked for and what the type always intended. Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01UWRpkQzKkNDYz3WiNgWvWU --- src/CanKit.Pro.Actor/ProtocolActor.cs | 49 +++++++++++++-------------- 1 file changed, 24 insertions(+), 25 deletions(-) diff --git a/src/CanKit.Pro.Actor/ProtocolActor.cs b/src/CanKit.Pro.Actor/ProtocolActor.cs index 4dd98ab9..ec74e28c 100644 --- a/src/CanKit.Pro.Actor/ProtocolActor.cs +++ b/src/CanKit.Pro.Actor/ProtocolActor.cs @@ -309,19 +309,19 @@ public Task PostAsync(Action work, CancellationToken cancellationToken = default var call = new WithdrawableCall(cancellationToken); return PostWithdrawable(call, () => { - if (!call.TryStart()) return; - try - { - work(); - call.Completion.TrySetResult(null); - } - catch (Exception ex) - { - call.Completion.TrySetException(ex); - } - finally + // Disposed on every path out, including the one where the token withdrew the item. + using (call) { - call.Dispose(); + if (!call.TryStart()) return; + try + { + work(); + call.Completion.TrySetResult(null); + } + catch (Exception ex) + { + call.Completion.TrySetException(ex); + } } }); } @@ -334,18 +334,17 @@ public Task PostAsync(Func work, CancellationToken cancellationToken = var call = new WithdrawableCall(cancellationToken); return PostWithdrawable(call, () => { - if (!call.TryStart()) return; - try - { - call.Completion.TrySetResult(work()); - } - catch (Exception ex) - { - call.Completion.TrySetException(ex); - } - finally + using (call) { - call.Dispose(); + if (!call.TryStart()) return; + try + { + call.Completion.TrySetResult(work()); + } + catch (Exception ex) + { + call.Completion.TrySetException(ex); + } } }); } @@ -390,7 +389,7 @@ private sealed class WithdrawableCall : IDisposable private const int Queued = 0, Started = 1, Withdrawn = 2; private readonly CancellationToken _token; - private CancellationTokenRegistration _registration; + private readonly CancellationTokenRegistration _registration; private int _state; internal WithdrawableCall(CancellationToken token) @@ -398,7 +397,7 @@ internal WithdrawableCall(CancellationToken token) _token = token; Completion = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); // A token that can never be cancelled needs no registration and allocates none. - if (token.CanBeCanceled) _registration = token.Register(Withdraw); + _registration = token.CanBeCanceled ? token.Register(Withdraw) : default; } internal TaskCompletionSource Completion { get; } From 78b4cd8cc41cf8f4eccfc9a645f25e78a9191ebc Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 29 Sep 2026 21:25:40 +0000 Subject: [PATCH 4/4] fix(actor): let the withdraw transition decide a dispatch failure too In SynchronizationContext mode a failing Send completes the task from the dispatch-failure handler. That handler bypassed the queued-to-started transition the loop and the token race for, so a token that had already won it and was about to cancel the task could be overtaken by the fault: work withdrawn before it started would surface as faulted instead of cancelled. The handler now goes through the same transition and faults the task only when the token has not withdrawn the item first. Two tests cover the neighbouring paths that had none: PostAsync on a disposed actor with a token still throws synchronously instead of returning a task that never completes, and a failing Send with a live token faults the task. Co-Authored-By: Claude Sonnet 5.5 Claude-Session: https://claude.ai/code/session_01UWRpkQzKkNDYz3WiNgWvWU --- src/CanKit.Pro.Actor/ProtocolActor.cs | 16 +++++++++- .../TestCases/ProtocolActorTests.cs | 29 +++++++++++++++++++ 2 files changed, 44 insertions(+), 1 deletion(-) diff --git a/src/CanKit.Pro.Actor/ProtocolActor.cs b/src/CanKit.Pro.Actor/ProtocolActor.cs index ec74e28c..26e24668 100644 --- a/src/CanKit.Pro.Actor/ProtocolActor.cs +++ b/src/CanKit.Pro.Actor/ProtocolActor.cs @@ -361,7 +361,7 @@ private Task PostWithdrawable(WithdrawableCall call, Action work) // though PostAsync failures are documented to surface via the returned task. onDispatchFailure: ex => { - call.Completion.TrySetException(ex); + call.Fail(ex); call.Dispose(); }); } @@ -405,6 +405,20 @@ internal WithdrawableCall(CancellationToken token) /// 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) diff --git a/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs b/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs index 2f189ca9..279effae 100644 --- a/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs @@ -308,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() {