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()
{