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
12 changes: 12 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,18 @@ equivalent domain method) and return typed handles implementing both
a bounded buffer; a slow consumer terminates with
`SubscriptionBackpressureException` without terminating sibling handles.

Callback delivery uses an independent bounded queue configured with
`AsyncHandlerOptions.QueueCapacity` (default `1024`). Every subscription handle
exposes `Completion`: normal unsubscribe completes it, while callback queue
overflow faults it with `AsyncHandlerOverflowException`, terminates the local
registration, and is surfaced by async enumeration as the same typed failure.
RPC worker saturation continues to use broker-visible protocol backpressure.

Schedule backend unavailability and broker saturation use the distinct coded
error `FitzErrorCodes.ScheduleBackendError` (`7010`). `Retryability` classifies
it as retryable subject to operation safety; it is not reported as malformed
cron or parse input.

## Unreleased preview migration

This preview intentionally breaks the earlier callback subscription surface.
Expand Down
8 changes: 4 additions & 4 deletions src/Abstractions/Domains/Kv/KvSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,13 @@ namespace Cntryl.Fitz.Abstractions.Domains.Kv;

public sealed class KvSubscription : SubscriptionHandle<KvNotification>
{
public KvSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe)
: this(pattern, EmptyNotifications(), unsubscribe)
public KvSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: this(pattern, EmptyNotifications(), unsubscribe, completion)
{
}

public KvSubscription(string pattern, IAsyncEnumerable<KvNotification> notifications, Func<CancellationToken, ValueTask> unsubscribe)
: base(pattern, notifications, unsubscribe)
public KvSubscription(string pattern, IAsyncEnumerable<KvNotification> notifications, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: base(pattern, notifications, unsubscribe, completion)
{
}
}
9 changes: 5 additions & 4 deletions src/Abstractions/Domains/Lease/LeaseSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,16 @@ namespace Cntryl.Fitz.Abstractions.Domains.Lease;

public sealed class LeaseSubscription : SubscriptionHandle<LeaseChangeEvent>
{
public LeaseSubscription(string route, Func<CancellationToken, ValueTask> unsubscribe)
: this(route, EmptyNotifications(), unsubscribe)
public LeaseSubscription(string route, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: this(route, EmptyNotifications(), unsubscribe, completion)
{
}

public LeaseSubscription(
string route, IAsyncEnumerable<LeaseChangeEvent> notifications,
Func<CancellationToken, ValueTask> unsubscribe)
: base(route, notifications, unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
: base(route, notifications, unsubscribe, completion)
{
Route = route;
}
Expand Down
9 changes: 5 additions & 4 deletions src/Abstractions/Domains/Notice/NoticeSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,16 @@ namespace Cntryl.Fitz.Abstractions.Domains.Notice;

public sealed class NoticeSubscription : SubscriptionHandle<NoticeMessage>
{
public NoticeSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe)
: this(pattern, EmptyNotifications(), unsubscribe)
public NoticeSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: this(pattern, EmptyNotifications(), unsubscribe, completion)
{
}

public NoticeSubscription(
string pattern, IAsyncEnumerable<NoticeMessage> notifications,
Func<CancellationToken, ValueTask> unsubscribe)
: base(pattern, notifications, unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
: base(pattern, notifications, unsubscribe, completion)
{
}
}
9 changes: 5 additions & 4 deletions src/Abstractions/Domains/Queue/QueueSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,16 @@ namespace Cntryl.Fitz.Abstractions.Domains.Queue;

public sealed class QueueSubscription : SubscriptionHandle<QueueAvailabilityEvent>
{
public QueueSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe)
: this(pattern, EmptyNotifications(), unsubscribe)
public QueueSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: this(pattern, EmptyNotifications(), unsubscribe, completion)
{
}

public QueueSubscription(
string pattern, IAsyncEnumerable<QueueAvailabilityEvent> notifications,
Func<CancellationToken, ValueTask> unsubscribe)
: base(pattern, notifications, unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
: base(pattern, notifications, unsubscribe, completion)
{
}
}
9 changes: 5 additions & 4 deletions src/Abstractions/Domains/Schedule/ScheduleSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,16 @@ namespace Cntryl.Fitz.Abstractions.Domains.Schedule;

public sealed class ScheduleSubscription : SubscriptionHandle<ScheduleNotification>
{
public ScheduleSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe)
: this(pattern, EmptyNotifications(), unsubscribe)
public ScheduleSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: this(pattern, EmptyNotifications(), unsubscribe, completion)
{
}

public ScheduleSubscription(
string pattern, IAsyncEnumerable<ScheduleNotification> notifications,
Func<CancellationToken, ValueTask> unsubscribe)
: base(pattern, notifications, unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
: base(pattern, notifications, unsubscribe, completion)
{
}
}
9 changes: 5 additions & 4 deletions src/Abstractions/Domains/Stream/StreamSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,16 @@ namespace Cntryl.Fitz.Abstractions.Domains.Stream;

public sealed class StreamSubscription : SubscriptionHandle<StreamCommitEvent>
{
public StreamSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe)
: this(pattern, EmptyNotifications(), unsubscribe)
public StreamSubscription(string pattern, Func<CancellationToken, ValueTask> unsubscribe, Task? completion = null)
: this(pattern, EmptyNotifications(), unsubscribe, completion)
{
}

public StreamSubscription(
string pattern, IAsyncEnumerable<StreamCommitEvent> notifications,
Func<CancellationToken, ValueTask> unsubscribe)
: base(pattern, notifications, unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
: base(pattern, notifications, unsubscribe, completion)
{
}
}
1 change: 1 addition & 0 deletions src/Abstractions/FitzErrorCodes.cs
Original file line number Diff line number Diff line change
Expand Up @@ -21,4 +21,5 @@ public static class FitzErrorCodes
public const uint RpcSubscriptionLimit = 6013;
public const uint ScheduleInvalidSubscriptionPattern = 7006;
public const uint ScheduleSubscriptionLimit = 7007;
public const uint ScheduleBackendError = 7010;
}
70 changes: 66 additions & 4 deletions src/Abstractions/Runtime/SubscriptionHandle.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,18 +5,35 @@ namespace Cntryl.Fitz.Runtime;
public abstract class SubscriptionHandle : IAsyncDisposable
{
private readonly Func<CancellationToken, ValueTask> _unsubscribe;
private readonly TaskCompletionSource? _ownedCompletion;
private int _unsubscribed;

protected SubscriptionHandle(
string pattern,
Func<CancellationToken, ValueTask> unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
{
Pattern = pattern;
_unsubscribe = unsubscribe;
if (completion is null)
{
_ownedCompletion = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
Completion = _ownedCompletion.Task;
}
else
{
Completion = completion;
}
}

public string Pattern { get; }

/// <summary>
/// Completes after normal unsubscribe and faults if local notification
/// delivery terminates, including async-handler queue overflow.
/// </summary>
public Task Completion { get; }

public ValueTask UnsubscribeAsync(CancellationToken cancellationToken = default)
{
return UnsubscribeCoreAsync(cancellationToken);
Expand All @@ -35,7 +52,16 @@ private async ValueTask UnsubscribeCoreAsync(CancellationToken cancellationToken
return;
}

await _unsubscribe(cancellationToken).ConfigureAwait(false);
try
{
await _unsubscribe(cancellationToken).ConfigureAwait(false);
_ownedCompletion?.TrySetResult();
}
catch (Exception exception)
{
_ownedCompletion?.TrySetException(exception);
throw;
}
}
}

Expand All @@ -46,8 +72,9 @@ public abstract class SubscriptionHandle<T> : SubscriptionHandle, IAsyncEnumerab
protected SubscriptionHandle(
string pattern,
IAsyncEnumerable<T> notifications,
Func<CancellationToken, ValueTask> unsubscribe)
: base(pattern, unsubscribe)
Func<CancellationToken, ValueTask> unsubscribe,
Task? completion = null)
: base(pattern, unsubscribe, completion)
{
_notifications = notifications;
}
Expand Down Expand Up @@ -80,3 +107,38 @@ public SubscriptionBackpressureException(string message, Exception innerExceptio
{
}
}

public sealed class AsyncHandlerOverflowException : Exception
{
public const string ErrorCode = "ASYNC_HANDLER_OVERFLOW";

public AsyncHandlerOverflowException()
: this("The async handler queue overflowed.")
{
}

public AsyncHandlerOverflowException(string message)
: base(message)
{
Domain = "unknown";
Subscription = "unknown";
}

public AsyncHandlerOverflowException(string message, Exception innerException)
: base(message, innerException)
{
Domain = "unknown";
Subscription = "unknown";
}

public AsyncHandlerOverflowException(string domain, string subscription)
: base($"The async handler queue overflowed for {domain} subscription '{subscription}'.")
{
Domain = domain;
Subscription = subscription;
}

public string Code { get; } = ErrorCode;
public string Domain { get; }
public string Subscription { get; }
}
3 changes: 2 additions & 1 deletion src/Core/AsyncHandlerOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,5 +2,6 @@ namespace Cntryl.Fitz;

public sealed record AsyncHandlerOptions(
int? MaxConcurrency = null,
TimeSpan? Timeout = null
TimeSpan? Timeout = null,
int QueueCapacity = 1024
);
2 changes: 1 addition & 1 deletion src/Core/Connection/FitzConnection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ public FitzConnection(ClientConfig config, Func<ITransport> transportFactory)
_asyncHandlerDispatcher = new AsyncHandlerDispatcher(
_config.ResolvedAsyncHandlers.MaxConcurrency,
ResolveAsyncHandlerTimeout(),
Math.Min(Math.Max(_config.ResolvedMaxRequestQueueSize, 0), 1024),
Math.Max(_config.ResolvedAsyncHandlers.QueueCapacity, 0),
OnAsyncHandlerError,
onMetricsChanged: OnAsyncHandlerMetricsChanged,
onSaturated: OnAsyncHandlerSaturated);
Expand Down
15 changes: 12 additions & 3 deletions src/Core/Domains/Kv/KvClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -103,11 +103,12 @@ public async Task<KvSubscription> SubscribeAsync(
buffer.Write(notification);
return ValueTask.CompletedTask;
}, cancellationToken).ConfigureAwait(false);
buffer.ObserveCompletion(registration.Completion);
return new KvSubscription(pattern, buffer.ReadAllAsync(CancellationToken.None), async token =>
{
buffer.Complete();
await registration.UnsubscribeAsync(token).ConfigureAwait(false);
});
}, registration.Completion);
}

internal async Task<KvSubscription> SubscribeAsync(
Expand All @@ -127,8 +128,12 @@ internal async Task<KvSubscription> SubscribeAsync(
SingleReader = true,
SingleWriter = false
});
SubscriptionRegistration<KvNotification>? registration = new(channel);
var handleId = Interlocked.Increment(ref _nextHandleId);
SubscriptionRegistration<KvNotification>? registration = new(
channel,
"kv",
pattern,
token => UnsubscribeAsync(pattern, handleId, token));
await _subscriptionGate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
Expand All @@ -152,8 +157,12 @@ internal async Task<KvSubscription> SubscribeAsync(
state.Writers[handleId] = registration;
}
SubscriptionPump.Start(registration, handler, _dispatchAsyncHandler);
var completion = registration.Completion;
registration = null;
return new KvSubscription(pattern, token => UnsubscribeAsync(pattern, handleId, token));
return new KvSubscription(
pattern,
token => UnsubscribeAsync(pattern, handleId, token),
completion);
}
finally
{
Expand Down
18 changes: 12 additions & 6 deletions src/Core/Domains/Lease/LeaseClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -464,11 +464,12 @@ public async Task<LeaseSubscription> SubscribeAsync(
buffer.Write(notification);
return ValueTask.CompletedTask;
}, ct).ConfigureAwait(false);
buffer.ObserveCompletion(registration.Completion);
return new LeaseSubscription(route, buffer.ReadAllAsync(CancellationToken.None), async token =>
{
buffer.Complete();
await registration.UnsubscribeAsync(token).ConfigureAwait(false);
});
}, registration.Completion);
}

internal async Task<LeaseSubscription> SubscribeAsync(
Expand Down Expand Up @@ -500,14 +501,18 @@ internal async Task<LeaseSubscription> SubscribeAsync(

try
{
registration = new SubscriptionRegistration<LeaseChangeEvent>(channel);
registration = new SubscriptionRegistration<LeaseChangeEvent>(
channel,
"lease",
route,
token => UnsubscribeAsync(route, handleId, token));
await _subscriptionGate.WaitAsync(ct).ConfigureAwait(false);
gateAcquired = true;

if (_subscriptionsByRoute.TryGetValue(route, out var existingSubscription))
{
existingSubscription.Registrations[handleId] = registration;
var existingHandle = CreateSubscription(route, handleId);
var existingHandle = CreateSubscription(route, handleId, registration.Completion);
SubscriptionPump.Start(registration, handler, _dispatchAsyncHandler);
registration = null;
return existingHandle;
Expand All @@ -519,7 +524,7 @@ internal async Task<LeaseSubscription> SubscribeAsync(
_subscriptionsByRoute[route] = subscription;
_routesBySubscriptionId[subscriptionId] = route;

var handle = CreateSubscription(route, handleId);
var handle = CreateSubscription(route, handleId, registration.Completion);
SubscriptionPump.Start(registration, handler, _dispatchAsyncHandler);
registration = null;
return handle;
Expand All @@ -535,11 +540,12 @@ internal async Task<LeaseSubscription> SubscribeAsync(
}
}

private LeaseSubscription CreateSubscription(string route, long handleId)
private LeaseSubscription CreateSubscription(string route, long handleId, Task completion)
{
return new LeaseSubscription(
route,
cancellationToken => UnsubscribeAsync(route, handleId, cancellationToken));
cancellationToken => UnsubscribeAsync(route, handleId, cancellationToken),
completion);
}

private async Task<ulong> SubscribeWireAsync(string route, CancellationToken ct)
Expand Down
Loading
Loading