From 0d0442ca5924e39b255cd1e71119588e0f493dd3 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 15:56:57 +0200 Subject: [PATCH 1/7] perf(kurrentdb): pass the resolved event to ack and nack instead of context items Persistent subscriptions stashed the ResolvedEvent (a struct, so boxed) and the PersistentSubscription in the consume context items for every event, only for Ack and Nack to read them back. That boxed the event and forced the lazily allocated items dictionary into existence on each message. Both are now passed to Ack and Nack as parameters. The per-event Logger.Configure call is replaced with Logger.Current = Log. The base class Log is built from the same subscription id and logger factory, so it is equivalent, and it avoids building a new LogContext when the gRPC callback does not carry the subscription's AsyncLocal. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../PersistentSubscriptionBase.cs | 26 +++++-------------- 1 file changed, 6 insertions(+), 20 deletions(-) diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs index a3c10b44..34158d9f 100644 --- a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs +++ b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs @@ -99,9 +99,6 @@ protected PersistentSubscriptionBase( if (options is { FailureHandler: not null, ThrowOnError: false }) Log.ThrowOnErrorIncompatible(); } - const string ResolvedEventKey = "resolvedEvent"; - const string SubscriptionKey = "subscription"; - /// /// Execute an operation to set up a persistent subscription /// @@ -135,20 +132,18 @@ void HandleDrop(PersistentSubscription __, SubscriptionDroppedReason reason, Exc => run.Fail(KurrentDBMappings.AsDropReason(reason), exception); async Task HandleEvent(PersistentSubscription subscription, ResolvedEvent re, int? retryCount, CancellationToken ct) { - Logger.Configure(Options.SubscriptionId, LoggerFactory); + Logger.Current = Log; - var context = CreateContext(run, re, ct) - .WithItem(ResolvedEventKey, re) - .WithItem(SubscriptionKey, subscription); + var context = CreateContext(run, re, ct); try { await Handler(context).NoContext(); LastProcessed = EventPosition.FromContext(context); - await Ack(context).NoContext(); + await Ack(subscription, re).NoContext(); } catch (OperationCanceledException) when (ct.IsCancellationRequested) { // Its own token was cancelled: the supervisor already knows the run is over. } catch (Exception e) { - await Nack(context, e).NoContext(); + await Nack(context, subscription, re, e).NoContext(); } } } @@ -171,16 +166,9 @@ protected abstract Task LocalSubscribe( CancellationToken cancellationToken ); - // ReSharper disable once MemberCanBeMadeStatic.Local -#pragma warning disable CA1822 - async ValueTask Ack(MessageConsumeContext ctx) { -#pragma warning restore CA1822 - var re = ctx.Items.GetItem(ResolvedEventKey); - var subscription = ctx.Items.GetItem(SubscriptionKey)!; - await subscription.Ack(re).NoContext(); - } + static async ValueTask Ack(PersistentSubscription subscription, ResolvedEvent re) => await subscription.Ack(re).NoContext(); - async ValueTask Nack(MessageConsumeContext ctx, Exception exception) { + async ValueTask Nack(MessageConsumeContext ctx, PersistentSubscription subscription, ResolvedEvent re, Exception exception) { if (exception is OperationCanceledException && ctx.CancellationToken.IsCancellationRequested) { return; } @@ -191,8 +179,6 @@ async ValueTask Nack(MessageConsumeContext ctx, Exception exception) { ctx.LogContext.MessageHandlingFailed(Options.SubscriptionId, ctx, exception); } - var re = ctx.Items.GetItem(ResolvedEventKey); - var subscription = ctx.Items.GetItem(SubscriptionKey)!; await _handleEventProcessingFailure(Client, subscription, re, exception).NoContext(); } From 07b99adb12499f2d10721f3f0f6ffd1a33a4c9ed Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 16:07:59 +0200 Subject: [PATCH 2/7] perf(subscriptions): handle payload-less events without a scope or an async context A context without a payload never enters the consume pipe: it is marked ignored and, where positions are checkpointed, acknowledged so the checkpoint can move past it. It still paid for everything around the pipe: the logging scope and its array, the activity check, and on checkpointed subscriptions an AsyncConsumeContext wrapper that only existed to carry the ack. Handler now branches on the payload first and sends payload-less contexts to a small helper that only ignores and acknowledges, with the same failure handling as before. HandleInternal calls the same helper with the run's own ack, so the wrapper is no longer allocated. Measured on a checkpointed subscription: 609 -> 457 bytes per payload-less event. Events with a payload are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../EventSubscription.cs | 58 ++-- .../EventSubscriptionWithCheckpoint.cs | 12 +- .../PayloadlessContextTests.cs | 275 ++++++++++++++++++ 3 files changed, 325 insertions(+), 20 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/PayloadlessContextTests.cs diff --git a/src/Core/src/Eventuous.Subscriptions/EventSubscription.cs b/src/Core/src/Eventuous.Subscriptions/EventSubscription.cs index 8c639778..146ac15e 100644 --- a/src/Core/src/Eventuous.Subscriptions/EventSubscription.cs +++ b/src/Core/src/Eventuous.Subscriptions/EventSubscription.cs @@ -237,9 +237,42 @@ void ReportConnectionDropped(Session session, Failure failure) { string GetActivityName(string? messageType) => _activityNames.GetOrAdd(messageType ?? "", static (type, prefix) => prefix + type, _activityNamePrefix); + protected ValueTask Handler(IMessageConsumeContext context) + => context.Message == null + ? HandleWithoutPayload(context, context is AsyncConsumeContext ? AcknowledgeAsyncContext : null) + : HandleWithPayload(context); + + static readonly Acknowledge AcknowledgeAsyncContext = static context => ((AsyncConsumeContext)context).Acknowledge(); + + /// + /// A context without a payload is ignored and acknowledged without entering the pipe, so it skips + /// what only the pipe needs: the logging scope and the activity. Hot, since checkpoint-reached + /// contexts arrive payload-less. + /// + private protected async ValueTask HandleWithoutPayload(IMessageConsumeContext context, Acknowledge? acknowledge) { + // ReSharper disable once NullCoalescingConditionIsAlwaysNotNullAccordingToAPIContract + Logger.Current ??= Log; + + Log.MessageReceived(context); + + try { + context.Ignore(SubscriptionId); + + if (acknowledge != null) await acknowledge(context).NoContext(); + } catch (OperationCanceledException e) when (context.CancellationToken.IsCancellationRequested) { + Log.MessageIgnoredWhenStopping(e); + } catch (Exception e) { context.Nack(SubscriptionId, e); } + + if (context.HasFailed() && Options.ThrowOnError) { + var exception = context.HandlingResults.GetException(); + + throw new SubscriptionException(context.Stream, context.MessageType, context.Message, exception ?? new InvalidOperationException()); + } + } + // ReSharper disable once CognitiveComplexity // ReSharper disable once CyclomaticComplexity - protected async ValueTask Handler(IMessageConsumeContext context) { + async ValueTask HandleWithPayload(IMessageConsumeContext context) { // Use KeyValuePair array instead of Dictionary for 5x speedup and 3x less allocation var scope = new KeyValuePair[] { new("SubscriptionId", SubscriptionId), @@ -251,10 +284,7 @@ protected async ValueTask Handler(IMessageConsumeContext context) { Logger.Current ??= Log; using (Log.Logger.BeginScope(scope)) { - // No activity for payload-less contexts: they are ignored and acknowledged below without - // entering the pipe, so an activity would never be started or disposed on the async path — - // a pure allocation leak, hot since checkpoint-reached contexts arrive payload-less. - var activity = EventuousDiagnostics.Enabled && context.Message != null + var activity = EventuousDiagnostics.Enabled ? SubscriptionActivity.Create( GetActivityName(context.MessageType), ActivityKind.Internal, @@ -269,23 +299,13 @@ protected async ValueTask Handler(IMessageConsumeContext context) { Log.MessageReceived(context); try { - if (context.Message != null) { - if (activity != null) { - context.ParentContext = activity.Context; + if (activity != null) { + context.ParentContext = activity.Context; - if (isAsync) { context.Items.AddItem(ContextItemKeys.Activity, activity); } - } - - await Pipe.Send(context).NoContext(); + if (isAsync) { context.Items.AddItem(ContextItemKeys.Activity, activity); } } - else { - context.Ignore(SubscriptionId); - if (isAsync) { - var asyncContext = context as AsyncConsumeContext; - await asyncContext!.Acknowledge().NoContext(); - } - } + await Pipe.Send(context).NoContext(); if (context.WasIgnored() && activity != null) activity.ActivityTraceFlags = ActivityTraceFlags.None; } catch (OperationCanceledException e) when (context.CancellationToken.IsCancellationRequested) { diff --git a/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs b/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs index f22397fd..c7bf95e1 100644 --- a/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs +++ b/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs @@ -113,7 +113,17 @@ protected async ValueTask HandleInternal(SubscriptionRun run, IMessageConsumeCon Logger.Current = Log; var checkpointedRun = (CheckpointedRun)run; - var ctx = new AsyncConsumeContext(context, checkpointedRun.AckMessage, checkpointedRun.NackMessage); + + // Never reaches the pipe, so nothing needs the wrapper: the run's own ack takes the context as is. + if (context.Message == null) { + // ReSharper disable once NullCoalescingConditionIsAlwaysNotNullAccordingToAPIContract + context.LogContext ??= Log; + await HandleWithoutPayload(context, checkpointedRun.AckMessage).NoContext(); + + return; + } + + var ctx = new AsyncConsumeContext(context, checkpointedRun.AckMessage, checkpointedRun.NackMessage); await Handler(ctx).NoContext(); } catch (OperationCanceledException e) when (context.CancellationToken.IsCancellationRequested) { context.LogContext.MessageHandlingFailed(Options.SubscriptionId, context, e); diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/PayloadlessContextTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/PayloadlessContextTests.cs new file mode 100644 index 00000000..0875bd3a --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/PayloadlessContextTests.cs @@ -0,0 +1,275 @@ +// Copyright (C) Eventuous HQ OÜ. All rights reserved +// Licensed under the Apache License, Version 2.0. + +using System.Collections.Concurrent; +using Eventuous.Subscriptions; +using Eventuous.Subscriptions.Checkpoints; +using Eventuous.Subscriptions.Context; +using Eventuous.Subscriptions.Filters; +using Eventuous.Tools; +using Microsoft.Extensions.Logging.Abstractions; +using Shouldly; + +namespace Eventuous.Tests.Subscriptions; + +/// +/// A context without a payload (empty data, unknown type, failed deserialization, a checkpoint-reached +/// marker) is never handled, but it still has to be accounted for: marked ignored and, where positions +/// are checkpointed, committed — a position that's never committed is a gap the checkpoint can't cross. +/// +public class PayloadlessContextTests { + const string Id = "payload-less"; + + [Test] + public async Task Checkpointed_subscription_ignores_it_and_commits_its_position(CancellationToken ct) { + var (subscription, handler, committed) = await StartCheckpointed(ct); + + var context = subscription.CreateContext(0, null); + await subscription.Deliver(context); + + handler.Positions.ShouldBeEmpty("a context without a payload must not reach the handler"); + ShouldBeIgnoredBySubscription(context); + context.LogContext.ShouldNotBeNull("a context delivered without a log context gets the subscription's"); + + (await Wait.Until(() => committed.Contains((ulong?)0), TimeSpan.FromSeconds(5))) + .ShouldBeTrue("the position of an ignored context should still be committed"); + + await subscription.Unsubscribe(_ => { }, ct); + } + + /// + /// Sequences are contiguous per run, so if the middle one were never committed the checkpoint would + /// stay parked at the first event however many followed. + /// + [Test] + public async Task Checkpoint_advances_past_one_sitting_between_two_events(CancellationToken ct) { + var (subscription, handler, committed) = await StartCheckpointed(ct); + + var first = subscription.CreateContext(0, new { Position = 0 }); + var middle = subscription.CreateContext(1, null); + var last = subscription.CreateContext(2, new { Position = 2 }); + + await subscription.Deliver(first); + await subscription.Deliver(middle); + await subscription.Deliver(last); + + (await Wait.Until(() => committed.Contains((ulong?)2), TimeSpan.FromSeconds(5))) + .ShouldBeTrue($"the checkpoint should reach the last event, but only got to [{string.Join(", ", committed)}]"); + + handler.Positions.ShouldBe([0UL, 2UL]); + ShouldBeIgnoredBySubscription(middle); + first.WasIgnored().ShouldBeFalse(); + last.WasIgnored().ShouldBeFalse(); + + await subscription.Unsubscribe(_ => { }, ct); + } + + [Test] + public async Task Plain_subscription_ignores_it_without_entering_the_pipe() { + var filter = new CountingFilter(); + var handler = new RecordingHandler(); + var logs = new CapturingLoggerFactory(LogLevel.Trace); + + await using var subscription = new PlainSubscription(new ConsumePipe().AddFilterFirst(filter).AddDefaultConsumer(handler), loggerFactory: logs); + + var context = CreateContext(0, 0, null); + await subscription.Deliver(context); + + filter.Seen.ShouldBe(0, "a context without a payload must not enter the pipe"); + handler.Positions.ShouldBeEmpty(); + ShouldBeIgnoredBySubscription(context); + logs.Contains("Received TestEvent").ShouldBeTrue(); + logs.Contains("payload-less ignored TestEvent").ShouldBeTrue(); + + var real = CreateContext(1, 1, new { Position = 1 }); + await subscription.Deliver(real); + + filter.Seen.ShouldBe(1); + handler.Positions.ShouldBe([1UL]); + } + + [Test] + public async Task Plain_subscription_acknowledges_one_that_carries_its_own_acknowledgement() { + var handler = new RecordingHandler(); + var acked = new List(); + + await using var subscription = new PlainSubscription(new ConsumePipe().AddDefaultConsumer(handler)); + + var context = new AsyncConsumeContext( + CreateContext(0, 0, null), + ctx => { + acked.Add(ctx); + + return default; + }, + (_, e) => throw e + ); + + await subscription.Deliver(context); + + acked.ShouldBe([context]); + handler.Positions.ShouldBeEmpty(); + ShouldBeIgnoredBySubscription(context); + } + + /// + /// Pins what happens today rather than what reads as intended: the failure is reported under the + /// subscription's name, which already holds the ignored result, and results keep one entry per name. + /// So the failure is logged but never recorded, and ThrowOnError has nothing to throw. + /// + [Test] + [Arguments(false)] + [Arguments(true)] + public async Task Failed_acknowledgement_is_logged_and_leaves_the_context_ignored(bool throwOnError) { + var logs = new CapturingLoggerFactory(); + + await using var subscription = new PlainSubscription(new ConsumePipe().AddDefaultConsumer(new RecordingHandler()), throwOnError, logs); + + var failure = new InvalidOperationException("ack failed"); + var context = new AsyncConsumeContext(CreateContext(0, 0, null), _ => throw failure, (_, e) => throw e); + + await subscription.Deliver(context); + + context.HasFailed().ShouldBeFalse(); + ShouldBeIgnoredBySubscription(context); + logs.Contains("Message handling failed at payload-less").ShouldBeTrue(); + logs.Contains("ack failed").ShouldBeTrue(); + } + + /// + /// An acknowledgement cancelled by the context's own token means the subscription is stopping, which + /// isn't a failure of the message. + /// + [Test] + public async Task Acknowledgement_cancelled_by_shutdown_is_not_a_failure() { + using var cts = new CancellationTokenSource(); + + await using var subscription = new PlainSubscription(new ConsumePipe().AddDefaultConsumer(new RecordingHandler()), true); + + var inner = CreateContext(0, 0, null); + inner.CancellationToken = cts.Token; + + var context = new AsyncConsumeContext( + inner, + _ => { + cts.Cancel(); + + throw new OperationCanceledException(cts.Token); + }, + (_, e) => throw e + ); + + await subscription.Deliver(context); + + context.HasFailed().ShouldBeFalse(); + ShouldBeIgnoredBySubscription(context); + } + + static void ShouldBeIgnoredBySubscription(IMessageConsumeContext context) { + context.WasIgnored().ShouldBeTrue(); + context.HandlingResults.GetResultsOf(EventHandlingStatus.Ignored).Single().HandlerType.ShouldBe(Id); + } + + static async Task<(CheckpointedSubscription, RecordingHandler, ConcurrentQueue)> StartCheckpointed(CancellationToken ct) { + var checkpointStore = new NoOpCheckpointStore(); + var committed = new ConcurrentQueue(); + checkpointStore.CheckpointStored += (_, cp) => committed.Enqueue(cp.Position); + + var handler = new RecordingHandler(); + + var options = new CheckpointedOptions { + SubscriptionId = Id, + CheckpointCommitBatchSize = 1, + CheckpointCommitDelayMs = 10 + }; + + var subscription = new CheckpointedSubscription(options, checkpointStore, new ConsumePipe().AddDefaultConsumer(handler)); + + await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); + + return (subscription, handler, committed); + } + + static MessageConsumeContext CreateContext(ulong position, ulong sequence, object? message) + => new( + Guid.NewGuid().ToString(), + "TestEvent", + "application/json", + "test-stream", + position, + position, + position, + sequence, + DateTime.UtcNow, + message, + null, + Id, + CancellationToken.None + ); + + record PlainOptions : SubscriptionOptions; + + record CheckpointedOptions : SubscriptionWithCheckpointOptions; + + sealed class PlainSubscription(ConsumePipe pipe, bool throwOnError = false, ILoggerFactory? loggerFactory = null) + : EventSubscription(new() { SubscriptionId = Id, ThrowOnError = throwOnError }, pipe, loggerFactory ?? NullLoggerFactory.Instance, null) { + protected override ValueTask Connect(SubscriptionRun run) => default; + + public ValueTask Deliver(IMessageConsumeContext context) { + context.LogContext = Log; + + return Handler(context); + } + } + + /// + /// Has no pump of its own: the test delivers into the connected run directly. + /// + sealed class CheckpointedSubscription(CheckpointedOptions options, ICheckpointStore checkpointStore, ConsumePipe pipe) + : EventSubscriptionWithCheckpoint(options, checkpointStore, pipe, 1, SubscriptionKind.All, NullLoggerFactory.Instance, null, null) { + SubscriptionRun? _run; + + protected override async ValueTask Connect(SubscriptionRun run) { + await GetCheckpoint(run).NoContext(); + + Volatile.Write(ref _run, run); + } + + // Without a log context, as a transport that never set one would deliver it. + public MessageConsumeContext CreateContext(ulong position, object? message) { + var run = Volatile.Read(ref _run)!; + var context = PayloadlessContextTests.CreateContext(position, run.NextSequence(), message); + + context.LogContext = null!; + context.CancellationToken = run.Token; + + return context; + } + + public ValueTask Deliver(IMessageConsumeContext context) => HandleInternal(Volatile.Read(ref _run)!, context); + } + + sealed class RecordingHandler : BaseEventHandler { + readonly ConcurrentQueue _positions = new(); + + public ulong[] Positions => _positions.ToArray(); + + public override ValueTask HandleEvent(IMessageConsumeContext context) { + _positions.Enqueue(context.GlobalPosition); + + return new(EventHandlingStatus.Success); + } + } + + sealed class CountingFilter : ConsumeFilter { + int _seen; + + public int Seen => Volatile.Read(ref _seen); + + protected override ValueTask Send(IMessageConsumeContext context, LinkedListNode? next) { + Interlocked.Increment(ref _seen); + + return next?.Value.Send(context, next.Next) ?? default; + } + } +} From a2cb547c465123e755e5f0240d8cf70fc495d885 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 16:26:59 +0200 Subject: [PATCH 3/7] perf(subscriptions): wait on the channel without allocating per wake The reader loop of ConcurrentChannelWorker, which AsyncHandlingFilter runs handlers on, read with the worker's cancellable token. A bounded channel serves a read on an empty channel from an operation it keeps and reuses only when the token cannot be cancelled; with a cancellable one it allocates a new operation and a cancellation registration each time. An empty channel is the normal state of a caught-up subscription, so that was paid per event. The loop now reads with CancellationToken.None and keeps checking the worker token between elements. Measured in isolation on an idle bounded channel, one item at a time (net10.0): 182 -> 82 B/item with a single reader, and 182 -> 155 B/item with four readers, where only one parked reader at a time can hold the reusable operation. WaitToReadAsync + TryRead does the same for one reader but wakes every parked reader per write and came out at 350 B/item with four, so it is not used. The allocation harness feeds events as fast as it can, the channel is rarely empty there, and its four scenarios are unchanged within noise. Shutdown now relies on an invariant that already held: the worker token is never cancelled before the channel is completed. An idle reader is woken by completion alone, no longer by the token. The token is private to ChannelWorkerBase and is only cancelled by Stop's drain deadline and StopWorker's finally, both after Stop has completed the channel. Elements still queued when the token is cancelled are left unprocessed as before. Tests pin the worker's shutdown directly: an idle worker stops promptly, queued elements are drained, several writers feed several readers, and a drain that outlives its deadline leaves the rest unprocessed. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Channels/ChannelExtensions.cs | 7 +- .../ConcurrentChannelWorkerTests.cs | 154 ++++++++++++++++++ 2 files changed, 160 insertions(+), 1 deletion(-) create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/ConcurrentChannelWorkerTests.cs diff --git a/src/Core/src/Eventuous.Subscriptions/Channels/ChannelExtensions.cs b/src/Core/src/Eventuous.Subscriptions/Channels/ChannelExtensions.cs index b45c7ecf..f00cb5e5 100644 --- a/src/Core/src/Eventuous.Subscriptions/Channels/ChannelExtensions.cs +++ b/src/Core/src/Eventuous.Subscriptions/Channels/ChannelExtensions.cs @@ -13,7 +13,11 @@ static class ChannelExtensions { public async Task Read(ProcessElement process, CancellationToken cancellationToken) { try { while (!cancellationToken.IsCancellationRequested) { - var element = await channel.Reader.ReadAsync(cancellationToken).NoContext(); + // Not cancellable on purpose: a bounded channel parks such a read on an operation it keeps + // and reuses, where a cancellable one allocates per wake — per event, on a caught-up + // subscription. So an idle reader is woken only by the channel completing, never by the + // token, which must therefore not be cancelled before the channel is completed (see Stop). + var element = await channel.Reader.ReadAsync(CancellationToken.None).NoContext(); await process(element, cancellationToken).NoContext(); } } catch (OperationCanceledException) { @@ -39,6 +43,7 @@ public async ValueTask Stop( Task[] readers, Func? finalize = null ) { + // First, before anything cancels cts: completion is the only thing that wakes an idle Read loop. channel.Writer.TryComplete(); var incompleteReaders = readers.Where(r => !r.IsCompleted).ToArray(); diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/ConcurrentChannelWorkerTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/ConcurrentChannelWorkerTests.cs new file mode 100644 index 00000000..cfa094dd --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/ConcurrentChannelWorkerTests.cs @@ -0,0 +1,154 @@ +// Copyright (C) Eventuous HQ OÜ. All rights reserved +// Licensed under the Apache License, Version 2.0. + +using System.Collections.Concurrent; +using System.Threading.Channels; +using Eventuous.Subscriptions.Channels; +using Shouldly; + +namespace Eventuous.Tests.Subscriptions; + +/// +/// Pins the shutdown contract of the reader loop behind AsyncHandlingFilter: an idle reader is woken +/// by the channel completing rather than by the worker's token, so each of these would hang or lose +/// elements if the loop stopped honouring one of the two. +/// +public class ConcurrentChannelWorkerTests { + static readonly TimeSpan Prompt = TimeSpan.FromSeconds(2); + + [Test] + [Arguments(1)] + [Arguments(4)] + public async Task Idle_worker_stops_promptly_on_dispose(int readers, CancellationToken ct) { + var worker = new ConcurrentChannelWorker(CreateChannel(readers), (_, _) => default, readers); + + // Long enough for every reader to be parked on the empty channel. + await Task.Delay(100, ct); + + // Well inside the ten seconds a reader that missed the completion would take to be cancelled. + await worker.DisposeAsync().AsTask().WaitAsync(Prompt, ct); + } + + [Test] + public async Task Idle_worker_stops_promptly_after_processing(CancellationToken ct) { + var processed = 0; + + var worker = new ConcurrentChannelWorker( + CreateChannel(1), + (_, _) => { + Interlocked.Increment(ref processed); + + return default; + }, + 1 + ); + + for (var i = 0; i < 3; i++) { + (await worker.Write(i, ct)).ShouldBeTrue(); + var expected = i + 1; + (await Wait.Until(() => Volatile.Read(ref processed) == expected, Prompt)).ShouldBeTrue("each element should be picked up from an idle channel"); + } + + await worker.DisposeAsync().AsTask().WaitAsync(Prompt, ct); + } + + [Test] + public async Task Queued_elements_are_drained_on_graceful_dispose(CancellationToken ct) { + const int count = 5; + + var gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var processed = new ConcurrentQueue(); + var cancelled = 0; + + var worker = new ConcurrentChannelWorker( + CreateChannel(1), + async (element, token) => { + await gate.Task; + if (token.IsCancellationRequested) Interlocked.Increment(ref cancelled); + processed.Enqueue(element); + }, + 1 + ); + + for (var i = 0; i < count; i++) (await worker.Write(i, ct)).ShouldBeTrue(); + + var disposing = worker.DisposeAsync().AsTask(); + + await Task.Delay(100, ct); + disposing.IsCompleted.ShouldBeFalse("dispose must wait for the elements queued before it"); + (await worker.Write(count, ct)).ShouldBeFalse("a stopping worker takes nothing new"); + + gate.SetResult(); + await disposing.WaitAsync(Prompt, ct); + + processed.ShouldBe(Enumerable.Range(0, count)); + cancelled.ShouldBe(0, "a drain inside the deadline runs on a live token"); + } + + [Test] + public async Task Elements_from_several_writers_are_all_processed_by_several_readers(CancellationToken ct) { + const int readers = 4; + const int writers = 4; + const int perWriter = 500; + + var processed = new ConcurrentBag(); + + var worker = new ConcurrentChannelWorker( + CreateChannel(readers), + async (element, _) => { + // Yields now and then, so readers go back to the channel both synchronously and not. + if (element % 7 == 0) await Task.Yield(); + processed.Add(element); + }, + readers + ); + + await Task.WhenAll( + Enumerable.Range(0, writers) + .Select(w => Task.Run( + async () => { + for (var i = 0; i < perWriter; i++) (await worker.Write(w * perWriter + i, ct)).ShouldBeTrue(); + }, + ct + ) + ) + ); + + await worker.DisposeAsync().AsTask().WaitAsync(Prompt, ct); + + processed.Order().ShouldBe(Enumerable.Range(0, writers * perWriter)); + } + + /// + /// The drain is bounded: once the deadline cancels the worker's token, what's still queued is left + /// unprocessed (for the handling filter, never acknowledged and so redelivered) rather than handed + /// to the processor with a dead token, and dispose still completes. + /// + [Test] + public async Task Elements_still_queued_when_the_drain_deadline_passes_are_left_unprocessed(CancellationToken ct) { + var started = new ConcurrentQueue(); + + var worker = new ConcurrentChannelWorker( + CreateChannel(1), + async (element, token) => { + started.Enqueue(element); + + // Returns normally once cancelled, so it's the loop that has to notice the token. + try { + await Task.Delay(Timeout.Infinite, token); + } catch (OperationCanceledException) { } + }, + 1 + ); + + for (var i = 0; i < 3; i++) (await worker.Write(i, ct)).ShouldBeTrue(); + + // Ten seconds of drain deadline, plus room. + await worker.DisposeAsync().AsTask().WaitAsync(TimeSpan.FromSeconds(20), ct); + + started.ShouldBe([0]); + } + + static Channel CreateChannel(int readers) + => Channel.CreateBounded(new BoundedChannelOptions(10 * readers) { SingleReader = readers == 1 }); +} From ba9ad3f7a54d2d378e124ca3c77a8b7d2f781de8 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 16:33:44 +0200 Subject: [PATCH 4/7] perf(subscriptions): find the checkpoint gap and trim the sequence without LINQ FirstBeforeGap walks the set once with the struct enumerator instead of Zip/Skip/FirstOrDefault, and the commit handler trims the committed prefix with a RemoveUpTo helper instead of RemoveWhere with a closure. Behaviour is unchanged, including the gap diagnostic. Tests pin the gap cases and the prefix removal; the gap tests also pass on the old implementation. The harness shows payload-less events at about 409 bytes/event against a baseline of about 460; the other scenarios are within noise. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Checkpoints/CheckpointCommitHandler.cs | 4 +- .../Checkpoints/CommitPositionSequence.cs | 29 +++++-- .../SequenceTests.cs | 85 +++++++++++++++++++ 3 files changed, 110 insertions(+), 8 deletions(-) diff --git a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CheckpointCommitHandler.cs b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CheckpointCommitHandler.cs index a95e52e5..03a82c0d 100644 --- a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CheckpointCommitHandler.cs +++ b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CheckpointCommitHandler.cs @@ -143,10 +143,10 @@ async Task CommitInternal(CommitPosition position, bool force, CancellationToken position.LogContext?.CommittingPosition(position); await _commitCheckpoint(new(_subscriptionId, position.Position), force, cancellationToken).NoContext(); _lastCommit = position; - _positions.RemoveWhere(x => x.Sequence <= position.Sequence); + _positions.RemoveUpTo(position.Sequence); } catch (OperationCanceledException) { await _commitCheckpoint(new(_subscriptionId, position.Position), true, default).NoContext(); - _positions.RemoveWhere(x => x.Sequence <= position.Sequence); + _positions.RemoveUpTo(position.Sequence); } catch (Exception e) { position.LogContext?.UnableToCommitPosition(position, e); } diff --git a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs index 43ef2ff3..0fc8bedc 100644 --- a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs +++ b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs @@ -19,14 +19,31 @@ public CommitPosition FirstBeforeGap() }; CommitPosition Get() { - var result = this - .Zip(this.Skip(1), (position1, position2) => (position1, position2)) - .FirstOrDefault(tup => tup.position1.Sequence + 1 != tup.position2.Sequence); + var previous = default(CommitPosition); + bool first = true; - if (result == default) return Max; + // The set is ordered by Sequence; the struct enumerator avoids the boxed IEnumerable one + foreach (var current in this) { + if (!first && previous.Sequence + 1 != current.Sequence) { + SubscriptionsEventSource.Log.CheckpointGapDetected(previous, current); - SubscriptionsEventSource.Log.CheckpointGapDetected(result.position1, result.position2); - return result.position1; + return previous; + } + + previous = current; + first = false; + } + + return Max; + } + + /// + /// Removes all positions with a sequence up to and including the given one. The set is ordered + /// by sequence, so these are always a prefix; no predicate (and no closure) needed. + /// + internal void RemoveUpTo(ulong sequence) { + // Min is default on an empty set, so the Count check is what stops the loop + while (Count > 0 && Min.Sequence <= sequence) Remove(Min); } class PositionsComparer : IComparer { diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs index cc43beba..d3d57d53 100644 --- a/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs +++ b/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs @@ -61,6 +61,91 @@ public void ShouldWorkForNormalCase() { first.ShouldBe(new(9, 9, timestamp)); } + [Test] + public void ShouldReturnMaxWhenNoGap() { + var sequence = Sequence(3, 4, 5); + sequence.FirstBeforeGap().Sequence.ShouldBe(5UL); + } + + [Test] + public void ShouldReturnEmptyForEmptySet() { + new CommitPositionSequence().FirstBeforeGap().ShouldBe(CommitPosition.None); + } + + [Test] + public void ShouldFindGapAtTheStart() { + var sequence = Sequence(1, 3, 4); + sequence.FirstBeforeGap().Sequence.ShouldBe(1UL); + } + + [Test] + public void ShouldFindGapInTheMiddle() { + var sequence = Sequence(0, 1, 2, 4, 5); + sequence.FirstBeforeGap().Sequence.ShouldBe(2UL); + } + + [Test] + public void ShouldReturnFirstOfTwoGaps() { + var sequence = Sequence(0, 1, 3, 4, 6, 7); + sequence.FirstBeforeGap().Sequence.ShouldBe(1UL); + } + + [Test] + public void ShouldWorkForTwoWithoutGap() { + var sequence = Sequence(7, 8); + sequence.FirstBeforeGap().Sequence.ShouldBe(8UL); + } + + [Test] + public void ShouldWorkForTwoWithGap() { + var sequence = Sequence(7, 9); + sequence.FirstBeforeGap().Sequence.ShouldBe(7UL); + } + + [Test] + public void RemoveUpTo_removes_exactly_the_prefix() { + var sequence = Sequence(2, 3, 4, 6, 7); + sequence.RemoveUpTo(4); + sequence.Select(x => x.Sequence).ShouldBe([6UL, 7UL]); + } + + [Test] + public void RemoveUpTo_removes_across_a_gap() { + var sequence = Sequence(2, 3, 6, 7); + sequence.RemoveUpTo(5); + sequence.Select(x => x.Sequence).ShouldBe([6UL, 7UL]); + } + + [Test] + public void RemoveUpTo_is_a_noop_on_empty_set() { + var sequence = new CommitPositionSequence(); + sequence.RemoveUpTo(10); + sequence.Count.ShouldBe(0); + } + + [Test] + public void RemoveUpTo_below_minimum_removes_nothing() { + var sequence = Sequence(5, 6, 7); + sequence.RemoveUpTo(4); + sequence.Count.ShouldBe(3); + } + + [Test] + public void RemoveUpTo_above_maximum_empties_the_set() { + var sequence = Sequence(5, 6, 7); + sequence.RemoveUpTo(100); + sequence.Count.ShouldBe(0); + } + + static CommitPositionSequence Sequence(params ulong[] sequences) { + var result = new CommitPositionSequence(); + + // Positions differ from sequences so a mix-up between the two shows + foreach (var seq in sequences) result.Add(new(seq + 100, seq, DateTime.Now)); + + return result; + } + public static IEnumerable> TestData() { var timestamp = DateTime.Now; From 36ee3522c38895a93983f50debed796fee40d184 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 16:45:49 +0200 Subject: [PATCH 5/7] perf(subscriptions): set the logger context once per loop instead of per event Logger.Current skips the AsyncLocal write when the value is unchanged, but on the consume path the write sat inside async methods (DelayedConsume, HandleInternal). An async method's ExecutionContext is restored when it returns, so the calling loop never saw the value and the next event wrote it again: a new ExecutionContext and value map, 72 bytes, per write. The write now lands in the loop's own context. AsyncHandlingFilter sets it in a non-async shim in front of the unchanged async body, so the channel reader keeps it and writes again only when the next message carries another LogContext (a pipe shared by two subscriptions). The KurrentDB $all pump sets it before its loop, and the stream subscription's callback sets it in a non-async shim, which puts it in the client's read loop. SQL and Redis already set it from a synchronous method inside their polling loops. Measured with the allocation harness, bytes per event: typed async 3331-3360 -> 3245-3274 from the filter alone, and 3208-3216 once the calling loop holds the context as the $all pump now does; payload-less 415 -> 346 on the same condition. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Filters/AsyncHandlingFilter.cs | 13 +++- .../AsyncHandlingFilterLogContextTests.cs | 77 +++++++++++++++++++ .../Subscriptions/AllStreamSubscription.cs | 5 ++ .../Subscriptions/StreamSubscription.cs | 12 ++- 4 files changed, 103 insertions(+), 4 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/AsyncHandlingFilterLogContextTests.cs diff --git a/src/Core/src/Eventuous.Subscriptions/Filters/AsyncHandlingFilter.cs b/src/Core/src/Eventuous.Subscriptions/Filters/AsyncHandlingFilter.cs index bc8ed8fe..91d40022 100644 --- a/src/Core/src/Eventuous.Subscriptions/Filters/AsyncHandlingFilter.cs +++ b/src/Core/src/Eventuous.Subscriptions/Filters/AsyncHandlingFilter.cs @@ -24,8 +24,17 @@ public AsyncHandlingFilter(uint concurrencyLimit, uint bufferSize = 10) { _worker = new(Channel.CreateBounded(options), DelayedConsume, (int)concurrencyLimit); } + // Not async on purpose: an async method's ExecutionContext is restored when it returns, so a logger + // context set inside one never reaches the reader loop, and every message pays for a new context. + // Set here it stays with the loop, and is written again only when the next message brings another. + static ValueTask DelayedConsume(WorkerTask workerTask, CancellationToken ct) { + Logger.Current = workerTask.Context.LogContext; + + return Consume(workerTask, ct); + } + // ReSharper disable once CognitiveComplexity - static async ValueTask DelayedConsume(WorkerTask workerTask, CancellationToken ct) { + static async ValueTask Consume(WorkerTask workerTask, CancellationToken ct) { var ctx = workerTask.Context; using var activity = ctx.Items.GetItem(ContextItemKeys.Activity)?.Start(); @@ -48,8 +57,6 @@ static async ValueTask DelayedConsume(WorkerTask workerTask, CancellationToken c ctx.CancellationToken = cts.Token; } - Logger.Current = ctx.LogContext; - try { try { await workerTask.Filter.Value.Send(ctx, workerTask.Filter.Next).NoContext(); diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/AsyncHandlingFilterLogContextTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/AsyncHandlingFilterLogContextTests.cs new file mode 100644 index 00000000..a532467a --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/AsyncHandlingFilterLogContextTests.cs @@ -0,0 +1,77 @@ +// Copyright (C) Eventuous HQ OÜ. All rights reserved +// Licensed under the Apache License, Version 2.0. + +using System.Collections.Concurrent; +using Eventuous.Subscriptions; +using Eventuous.Subscriptions.Context; +using Eventuous.Subscriptions.Filters; +using Eventuous.Subscriptions.Logging; +using Shouldly; + +namespace Eventuous.Tests.Subscriptions; + +/// +/// The filter's reader keeps the logger context between messages instead of setting it anew for each, so +/// what's pinned here is that a message still never runs under the context of the one before it — the case +/// of a pipe shared by two subscriptions, whose messages interleave on one reader. +/// +public class AsyncHandlingFilterLogContextTests { + [Test] + public async Task Handler_sees_the_log_context_of_its_own_message() { + var handler = new RecordingHandler(); + + await using var pipe = new ConsumePipe().AddDefaultConsumer(handler).AddFilterFirst(new AsyncHandlingFilter(1)); + + var first = Logger.CreateContext("first", null); + var second = Logger.CreateContext("second", null); + + // Repeats and switches both: a context kept from the previous message must be neither lost nor stale. + LogContext[] sent = [first, first, second, first, second, second]; + + for (var i = 0; i < sent.Length; i++) { + await pipe.Send(CreateContext(i, sent[i])); + } + + var handled = await Wait.Until(() => handler.Seen.Count == sent.Length, TimeSpan.FromSeconds(5)); + + handled.ShouldBeTrue($"{handler.Seen.Count} of {sent.Length} messages handled"); + + // One reader, so the order is the order sent. + var seen = handler.Seen.ToArray(); + + for (var i = 0; i < sent.Length; i++) { + seen[i].Message.ShouldBeSameAs(sent[i], $"message {i} arrived out of order"); + seen[i].Current.ShouldBeSameAs(sent[i], $"message {i} was handled under another message's log context"); + } + } + + static AsyncConsumeContext CreateContext(int position, LogContext logContext) { + var context = new MessageConsumeContext( + Guid.NewGuid().ToString(), + "TestEvent", + "application/json", + "test-stream", + (ulong)position, + (ulong)position, + (ulong)position, + (ulong)position, + DateTime.UtcNow, + new { Number = position }, + new(), + logContext.SubscriptionId, + CancellationToken.None + ) { LogContext = logContext }; + + return new(context, _ => default, (_, _) => default); + } + + sealed class RecordingHandler : BaseEventHandler { + public ConcurrentQueue<(LogContext Message, LogContext Current)> Seen { get; } = new(); + + public override ValueTask HandleEvent(IMessageConsumeContext context) { + Seen.Enqueue((context.LogContext, Logger.Current)); + + return new(EventHandlingStatus.Success); + } + } +} diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs index c089fc69..34e7dff1 100644 --- a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs +++ b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs @@ -6,6 +6,7 @@ using Eventuous.Subscriptions.Context; using Eventuous.Subscriptions.Diagnostics; using Eventuous.Subscriptions.Filters; +using Eventuous.Subscriptions.Logging; using Eventuous.Tools; namespace Eventuous.KurrentDB.Subscriptions; @@ -153,6 +154,10 @@ async Task Consume( ulong? headPosition, ulong? lastScannedPosition ) { + // Once, for the whole loop: set by HandleInternal instead it would be lost when that returns, and + // written again for every event. + Logger.Current = Log; + while (await messages.MoveNextAsync().NoContext()) { // Falling behind re-enters catch-up mode: re-read the head as the new commit candidate, since // reading it later from the caught-up message could race matches still in flight. Kept outside diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs index 8dfd106e..5fd40357 100644 --- a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs +++ b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs @@ -5,6 +5,7 @@ using Eventuous.Subscriptions.Context; using Eventuous.Subscriptions.Diagnostics; using Eventuous.Subscriptions.Filters; +using Eventuous.Subscriptions.Logging; using Eventuous.Tools; namespace Eventuous.KurrentDB.Subscriptions; @@ -113,7 +114,16 @@ protected override async ValueTask Connect(SubscriptionRun run) { _ => FromStream.After(StreamPosition.FromInt64((long)position)) }; - async Task HandleEvent(ResolvedEvent re, CancellationToken ct) { + // Not async on purpose: the client calls this from its own read loop, so a logger context set here + // stays with that loop, where one set inside an async method would be lost when it returns and + // written again for every event. + Task HandleEvent(ResolvedEvent re, CancellationToken ct) { + Logger.Current = Log; + + return HandleResolvedEvent(re, ct); + } + + async Task HandleResolvedEvent(ResolvedEvent re, CancellationToken ct) { // Despite ResolvedEvent.Event being not marked as nullable, it returns null for deleted events // ReSharper disable once ConditionIsAlwaysTrueOrFalse // ReSharper disable once ConditionIsAlwaysTrueOrFalseAccordingToNullableAPIContract From 37937684c2982444367b9030619dd2f444f927e5 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 16:56:19 +0200 Subject: [PATCH 6/7] perf(subscriptions): return forwarded tasks instead of awaiting them When a handler really suspends, every async method above it boxes its own state machine. Two layers on the consume path were async only to forward the task of the next one: EventHandler.HandleEvent awaited the typed handler and returned its status unchanged. It now returns the handler's task, or the cached Ignored one for a message type without a handler. TracingFilter.Send has nothing to do after the next filter when there is no activity: none started here because nobody listens or the span is sampled out, and no current one to reuse. It now returns the next filter's task in that case, and the part that carries an activity moved to an async helper. The activity is created before the split and started inside the helper, so it is current for exactly the same scope as before. With the dummy listener in place an activity always exists, so this only pays off once OpenTelemetry tracing is wired and drops the span. Allocation harness, bytes per event with a handler that suspends: 3280-3350 before, 3125-3150 after. In a copy of the harness with the dummy listener removed and a suspending handler, the tracing pipe goes from 2547 to 2340-2350. The rows with a synchronous handler and the traced row with a listener are unchanged, within noise. An exception thrown before the first await (a context without a message) now leaves HandleEvent directly instead of in the returned task. The callers await inside their try blocks, so the outcome is the same. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Filters/TracingFilter.cs | 24 +++++++++--- .../Handlers/EventHandler.cs | 5 ++- .../DefaultConsumerTests.cs | 39 +++++++++++++++++++ 3 files changed, 61 insertions(+), 7 deletions(-) diff --git a/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs b/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs index 66802450..d906d943 100644 --- a/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs +++ b/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs @@ -19,23 +19,37 @@ public TracingFilter(string consumerName) { _defaultTags = [.. tags, .. EventuousDiagnostics.Tags]; } - protected override async ValueTask Send(IMessageConsumeContext context, LinkedListNode? next) { - if (context.Message == null || next == null) return; + protected override ValueTask Send(IMessageConsumeContext context, LinkedListNode? next) { + if (context.Message == null || next == null) return default; // The subscription's own activity is reused, not owned: disposing it would stop it before the // subscription is done with it, so only an activity started here gets disposed. var reuseCurrent = Activity.Current?.Context == context.ParentContext; - using var started = reuseCurrent + var created = reuseCurrent ? null - : SubscriptionActivity.Start( + : SubscriptionActivity.Create( $"{Constants.Components.Consumer}.{context.SubscriptionId}/{context.MessageType}", ActivityKind.Consumer, context, _defaultTags ); - var activity = reuseCurrent ? Activity.Current : started; + var reused = reuseCurrent ? Activity.Current : null; + + // Nobody listening, or sampled out: nothing to record after the next filter, so its task is returned + // as is rather than awaited. + if (reused == null && created == null) return next.Value.Send(context, next.Next); + + return SendTraced(context, next, reused, created); + } + + static async ValueTask SendTraced(IMessageConsumeContext context, LinkedListNode next, Activity? reused, Activity? created) { + // Started here, not by the caller: this method's execution context is restored when it returns, so the + // started activity doesn't stay current for whoever called the filter. + using var started = created?.Start(); + + var activity = reused ?? started; if (activity?.IsAllDataRequested == true && context is AsyncConsumeContext asyncConsumeContext) { activity.SetContextTags(context)?.SetTag(TelemetryTags.Eventuous.Partition, asyncConsumeContext.PartitionId); diff --git a/src/Core/src/Eventuous.Subscriptions/Handlers/EventHandler.cs b/src/Core/src/Eventuous.Subscriptions/Handlers/EventHandler.cs index 7f2554b5..4276ed31 100644 --- a/src/Core/src/Eventuous.Subscriptions/Handlers/EventHandler.cs +++ b/src/Core/src/Eventuous.Subscriptions/Handlers/EventHandler.cs @@ -57,8 +57,9 @@ ValueTask NoHandler() { } } - public override async ValueTask HandleEvent(IMessageConsumeContext context) - => !_handlersMap.TryGetValue(context.Message!.GetType(), out var handler) ? EventHandlingStatus.Ignored : await handler(context).NoContext(); + // Not async: there is nothing to do after the handler, so its task is returned as is rather than awaited + public override ValueTask HandleEvent(IMessageConsumeContext context) + => !_handlersMap.TryGetValue(context.Message!.GetType(), out var handler) ? Ignored : handler(context); public override string ToString() { var sb = new StringBuilder(); diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/DefaultConsumerTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/DefaultConsumerTests.cs index b802521c..3fde2301 100644 --- a/src/Core/test/Eventuous.Tests.Subscriptions/DefaultConsumerTests.cs +++ b/src/Core/test/Eventuous.Tests.Subscriptions/DefaultConsumerTests.cs @@ -3,6 +3,8 @@ using Eventuous.Subscriptions.Consumers; using Eventuous.Subscriptions.Context; using Eventuous.TestHelpers.TUnit; +using Eventuous.TestHelpers.TUnit.Logging; +using EventHandler = Eventuous.Subscriptions.EventHandler; namespace Eventuous.Tests.Subscriptions; @@ -21,7 +23,44 @@ public async Task ShouldFailWhenHandlerNacks() { await Assert.That(ctx.HandlingResults.GetFailureStatus()).IsEqualTo(EventHandlingStatus.Failure); } + [Test] + public async Task ShouldNackWhenTypedHandlerThrowsSynchronously() { + var error = new InvalidOperationException("handler failed"); + var consumer = new DefaultConsumer([new TypedHandler(error)]); + var ctx = CreateContext(new Handled()); + + await consumer.Consume(ctx); + + await Assert.That(ctx.HasFailed()).IsTrue(); + await Assert.That(ctx.HandlingResults.GetException()).IsSameReferenceAs(error); + } + + [Test] + public async Task ShouldIgnoreMessageWithoutTypedHandler() { + var consumer = new DefaultConsumer([new TypedHandler(new InvalidOperationException())]); + var ctx = CreateContext(new NotHandled()); + + await consumer.Consume(ctx); + + await Assert.That(ctx.WasIgnored()).IsTrue(); + await Assert.That(ctx.HasFailed()).IsFalse(); + } + + static MessageConsumeContext CreateContext(object message) + => new("id", "type", "application/json", "stream", 0, 0, 0, 0, DateTime.UtcNow, message, null, "test", CancellationToken.None) { + LogContext = new("test", new LoggerFactory().AddTUnit(LogLevel.Information)) + }; + public void Dispose() => _listener.Dispose(); + + record Handled; + + record NotHandled; + + class TypedHandler : EventHandler { + // Not async: the exception leaves the handler delegate before it returns a task + public TypedHandler(Exception error) => On(_ => throw error); + } } class FailingHandler : IEventHandler { From 6e97a3d8fcf0262de104ed5024b63b9d5969fe87 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Fri, 2 Oct 2026 20:19:12 +0200 Subject: [PATCH 7/7] refactor(subscriptions): compare commit positions with CompareTo The hand-written comparison returned the same -1, 0 or 1 that ulong.CompareTo gives for the two sequences. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Checkpoints/CommitPositionSequence.cs | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs index 0fc8bedc..34835dbf 100644 --- a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs +++ b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs @@ -48,10 +48,6 @@ internal void RemoveUpTo(ulong sequence) { class PositionsComparer : IComparer { [MethodImpl(MethodImplOptions.AggressiveInlining)] - public int Compare(CommitPosition x, CommitPosition y) { - if (x.Sequence == y.Sequence) return 0; - - return x.Sequence > y.Sequence ? 1 : -1; - } + public int Compare(CommitPosition x, CommitPosition y) => x.Sequence.CompareTo(y.Sequence); } }