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/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..34835dbf 100644 --- a/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs +++ b/src/Core/src/Eventuous.Subscriptions/Checkpoints/CommitPositionSequence.cs @@ -19,22 +19,35 @@ 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 { [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); } } 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/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/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/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/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 }); +} 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 { 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; + } + } +} 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; 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/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(); } 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