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
24 changes: 19 additions & 5 deletions src/CanKit.Pro.Actor/IProtocolActor.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System;
using System.Threading;
using System.Threading.Tasks;

namespace CanKit.Pro.Actor
Expand Down Expand Up @@ -30,13 +31,26 @@ public interface IProtocolActor : IDisposable
/// <paramref name="work"/> is surfaced through the returned task's fault, not through
/// <see cref="BackgroundExceptionOccurred"/> — the caller is already positioned to observe
/// it by awaiting.
/// <para>
/// <paramref name="cancellationToken"/> withdraws <paramref name="work"/> 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.
/// </para>
/// </summary>
Task PostAsync(Action work);
/// <param name="work">The work to run on the mailbox loop.</param>
/// <param name="cancellationToken">Withdraws the work while it is still queued.</param>
Task PostAsync(Action work, CancellationToken cancellationToken = default);

/// <summary>
/// Same as <see cref="PostAsync(Action)"/> but returns <paramref name="work"/>'s result.
/// Same as <see cref="PostAsync(Action, CancellationToken)"/> but returns
/// <paramref name="work"/>'s result.
/// </summary>
Task<T> PostAsync<T>(Func<T> work);
/// <param name="work">The work to run on the mailbox loop.</param>
/// <param name="cancellationToken">Withdraws the work while it is still queued.</param>
Task<T> PostAsync<T>(Func<T> work, CancellationToken cancellationToken = default);

/// <summary>
/// Schedules <paramref name="callback"/> to run on the actor's mailbox loop once
Expand All @@ -52,8 +66,8 @@ public interface IProtocolActor : IDisposable
/// Raised whenever a posted work item (via <see cref="Post"/>) or a scheduled callback
/// (via <see cref="Schedule"/>) 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 <see cref="PostAsync(Action)"/>/
/// <see cref="PostAsync{T}(Func{T})"/> failures, which surface through their own returned
/// failing item never stops the actor. Never raised for <see cref="PostAsync(Action, CancellationToken)"/>/
/// <see cref="PostAsync{T}(Func{T}, CancellationToken)"/> failures, which surface through their own returned
/// task instead.
/// </summary>
event EventHandler<Exception> BackgroundExceptionOccurred;
Expand Down
126 changes: 103 additions & 23 deletions src/CanKit.Pro.Actor/ProtocolActor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,7 @@ public sealed class ProtocolActor : IProtocolActor
/// <remarks>
/// <para>
/// 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 <see cref="PostAsync(Action)"/> and
/// requested work inline instead of routing it through <see cref="PostAsync(Action, CancellationToken)"/> 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.
Expand Down Expand Up @@ -302,50 +302,130 @@ public void Post(Action work)
}

/// <inheritdoc />
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<object?>(TaskCreationOptions.RunContinuationsAsynchronously);
PostInternal(
() =>
if (cancellationToken.IsCancellationRequested) return Task.FromCanceled(cancellationToken);
var call = new WithdrawableCall<object?>(cancellationToken);
Comment thread
dborgards marked this conversation as resolved.
Fixed
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
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;
}
});
}

/// <inheritdoc />
public Task<T> PostAsync<T>(Func<T> work)
public Task<T> PostAsync<T>(Func<T> work, CancellationToken cancellationToken = default)
{
if (work is null) throw new ArgumentNullException(nameof(work));
var tcs = new TaskCompletionSource<T>(TaskCreationOptions.RunContinuationsAsynchronously);
PostInternal(
() =>
if (cancellationToken.IsCancellationRequested) return Task.FromCanceled<T>(cancellationToken);
var call = new WithdrawableCall<T>(cancellationToken);
Comment thread
dborgards marked this conversation as resolved.
Fixed
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
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<T> PostWithdrawable<T>(WithdrawableCall<T> 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;
}

/// <summary>
/// One <c>PostAsync</c> 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.
/// </summary>
private sealed class WithdrawableCall<T> : 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<T>(TaskCreationOptions.RunContinuationsAsynchronously);
// A token that can never be cancelled needs no registration and allocates none.
_registration = token.CanBeCanceled ? token.Register(Withdraw) : default;
}

internal TaskCompletionSource<T> Completion { get; }

/// <summary>Claims the work for the loop; <c>false</c> if the token withdrew it first.</summary>
internal bool TryStart() => Interlocked.CompareExchange(ref _state, Started, Queued) == Queued;

/// <summary>
/// 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
/// <see cref="TryStart"/> and <see cref="Withdraw"/> race for, not by whichever of
/// <see cref="TaskCompletionSource{TResult}.TrySetCanceled()"/> and
/// <see cref="TaskCompletionSource{TResult}.TrySetException(Exception)"/> gets there first.
/// </summary>
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<Exception>? onDispatchFailure)
Expand Down
4 changes: 4 additions & 0 deletions src/CanKit.Pro.Actor/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@ namespace CanKit.Pro.Actor
{
event System.EventHandler<System.Exception> BackgroundExceptionOccurred;
void Post(System.Action work);
System.Threading.Tasks.Task PostAsync(System.Action work);
System.Threading.Tasks.Task<T> PostAsync<T>(System.Func<T> work);
System.Threading.Tasks.Task PostAsync(System.Action work, System.Threading.CancellationToken cancellationToken = default);
System.Threading.Tasks.Task<T> PostAsync<T>(System.Func<T> 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
Expand All @@ -21,8 +21,8 @@ namespace CanKit.Pro.Actor
public event System.EventHandler<System.Exception>? BackgroundExceptionOccurred;
public void Dispose() { }
public void Post(System.Action work) { }
public System.Threading.Tasks.Task PostAsync(System.Action work) { }
public System.Threading.Tasks.Task<T> PostAsync<T>(System.Func<T> work) { }
public System.Threading.Tasks.Task PostAsync(System.Action work, System.Threading.CancellationToken cancellationToken = default) { }
public System.Threading.Tasks.Task<T> PostAsync<T>(System.Func<T> work, System.Threading.CancellationToken cancellationToken = default) { }
public System.IDisposable Schedule(System.TimeSpan delay, System.Action callback) { }
}
}
3 changes: 2 additions & 1 deletion tests/CanKit.Pro.Tests/Infrastructure/VirtualClock.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using CanKit.Pro.Actor;

Expand All @@ -23,7 +24,7 @@ namespace CanKit.Pro.Tests.Infrastructure;
/// <para>
/// <b>Why two round-trips.</b> <see cref="ProtocolActor"/>'s loop runs
/// <c>wait → DrainMailbox → DrainPendingTimerInserts → FireDueTimers</c>. A single
/// <see cref="IProtocolActor.PostAsync(Action)"/> completes during the drain, so awaiting it can
/// <see cref="IProtocolActor.PostAsync(Action, CancellationToken)"/> completes during the drain, so awaiting it can
/// resume while that same iteration is still inside <c>FireDueTimers</c>. The second one is
/// drained in the <em>next</em> 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
Expand Down
4 changes: 2 additions & 2 deletions tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<T> PostAsync<T>(Func<T> work)
public Task<T> PostAsync<T>(Func<T> 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.
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using AwesomeAssertions;
using CanKit.Pro.Actor;
Expand Down Expand Up @@ -136,8 +137,8 @@ private sealed class RecordingActor : IProtocolActor
public event EventHandler<Exception>? BackgroundExceptionOccurred { add { } remove { } }

public void Post(Action work) => work();
public Task PostAsync(Action work) { work(); return Task.CompletedTask; }
public Task<T> PostAsync<T>(Func<T> work) => Task.FromResult(work());
public Task PostAsync(Action work, CancellationToken cancellationToken = default) { work(); return Task.CompletedTask; }
public Task<T> PostAsync<T>(Func<T> work, CancellationToken cancellationToken = default) => Task.FromResult(work());

public IDisposable Schedule(TimeSpan delay, Action callback)
{
Expand Down
4 changes: 2 additions & 2 deletions tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<T> PostAsync<T>(Func<T> work) => throw new NotSupportedException();
public Task<T> PostAsync<T>(Func<T> work, CancellationToken cancellationToken = default) => throw new NotSupportedException();

#pragma warning disable CS0067
public event EventHandler<Exception>? BackgroundExceptionOccurred;
Expand Down
Loading
Loading