From 4f2c21b1d4195a0d373a2ba872c60bf2f8c19e16 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 00:46:18 +0200 Subject: [PATCH 1/7] fix(subscriptions): fix consume-path bugs and trim per-message allocations - MurmurHash3 hashes the partition key through a span byte view instead of reading an unpinned string; AllowUnsafeBlocks is no longer needed - MessageConsumeContextConverter caches are thread-safe - TracingFilter disposes only the activities it created, so trace-flag changes made after handling take effect, and it keeps the error status a failed handler set instead of overwriting it with OK - ConsumePipe validates the second filter's context type when the pipe is composed, not on every message - DefaultConsumer logging scope is a KeyValuePair array - MessageConsumeContext.Items is allocated on first use and published atomically, so racing first accesses share one bag - Checkpointed runs build their ack and nack delegates once per run Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Consumers/DefaultConsumer.cs | 8 +- .../MessageConsumeContextConverter.cs | 43 ++++---- .../Context/MessageConsumeContext.cs | 7 +- .../EventSubscriptionWithCheckpoint.cs | 27 ++++- .../Eventuous.Subscriptions.csproj | 1 - .../Filters/ConsumePipe.cs | 2 +- .../Filters/Partitioning/MurmurHash3.cs | 17 +-- .../Filters/TracingFilter.cs | 17 ++- .../ConsumePipeTests.cs | 14 +++ .../ContextItemsTests.cs | 39 +++++++ .../MessageConsumeContextConverterTests.cs | 90 ++++++++++++++++ .../MurmurHash3Tests.cs | 23 ++++ .../TracingFilterTests.cs | 101 ++++++++++++++++++ 13 files changed, 334 insertions(+), 55 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/ContextItemsTests.cs create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/MessageConsumeContextConverterTests.cs create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/MurmurHash3Tests.cs create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/TracingFilterTests.cs diff --git a/src/Core/src/Eventuous.Subscriptions/Consumers/DefaultConsumer.cs b/src/Core/src/Eventuous.Subscriptions/Consumers/DefaultConsumer.cs index 17339fbfd..e0de62102 100644 --- a/src/Core/src/Eventuous.Subscriptions/Consumers/DefaultConsumer.cs +++ b/src/Core/src/Eventuous.Subscriptions/Consumers/DefaultConsumer.cs @@ -8,10 +8,10 @@ namespace Eventuous.Subscriptions.Consumers; // ReSharper disable once ParameterTypeCanBeEnumerable.Local public class DefaultConsumer(IEventHandler[] eventHandlers) : IMessageConsumer { public async ValueTask Consume(IMessageConsumeContext context) { - var scope = new Dictionary { - {"SubscriptionId", context.SubscriptionId}, - {"Stream", context.Stream}, - {"MessageType", context.MessageType} + var scope = new KeyValuePair[] { + new("SubscriptionId", context.SubscriptionId), + new("Stream", context.Stream), + new("MessageType", context.MessageType) }; using var _ = context.LogContext.Logger.BeginScope(scope); diff --git a/src/Core/src/Eventuous.Subscriptions/Consumers/MessageConsumeContextConverter.cs b/src/Core/src/Eventuous.Subscriptions/Consumers/MessageConsumeContextConverter.cs index f6c04f7e2..95b6fb92a 100644 --- a/src/Core/src/Eventuous.Subscriptions/Consumers/MessageConsumeContextConverter.cs +++ b/src/Core/src/Eventuous.Subscriptions/Consumers/MessageConsumeContextConverter.cs @@ -1,6 +1,7 @@ // Copyright (C) Eventuous HQ OÜ. All rights reserved // Licensed under the Apache License, Version 2.0. +using System.Collections.Concurrent; using System.Linq.Expressions; using System.Runtime.CompilerServices; using Eventuous.Subscriptions.Logging; @@ -11,9 +12,6 @@ namespace Eventuous.Subscriptions.Consumers; using System.Diagnostics.CodeAnalysis; using Context; -#if NET8_0 -using Lock = object; -#endif /// /// Converts non-generic IMessageConsumeContext to a typed IMessageConsumeContext. @@ -22,9 +20,15 @@ namespace Eventuous.Subscriptions.Consumers; /// via , which will be attempted before using reflection. /// public static class MessageConsumeContextConverter { - internal static readonly Dictionary ConversionCache = new(); - internal static readonly List RegisteredConverters = []; - static readonly Lock CacheLock = new(); + internal static readonly ConcurrentDictionary ConversionCache = new(); + + /// + /// Copy-on-write: swaps in a new array, so a conversion running while a module + /// initializer registers a converter reads a complete snapshot, never a list being resized. + /// + static volatile ContextConversion[] registeredConverters = []; + + internal static IReadOnlyList RegisteredConverters => registeredConverters; /// /// Registers a converter function to try before the fallback reflection-based conversion. @@ -34,38 +38,27 @@ public static class MessageConsumeContextConverter { /// A function that returns a typed context or null if not handled. [MethodImpl(MethodImplOptions.Synchronized)] public static void Register(ContextConversion converter) { - RegisteredConverters.Add(converter); + registeredConverters = [.. registeredConverters, converter]; } public static IMessageConsumeContext ConvertToGeneric(this IMessageConsumeContext context, InternalLogger? log = null) { var messageType = context.Message!.GetType(); + var converters = registeredConverters; - // ReSharper disable InconsistentlySynchronizedField - if (RegisteredConverters.Count > 0) { - for (var i = 0; i < RegisteredConverters.Count; i++) { - var converter = RegisteredConverters[i]; - - if (converter(context) is { } typedContext) { - return typedContext; - } + for (var i = 0; i < converters.Length; i++) { + if (converters[i](context) is { } typedContext) { + return typedContext; } } - // ReSharper restore InconsistentlySynchronizedField - // ReSharper disable once InconsistentlySynchronizedField if (!ConversionCache.TryGetValue(messageType, out var conversion)) { log?.Log("Static context conversion not found for message type {MessageType}, using reflections. Consider opening a GitHub issue to help improving the generator", messageType); - lock (CacheLock) { - if (!ConversionCache.TryGetValue(messageType, out conversion)) { - conversion = CreateConversionFunction(messageType); - - ConversionCache[messageType] = conversion; - } - } + // Racing callers may both compile a conversion; only one gets cached, and either works + conversion = ConversionCache.GetOrAdd(messageType, CreateConversionFunction); } - return (IMessageConsumeContext)conversion!(context); + return (IMessageConsumeContext)conversion(context); } [UnconditionalSuppressMessage("AOT", "IL3050", Justification = "This should not be used because all the conversions should be pre-generated")] diff --git a/src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs b/src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs index 23fd9c6c2..b0e0779fb 100644 --- a/src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs +++ b/src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs @@ -23,6 +23,8 @@ public class MessageConsumeContext( CancellationToken cancellationToken ) : IMessageConsumeContext { + ContextItems? _items; + /// public string MessageId { get; } = eventId; /// @@ -44,7 +46,7 @@ CancellationToken cancellationToken /// public object? Message { get; } = message; /// - public ContextItems Items { get; } = new(); + public ContextItems Items => _items ?? CreateItems(); /// public ActivityContext? ParentContext { get; set; } /// @@ -57,6 +59,9 @@ CancellationToken cancellationToken public string SubscriptionId { get; } = subscriptionId; /// public LogContext LogContext { get; set; } = Logger.Current; + + // Published atomically: two first accesses racing must share one bag, or an item added to the losing one is lost + ContextItems CreateItems() => Interlocked.CompareExchange(ref _items, new(), null) ?? _items; } public class MessageConsumeContext(IMessageConsumeContext innerContext) : WrappedConsumeContext(innerContext), IMessageConsumeContext diff --git a/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs b/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs index 70b8fcd5f..f22397fd1 100644 --- a/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs +++ b/src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs @@ -52,8 +52,25 @@ EventPosition GetPositionFromContext(IMessageConsumeContext context) /// A run carrying this attempt's own commit handler, so an acknowledgement reaches the handler that /// dispatched it, and the base class never has to know checkpoints exist. /// - sealed class CheckpointedRun(CancellationToken lifetime, CheckpointCommitHandler checkpoint) : SubscriptionRun(lifetime) { - internal CheckpointCommitHandler Checkpoint { get; } = checkpoint; + sealed class CheckpointedRun : SubscriptionRun { + // Not a primary constructor: the ack and nack delegates close over this run, which an initializer + // can't reference. + public CheckpointedRun(CancellationToken lifetime, CheckpointCommitHandler checkpoint, EventSubscriptionWithCheckpoint subscription) + : base(lifetime) { + Checkpoint = checkpoint; + AckMessage = ctx => subscription.Ack(this, ctx); + NackMessage = (ctx, exception) => subscription.NackOnAsyncWorker(this, ctx, exception); + } + + internal CheckpointCommitHandler Checkpoint { get; } + + /// + /// Created once per run rather than per message, and bound to this run for the reason given on + /// . + /// + internal Acknowledge AckMessage { get; } + + internal Fail NackMessage { get; } } /// @@ -69,7 +86,8 @@ protected sealed override SubscriptionRun CreateRun(CancellationToken lifetime) TimeSpan.FromMilliseconds(Options.CheckpointCommitDelayMs), Options.CheckpointCommitBatchSize, LoggerFactory - ) + ), + this ); // Registered first, before Connect, so it's the first OnDisconnect registration — and since release @@ -94,7 +112,8 @@ protected async ValueTask HandleInternal(SubscriptionRun run, IMessageConsumeCon try { Logger.Current = Log; - var ctx = new AsyncConsumeContext(context, c => Ack(run, c), (c, e) => NackOnAsyncWorker(run, c, e)); + var checkpointedRun = (CheckpointedRun)run; + 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/Eventuous.Subscriptions.csproj b/src/Core/src/Eventuous.Subscriptions/Eventuous.Subscriptions.csproj index 8142706ad..946d0c063 100644 --- a/src/Core/src/Eventuous.Subscriptions/Eventuous.Subscriptions.csproj +++ b/src/Core/src/Eventuous.Subscriptions/Eventuous.Subscriptions.csproj @@ -1,6 +1,5 @@ - true true diff --git a/src/Core/src/Eventuous.Subscriptions/Filters/ConsumePipe.cs b/src/Core/src/Eventuous.Subscriptions/Filters/ConsumePipe.cs index fbe54a5a4..1f4ec7f2c 100644 --- a/src/Core/src/Eventuous.Subscriptions/Filters/ConsumePipe.cs +++ b/src/Core/src/Eventuous.Subscriptions/Filters/ConsumePipe.cs @@ -53,7 +53,7 @@ public ConsumePipe AddFilterLast(IConsumeFilter filter) // Deny adding filter of the same type twice if (_filters.Any(x => x.GetType() == filter.GetType())) throw new DuplicateFilterException(filter); - if (_filters.Count > 1 && !typeof(TIn).IsAssignableFrom(_filters.Last().Produces)) { + if (_filters.Count > 0 && !typeof(TIn).IsAssignableFrom(_filters.Last().Produces)) { throw new InvalidContextTypeException(_filters.Last().Produces, typeof(TIn)); } diff --git a/src/Core/src/Eventuous.Subscriptions/Filters/Partitioning/MurmurHash3.cs b/src/Core/src/Eventuous.Subscriptions/Filters/Partitioning/MurmurHash3.cs index 3e3e830ff..f922e4d0d 100644 --- a/src/Core/src/Eventuous.Subscriptions/Filters/Partitioning/MurmurHash3.cs +++ b/src/Core/src/Eventuous.Subscriptions/Filters/Partitioning/MurmurHash3.cs @@ -2,6 +2,7 @@ // Licensed under the Apache License, Version 2.0. using System.Runtime.CompilerServices; +using System.Runtime.InteropServices; namespace Eventuous.Subscriptions.Filters.Partitioning; @@ -11,21 +12,9 @@ static class MurmurHash3 { const uint Seed = 0xc58f1a7b; - const int CharSize = sizeof(char); - - static unsafe Span GetBytes(string data) - { - ArgumentNullException.ThrowIfNull(data); - if (data.Length == 0) return Span.Empty; - - fixed (char* p = data) - { - return new(p, data.Length * CharSize); - } - } - public static uint Hash(string partitionKey) { - var bytes = GetBytes(partitionKey); + ArgumentNullException.ThrowIfNull(partitionKey); + var bytes = MemoryMarshal.AsBytes(partitionKey.AsSpan()); var length = bytes.Length; var h1 = Seed; var tailLen = length & 3; diff --git a/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs b/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs index 32efda5ce..668024501 100644 --- a/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs +++ b/src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs @@ -22,14 +22,20 @@ public TracingFilter(string consumerName) { protected override async ValueTask Send(IMessageConsumeContext context, LinkedListNode? next) { if (context.Message == null || next == null) return; - using var activity = Activity.Current?.Context != context.ParentContext - ? SubscriptionActivity.Start( + // 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 + ? null + : SubscriptionActivity.Start( $"{Constants.Components.Consumer}.{context.SubscriptionId}/{context.MessageType}", ActivityKind.Consumer, context, _defaultTags - ) - : Activity.Current; + ); + + var activity = reuseCurrent ? Activity.Current : started; if (activity?.IsAllDataRequested == true && context is AsyncConsumeContext asyncConsumeContext) { activity.SetContextTags(context)?.SetTag(TelemetryTags.Eventuous.Partition, asyncConsumeContext.PartitionId); @@ -43,7 +49,8 @@ protected override async ValueTask Send(IMessageConsumeContext context, LinkedLi activity.ActivityTraceFlags = ActivityTraceFlags.None; } - activity.SetActivityStatus(ActivityStatus.Ok()); + // A handler failure is recorded with Nack, not thrown, and Nack has already set the error status + if (!context.HasFailed()) activity.SetActivityStatus(ActivityStatus.Ok()); } } catch (Exception e) { diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/ConsumePipeTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/ConsumePipeTests.cs index 7f6d0e28f..16ab91c12 100644 --- a/src/Core/test/Eventuous.Tests.Subscriptions/ConsumePipeTests.cs +++ b/src/Core/test/Eventuous.Tests.Subscriptions/ConsumePipeTests.cs @@ -53,6 +53,20 @@ public async Task ShouldMakeSecondDisposalWaitForTheFirst() { await Assert.That(filter.Disposals).IsEqualTo(1); } + [Test] + public async Task ShouldRejectSecondFilterThatCannotConsumeWhatTheFirstProduces() { + // The first filter produces IMessageConsumeContext, the second needs AsyncConsumeContext + var pipe = new ConsumePipe().AddFilterLast(new TestFilter(Key, "payload")); + + await Assert.That(() => pipe.AddFilterLast(new AsyncOnlyFilter())).Throws(); + await Assert.That(pipe.RegisteredFilters.Count()).IsEqualTo(1); + } + + class AsyncOnlyFilter : ConsumeFilter { + protected override ValueTask Send(AsyncConsumeContext context, LinkedListNode? next) + => next?.Value.Send(context, next.Next) ?? default; + } + /// /// Blocks inside until released, so a second disposal can be observed while the /// first one is still running. diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/ContextItemsTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/ContextItemsTests.cs new file mode 100644 index 000000000..ab96384dc --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/ContextItemsTests.cs @@ -0,0 +1,39 @@ +using Eventuous.Subscriptions.Context; + +namespace Eventuous.Tests.Subscriptions; + +public class ContextItemsTests { + [Test] + public async Task ShouldReturnSameItemsInstanceOnEveryAccess() { + var ctx = TestContext.CreateContext(); + + await Assert.That(ctx.Items).IsSameReferenceAs(ctx.Items); + } + + [Test] + public async Task ShouldShareOneItemsInstanceBetweenRacingFirstAccesses() { + for (var i = 0; i < 1000; i++) { + var ctx = TestContext.CreateContext(); + var barrier = new Barrier(2); + + var first = Task.Run(() => { barrier.SignalAndWait(); return ctx.Items; }); + var second = Task.Run(() => { barrier.SignalAndWait(); return ctx.Items; }); + + var items = await Task.WhenAll(first, second); + + await Assert.That(items[1]).IsSameReferenceAs(items[0]); + await Assert.That(ctx.Items).IsSameReferenceAs(items[0]); + } + } + + [Test] + public async Task ShouldKeepItemsAddedBeforeWrapping() { + var ctx = TestContext.CreateContext(); + ctx.Items.AddItem("key", "value"); + + var wrapped = new MessageConsumeContext(ctx); + + await Assert.That(wrapped.Items).IsSameReferenceAs(ctx.Items); + await Assert.That(wrapped.Items.GetItem("key")).IsEqualTo("value"); + } +} diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/MessageConsumeContextConverterTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/MessageConsumeContextConverterTests.cs new file mode 100644 index 000000000..53be852f6 --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/MessageConsumeContextConverterTests.cs @@ -0,0 +1,90 @@ +using Eventuous.Subscriptions.Consumers; +using Eventuous.Subscriptions.Context; + +namespace Eventuous.Tests.Subscriptions; + +public class MessageConsumeContextConverterTests { + /// + /// A regression smoke test, not proof of thread-safety. Many workers convert messages of types nobody converted + /// before while converters get registered mid-flight, the way a module initializer of a late-loaded assembly would. + /// Every conversion must still produce the right typed context and nothing may throw. A race, if reintroduced, + /// would show up here only probabilistically. + /// Converter registrations are process-wide and can't be undone without changing production code, so the + /// no-op converters stay registered; the conversion cache entries added here are removed again. + /// + [Test] + public async Task ShouldConvertConcurrentlyWhileCachingAndRegistering() { + // List over distinct T gives hundreds of message types that no other test converts + var messageTypes = typeof(object).Assembly.GetExportedTypes() + .Where(t => t is { ContainsGenericParameters: false, IsByRefLike: false, IsPointer: false } && t != typeof(void) && !(t.IsAbstract && t.IsSealed)) + .Take(500) + .Select(t => typeof(List<>).MakeGenericType(t)) + .ToArray(); + + const int workers = 8; + using var start = new Barrier(workers + 1); + + try { + await RunConcurrently(messageTypes, workers, start); + } finally { + foreach (var messageType in messageTypes) MessageConsumeContextConverter.ConversionCache.TryRemove(messageType, out _); + } + } + + static async Task RunConcurrently(Type[] messageTypes, int workers, Barrier start) { + var conversions = Enumerable.Range(0, workers) + .Select( + worker => Task.Run( + () => { + start.SignalAndWait(); + + // Each worker walks the types from a different offset, so the same types are + // being added by one worker while another is looking them up + for (var i = 0; i < messageTypes.Length; i++) { + var messageType = messageTypes[(i + worker * messageTypes.Length / workers) % messageTypes.Length]; + var message = Activator.CreateInstance(messageType)!; + var typed = CreateContext(message).ConvertToGeneric(); + + var expected = typeof(MessageConsumeContext<>).MakeGenericType(messageType); + + if (typed.GetType() != expected || !ReferenceEquals(typed.Message, message)) { + throw new InvalidOperationException($"Expected {expected}, got {typed.GetType()}"); + } + } + } + ) + ) + .ToArray(); + + var registrations = Task.Run( + () => { + start.SignalAndWait(); + + // Converters that handle nothing, so conversions still fall through to the cache + for (var i = 0; i < 8; i++) { + MessageConsumeContextConverter.Register(_ => null); + Thread.Yield(); + } + } + ); + + await Task.WhenAll([..conversions, registrations]); + } + + static MessageConsumeContext CreateContext(object message) + => new( + Guid.NewGuid().ToString(), + message.GetType().Name, + "application/json", + "test-stream", + 0, + 0, + 0, + 0, + DateTime.UtcNow, + message, + null, + "test-subscription", + CancellationToken.None + ); +} diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/MurmurHash3Tests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/MurmurHash3Tests.cs new file mode 100644 index 000000000..433be5255 --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/MurmurHash3Tests.cs @@ -0,0 +1,23 @@ +using Eventuous.Subscriptions.Filters.Partitioning; + +namespace Eventuous.Tests.Subscriptions; + +public class MurmurHash3Tests { + // Expected values were captured from the original pointer-based implementation, so this test + // proves the span-based rewrite produces identical hashes (and therefore identical partitions). + [Test] + [Arguments("", 2183108998u)] + [Arguments("a", 3484832574u)] + [Arguments("ab", 2947271106u)] + [Arguments("abc", 3529399110u)] + [Arguments("Order-123", 1952828587u)] + [Arguments("ünïcødé-ключ-😀", 2930121001u)] + public async Task ShouldProduceKnownHash(string key, uint expected) { + await Assert.That(MurmurHash3.Hash(key)).IsEqualTo(expected); + } + + [Test] + public async Task ShouldThrowOnNull() { + await Assert.That(() => MurmurHash3.Hash(null!)).Throws(); + } +} diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/TracingFilterTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/TracingFilterTests.cs new file mode 100644 index 000000000..bdc9e9ac8 --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/TracingFilterTests.cs @@ -0,0 +1,101 @@ +using System.Diagnostics; +using Eventuous.Diagnostics; +using Eventuous.Subscriptions.Context; +using Eventuous.Subscriptions.Filters; + +namespace Eventuous.Tests.Subscriptions; + +[NotInParallel] +public class TracingFilterTests : IDisposable { + const string TestSourceName = "eventuous.tests.tracing-filter"; + + readonly ActivitySource _testSource = new(TestSourceName); + readonly ActivityListener _listener; + + public TracingFilterTests() { + _listener = new() { + ShouldListenTo = source => source.Name == EventuousDiagnostics.InstrumentationName || source.Name == TestSourceName, + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded + }; + ActivitySource.AddActivityListener(_listener); + } + + [Test] + public async Task ShouldNotStopTheActivityItReuses() { + // The subscription's own activity, which the subscription stops once it's done with the message + using var subscriptionActivity = _testSource.StartActivity("subscription"); + await Assert.That(subscriptionActivity).IsNotNull(); + + var context = TestContext.CreateContext(); + context.ParentContext = subscriptionActivity!.Context; + + var recorder = new ActivityRecorder(); + var pipe = new ConsumePipe().AddFilterLast(new TracingFilter("test-consumer")).AddFilterLast(recorder); + + await pipe.Send(context); + + await Assert.That(recorder.Seen).IsSameReferenceAs(subscriptionActivity); + await Assert.That(subscriptionActivity.IsStopped).IsFalse(); + await Assert.That(subscriptionActivity.Duration).IsEqualTo(TimeSpan.Zero); + } + + [Test] + public async Task ShouldStopTheActivityItStarts() { + var context = TestContext.CreateContext(); + context.ParentContext = new ActivityContext(ActivityTraceId.CreateRandom(), ActivitySpanId.CreateRandom(), ActivityTraceFlags.Recorded); + + var recorder = new ActivityRecorder(); + var pipe = new ConsumePipe().AddFilterLast(new TracingFilter("test-consumer")).AddFilterLast(recorder); + + await pipe.Send(context); + + await Assert.That(recorder.Seen).IsNotNull(); + await Assert.That(recorder.Seen!.Source.Name).IsEqualTo(EventuousDiagnostics.InstrumentationName); + await Assert.That(recorder.Seen.IsStopped).IsTrue(); + } + + [Test] + public async Task ShouldKeepTheErrorStatusOfAFailedMessage() { + using var subscriptionActivity = _testSource.StartActivity("subscription"); + await Assert.That(subscriptionActivity).IsNotNull(); + + var context = TestContext.CreateContext(); + context.ParentContext = subscriptionActivity!.Context; + + var pipe = new ConsumePipe().AddFilterLast(new TracingFilter("test-consumer")).AddFilterLast(new NackingFilter()); + + await pipe.Send(context); + + await Assert.That(context.HasFailed()).IsTrue(); + await Assert.That(subscriptionActivity.Status).IsEqualTo(ActivityStatusCode.Error); + } + + /// + /// Terminal filter recording the activity the handler would run under. + /// + sealed class ActivityRecorder : ConsumeFilter { + public Activity? Seen { get; private set; } + + protected override ValueTask Send(IMessageConsumeContext context, LinkedListNode? next) { + Seen = Activity.Current; + + return default; + } + } + + /// + /// Terminal filter failing the message the way the default consumer does: with Nack, not by throwing. + /// + sealed class NackingFilter : ConsumeFilter { + protected override ValueTask Send(IMessageConsumeContext context, LinkedListNode? next) { + context.Nack("test-handler", new InvalidOperationException("handler failed")); + + return default; + } + } + + public void Dispose() { + _listener.Dispose(); + _testSource.Dispose(); + } +} From 3eab83204e0aaa602c3504218c2dcc9eb44a1b5d Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 00:46:18 +0200 Subject: [PATCH 2/7] fix(diagnostics): observe every traced instance and skip unobserved measures - TracedEventHandler and TracedCommandService use a shared static DiagnosticListener, so metrics come from every instance, not only the most recently created one - Handlers, command services and the traced event writer skip measures when nobody listens; Measure uses Stopwatch - Activity names are cached per message type - ActivityStatus.Ok() is a shared instance; GetParentTag reads the tag without enumerating - The subscription duration metric test runs for every store Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Diagnostics/CommandServiceActivity.cs | 14 +- .../Diagnostics/TracedCommandService.cs | 6 +- .../ActivityExtensions.cs | 3 +- .../Eventuous.Diagnostics/ActivityStatus.cs | 5 +- .../Eventuous.Diagnostics/Metrics/Measure.cs | 5 +- .../Diagnostics/Tracing/BaseTracer.cs | 18 ++- .../Diagnostics/Tracing/TracedEventWriter.cs | 10 +- .../EventSubscription.cs | 12 +- .../Handlers/TracedEventHandler.cs | 20 ++- .../TracedEventHandlerTests.cs | 143 ++++++++++++++++++ .../Eventuous.Tests/TracedEventWriterTests.cs | 104 +++++++++++++ .../ActivityHelpersTests.cs | 50 ++++++ .../TracedCommandServiceMetricsTests.cs | 58 +++++++ .../MetricsTests.cs | 22 ++- .../Metrics/MetricsTests.cs | 6 + .../Metrics/MetricsTests.cs | 6 + .../Metrics/MetricsTests.cs | 6 + .../Metrics/MetricsTests.cs | 6 + 18 files changed, 452 insertions(+), 42 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/TracedEventHandlerTests.cs create mode 100644 src/Core/test/Eventuous.Tests/TracedEventWriterTests.cs create mode 100644 src/Diagnostics/test/Eventuous.Tests.Diagnostics/ActivityHelpersTests.cs create mode 100644 src/Diagnostics/test/Eventuous.Tests.Diagnostics/TracedCommandServiceMetricsTests.cs diff --git a/src/Core/src/Eventuous.Application/Diagnostics/CommandServiceActivity.cs b/src/Core/src/Eventuous.Application/Diagnostics/CommandServiceActivity.cs index 4f9ef481e..5a7b6650c 100644 --- a/src/Core/src/Eventuous.Application/Diagnostics/CommandServiceActivity.cs +++ b/src/Core/src/Eventuous.Application/Diagnostics/CommandServiceActivity.cs @@ -9,27 +9,33 @@ namespace Eventuous.Diagnostics; using Tracing; static class CommandServiceActivity { + // Shared by every traced service, whatever its state type: the metrics listener only observes the latest + // listener announced under a name, so a listener per service would hide the metrics of all but one. + static readonly DiagnosticSource MetricsSource = new DiagnosticListener(CommandServiceMetrics.ListenerName); + public static async Task> TryExecute( string appServiceTypeName, TCommand command, - DiagnosticSource diagnosticSource, HandleCommand handleCommand, CancellationToken cancellationToken ) where TCommand : class where T : State, new() { var cmdName = command.GetType().Name; using var activity = StartActivity(appServiceTypeName, cmdName); - using var measure = Measure.Start(diagnosticSource, new CommandServiceMetricsContext(appServiceTypeName, cmdName)); + // Nobody listening means nothing to record, so skip the measure and its context rather than allocate both. + using var measure = MetricsSource.IsEnabled(Measure.EventName) + ? Measure.Start(MetricsSource, new CommandServiceMetricsContext(appServiceTypeName, cmdName)) + : null; try { var result = await handleCommand(command, cancellationToken).NoContext(); activity?.SetActivityStatus(result is { Success: true } ? ActivityStatus.Ok() : ActivityStatus.Error(result.Exception)); - if (!result.Success) measure.SetError(); + if (!result.Success) measure?.SetError(); return result; } catch (Exception e) { activity?.SetActivityStatus(ActivityStatus.Error(e)); - measure.SetError(); + measure?.SetError(); throw; } diff --git a/src/Core/src/Eventuous.Application/Diagnostics/TracedCommandService.cs b/src/Core/src/Eventuous.Application/Diagnostics/TracedCommandService.cs index 8de58ca9b..d4eb8575f 100644 --- a/src/Core/src/Eventuous.Application/Diagnostics/TracedCommandService.cs +++ b/src/Core/src/Eventuous.Application/Diagnostics/TracedCommandService.cs @@ -1,8 +1,6 @@ // Copyright (C) Eventuous HQ OÜ. All rights reserved // Licensed under the Apache License, Version 2.0. -using System.Diagnostics; - namespace Eventuous.Diagnostics; public class TracedCommandService : ICommandService where TState : State, new() { @@ -10,8 +8,7 @@ namespace Eventuous.Diagnostics; ICommandService InnerService { get; } - readonly string _appServiceTypeName; - readonly DiagnosticSource _metricsSource = new DiagnosticListener(CommandServiceMetrics.ListenerName); + readonly string _appServiceTypeName; TracedCommandService(ICommandService appService) { _appServiceTypeName = appService.GetType().Name; @@ -23,7 +20,6 @@ public Task> Handle(TCommand command, CancellationToken => CommandServiceActivity.TryExecute( _appServiceTypeName, command, - _metricsSource, InnerService.Handle, cancellationToken ); diff --git a/src/Core/src/Eventuous.Diagnostics/ActivityExtensions.cs b/src/Core/src/Eventuous.Diagnostics/ActivityExtensions.cs index 76d2cc704..fa5c4b46a 100644 --- a/src/Core/src/Eventuous.Diagnostics/ActivityExtensions.cs +++ b/src/Core/src/Eventuous.Diagnostics/ActivityExtensions.cs @@ -5,9 +5,10 @@ namespace Eventuous.Diagnostics; public static class ActivityExtensions { extension(Activity activity) { + // Same result as searching Tags, which only yields string values, without allocating an enumerator and a closure. [MethodImpl(MethodImplOptions.AggressiveInlining)] string? GetParentTag(string tag) - => activity.Parent?.Tags.FirstOrDefault(x => x.Key == tag).Value; + => activity.Parent?.GetTagItem(tag) as string; [MethodImpl(MethodImplOptions.AggressiveInlining)] public Activity CopyParentTag(string tag, string? parentTag = null) { diff --git a/src/Core/src/Eventuous.Diagnostics/ActivityStatus.cs b/src/Core/src/Eventuous.Diagnostics/ActivityStatus.cs index a21cc4468..aaefef9a2 100644 --- a/src/Core/src/Eventuous.Diagnostics/ActivityStatus.cs +++ b/src/Core/src/Eventuous.Diagnostics/ActivityStatus.cs @@ -4,8 +4,11 @@ namespace Eventuous.Diagnostics; public record ActivityStatus(ActivityStatusCode StatusCode, string? Description, Exception? Exception) { + static readonly ActivityStatus PlainOk = new(ActivityStatusCode.Ok, null, null); + + // Shared when there's no description: the record is immutable, and this is the status of every successful span. public static ActivityStatus Ok(string? description = null) - => new(ActivityStatusCode.Ok, description, null); + => description == null ? PlainOk : new(ActivityStatusCode.Ok, description, null); public static ActivityStatus Error(Exception? exception = null, string? description = null) => new(ActivityStatusCode.Error, description ?? exception?.Message, exception); diff --git a/src/Core/src/Eventuous.Diagnostics/Metrics/Measure.cs b/src/Core/src/Eventuous.Diagnostics/Metrics/Measure.cs index 19f71b971..c5ac0bc9f 100644 --- a/src/Core/src/Eventuous.Diagnostics/Metrics/Measure.cs +++ b/src/Core/src/Eventuous.Diagnostics/Metrics/Measure.cs @@ -16,12 +16,11 @@ public sealed class Measure(DiagnosticSource diagnosticSource, object context) : Justification = "MeasureContext is not referencing anything." )] void Record() { - var stoppedAt = DateTime.UtcNow; - var duration = stoppedAt - _startedAt; + var duration = Stopwatch.GetElapsedTime(_startedAt); diagnosticSource.Write(EventName, new MeasureContext(duration, _error, context)); } - readonly DateTime _startedAt = DateTime.UtcNow; + readonly long _startedAt = Stopwatch.GetTimestamp(); bool _error; diff --git a/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/BaseTracer.cs b/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/BaseTracer.cs index 8f899bc25..b3f3038a2 100644 --- a/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/BaseTracer.cs +++ b/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/BaseTracer.cs @@ -18,7 +18,7 @@ public abstract class BaseTracer { protected async Task Trace(StreamName stream, string operation, Func> task) { using var activity = StartActivity(stream, operation); - using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation)); + using var measure = StartMeasure(operation); try { var result = await task().NoContext(); @@ -27,7 +27,7 @@ protected async Task Trace(StreamName stream, string operation, Func Trace(StreamName stream, string operation, Func task) { using var activity = StartActivity(stream, operation); - using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation)); + using var measure = StartMeasure(operation); try { await task().NoContext(); activity?.SetActivityStatus(ActivityStatus.Ok()); } catch (Exception e) { activity?.SetActivityStatus(ActivityStatus.Error(e)); - measure.SetError(); + measure?.SetError(); throw; } @@ -55,7 +55,7 @@ protected async IAsyncEnumerable TraceEnumerable( [EnumeratorCancellation] CancellationToken cancellationToken = default ) { using var activity = StartActivity(stream, operation); - using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation)); + using var measure = StartMeasure(operation); var enumerator = source.GetAsyncEnumerator(cancellationToken); @@ -67,7 +67,7 @@ protected async IAsyncEnumerable TraceEnumerable( moved = await enumerator.MoveNextAsync().NoContext(); } catch (Exception e) { activity?.SetActivityStatus(ActivityStatus.Error(e)); - measure.SetError(); + measure?.SetError(); throw; } @@ -81,6 +81,12 @@ protected async IAsyncEnumerable TraceEnumerable( activity?.SetActivityStatus(ActivityStatus.Ok()); } + // Nobody listening means nothing to record, so skip the measure and its context rather than allocate both. + private protected Measure? StartMeasure(string operation) + => MetricsSource.IsEnabled(Measure.EventName) + ? Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation)) + : null; + protected static Activity? StartActivity(StreamName stream, string operationName) { if (!EventuousDiagnostics.Enabled) return null; diff --git a/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs b/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs index b5b387449..0dfccaa9b 100644 --- a/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs +++ b/src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs @@ -5,7 +5,6 @@ namespace Eventuous.Diagnostics.Tracing; -using Metrics; using static Constants; public class TracedEventWriter(IEventWriter writer) : BaseTracer, IEventWriter { @@ -21,8 +20,9 @@ CancellationToken cancellationToken ) { using var activity = StartActivity(stream, Operations.AppendEvents); - using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, Operations.AppendEvents)); + using var measure = StartMeasure(Operations.AppendEvents); + // Snapshot, so the inner writer can't see the caller change the collection while the append is in flight var tracedEvents = events .Select(x => x with { Metadata = x.Metadata.AddActivityTags(activity) }) .ToArray(); @@ -34,7 +34,7 @@ CancellationToken cancellationToken return result; } catch (Exception e) { activity?.SetActivityStatus(ActivityStatus.Error(e)); - measure.SetError(); + measure?.SetError(); throw; } @@ -46,7 +46,7 @@ public async Task AppendEvents(IReadOnlyCollection a.StreamName.ToString()))); using var activity = StartActivity(streamNames, Operations.AppendEvents); - using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, Operations.AppendEvents)); + using var measure = StartMeasure(Operations.AppendEvents); var tracedAppends = appends.Select(a => a with { Events = [.. a.Events.Select(x => x with { Metadata = x.Metadata.AddActivityTags(activity) })] }).ToArray(); @@ -57,7 +57,7 @@ public async Task AppendEvents(IReadOnlyCollection : IMessageSubscription, IAsyncDisposa Session? _session; int _disposed; + readonly string _activityNamePrefix; + readonly ConcurrentDictionary _activityNames = new(); + [PublicAPI] public bool IsRunning => Volatile.Read(ref _session) is not null; @@ -43,6 +47,8 @@ protected EventSubscription( EventSerializer = eventSerializer ?? Eventuous.EventSerializer.Default; Options = options; Log = Logger.CreateContext(options.SubscriptionId, loggerFactory); + + _activityNamePrefix = $"{Constants.Components.Subscription}.{options.SubscriptionId}/"; } public string SubscriptionId => Options.SubscriptionId; @@ -227,6 +233,10 @@ void ReportConnectionDropped(Session session, Failure failure) { /// protected virtual SubscriptionRun CreateRun(CancellationToken lifetime) => new(lifetime); + // Keyed defensively: a transport can still hand over a null type, which the dictionary would reject. + string GetActivityName(string? messageType) + => _activityNames.GetOrAdd(messageType ?? "", static (type, prefix) => prefix + type, _activityNamePrefix); + // ReSharper disable once CognitiveComplexity // ReSharper disable once CyclomaticComplexity protected async ValueTask Handler(IMessageConsumeContext context) { @@ -246,7 +256,7 @@ protected async ValueTask Handler(IMessageConsumeContext context) { // a pure allocation leak, hot since checkpoint-reached contexts arrive payload-less. var activity = EventuousDiagnostics.Enabled && context.Message != null ? SubscriptionActivity.Create( - $"{Constants.Components.Subscription}.{SubscriptionId}/{context.MessageType}", + GetActivityName(context.MessageType), ActivityKind.Internal, context, EventuousDiagnostics.Tags diff --git a/src/Core/src/Eventuous.Subscriptions/Handlers/TracedEventHandler.cs b/src/Core/src/Eventuous.Subscriptions/Handlers/TracedEventHandler.cs index b2ea7b6eb..d5865e4cb 100644 --- a/src/Core/src/Eventuous.Subscriptions/Handlers/TracedEventHandler.cs +++ b/src/Core/src/Eventuous.Subscriptions/Handlers/TracedEventHandler.cs @@ -1,6 +1,7 @@ // Copyright (C) Eventuous HQ OÜ. All rights reserved // Licensed under the Apache License, Version 2.0. +using System.Collections.Concurrent; using System.Diagnostics; using Eventuous.Diagnostics; using Eventuous.Diagnostics.Metrics; @@ -12,19 +13,26 @@ namespace Eventuous.Subscriptions; using Diagnostics; public class TracedEventHandler(IEventHandler eventHandler) : IEventHandler { - readonly DiagnosticSource _metricsSource = new DiagnosticListener(SubscriptionMetrics.ListenerName); + // One listener for all handlers: the metrics listener only observes the latest listener announced under a name. + static readonly DiagnosticSource MetricsSource = new DiagnosticListener(SubscriptionMetrics.ListenerName); readonly KeyValuePair[] _defaultTags = [new (TelemetryTags.Eventuous.EventHandler, eventHandler.GetType().Name)]; public string DiagnosticName { get; } = eventHandler.DiagnosticName; + readonly string _activityNamePrefix = $"{Constants.Components.EventHandler}.{eventHandler.DiagnosticName}/"; + readonly ConcurrentDictionary _activityNames = new(); + public async ValueTask HandleEvent(IMessageConsumeContext context) { using var activity = SubscriptionActivity - .Create($"{Constants.Components.EventHandler}.{DiagnosticName}/{context.MessageType}", ActivityKind.Internal, tags: _defaultTags) + .Create(GetActivityName(context.MessageType), ActivityKind.Internal, tags: _defaultTags) ?.SetContextTags(context) ?.Start(); - using var measure = Measure.Start(_metricsSource, new SubscriptionMetrics.SubscriptionMetricsContext(DiagnosticName, context)); + // Nobody listening means nothing to record, so skip the measure and its context rather than allocate both. + using var measure = MetricsSource.IsEnabled(Measure.EventName) + ? Measure.Start(MetricsSource, new SubscriptionMetrics.SubscriptionMetricsContext(DiagnosticName, context)) + : null; try { var status = await eventHandler.HandleEvent(context).NoContext(); @@ -38,9 +46,13 @@ public async ValueTask HandleEvent(IMessageConsumeContext c return EventHandlingStatus.Pending; } catch (Exception e) { activity?.SetActivityStatus(ActivityStatus.Error(e, $"Error handling {context.MessageType}")); - measure.SetError(); + measure?.SetError(); throw; } } + + // Keyed defensively: a transport can still hand over a null type, which the dictionary would reject. + string GetActivityName(string? messageType) + => _activityNames.GetOrAdd(messageType ?? "", static (type, prefix) => prefix + type, _activityNamePrefix); } diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/TracedEventHandlerTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/TracedEventHandlerTests.cs new file mode 100644 index 000000000..a6e0fec7e --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/TracedEventHandlerTests.cs @@ -0,0 +1,143 @@ +using System.Diagnostics; +using System.Diagnostics.Metrics; +using Eventuous.Diagnostics; +using Eventuous.Subscriptions; +using Eventuous.Subscriptions.Context; +using Eventuous.Subscriptions.Diagnostics; +using Eventuous.Subscriptions.Filters; +using Microsoft.Extensions.Logging.Abstractions; +using TUnit.Assertions.Enums; + +namespace Eventuous.Tests.Subscriptions; + +[NotInParallel] +public class TracedEventHandlerTests { + static MessageConsumeContext CreateContext(string messageType = "TestEvent") + => new( + Guid.NewGuid().ToString(), + messageType, + "application/json", + "test-stream", + 0, + 0, + 0, + 0, + DateTime.UtcNow, + new object(), + null, + "test-subscription", + CancellationToken.None + ); + + [Test] + public async Task ShouldMeasureEveryTracedHandler() { + var handlers = new List(); + + using var meterListener = new MeterListener { + InstrumentPublished = (instrument, listener) => { + if (instrument.Name == SubscriptionMetrics.ProcessingRateName) listener.EnableMeasurementEvents(instrument); + } + }; + + meterListener.SetMeasurementEventCallback( + (_, _, tags, _) => { + foreach (var tag in tags) { + if (tag.Key != SubscriptionMetrics.EventHandlerTag) continue; + + lock (handlers) handlers.Add(tag.Value); + } + } + ); + meterListener.Start(); + + // The metrics listener has to exist before the handlers: each handler used to announce its own diagnostic + // listener, and the metrics listener only kept the last one it saw. + using var metrics = new SubscriptionMetrics([]); + + var first = new TracedEventHandler(new TracedFirstHandler()); + var second = new TracedEventHandler(new TracedSecondHandler()); + + await first.HandleEvent(CreateContext()); + await second.HandleEvent(CreateContext()); + + await Assert.That(handlers).Contains(nameof(TracedFirstHandler)); + await Assert.That(handlers).Contains(nameof(TracedSecondHandler)); + } + + [Test] + public async Task ShouldNameHandlerActivitiesAfterHandlerAndMessageType() { + var names = new List(); + + using var listener = StartActivityListener(names, "handler."); + + var handler = new TracedEventHandler(new TracedFirstHandler()); + + // Repeated types exercise the cached names, which must match the names built per event before them. + await handler.HandleEvent(CreateContext("TestEvent")); + await handler.HandleEvent(CreateContext("OtherEvent")); + await handler.HandleEvent(CreateContext("TestEvent")); + + await Assert.That(names).IsEquivalentTo( + ["handler.TracedFirstHandler/TestEvent", "handler.TracedFirstHandler/OtherEvent", "handler.TracedFirstHandler/TestEvent"], + CollectionOrdering.Matching + ); + } + + [Test] + public async Task ShouldNameSubscriptionActivitiesAfterSubscriptionAndMessageType() { + var names = new List(); + + using var listener = StartActivityListener(names, "sub."); + + await using var subscription = new TestSubscription(new ConsumePipe().AddDefaultConsumer(new TracedFirstHandler())); + + await subscription.Deliver(CreateContext("TestEvent")); + await subscription.Deliver(CreateContext("OtherEvent")); + await subscription.Deliver(CreateContext("TestEvent")); + + await Assert.That(names).IsEquivalentTo( + ["sub.traced-sub/TestEvent", "sub.traced-sub/OtherEvent", "sub.traced-sub/TestEvent"], + CollectionOrdering.Matching + ); + } + + static ActivityListener StartActivityListener(List names, string prefix) { + var listener = new ActivityListener { + ShouldListenTo = source => source.Name == EventuousDiagnostics.InstrumentationName, + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + ActivityStarted = activity => { + if (!activity.OperationName.StartsWith(prefix, StringComparison.Ordinal)) return; + + lock (names) names.Add(activity.OperationName); + } + }; + ActivitySource.AddActivityListener(listener); + + return listener; + } + + record TestOptions : SubscriptionOptions; + + sealed class TestSubscription(ConsumePipe pipe) + : EventSubscription(new() { SubscriptionId = "traced-sub" }, pipe, NullLoggerFactory.Instance, null) { + protected override ValueTask Connect(SubscriptionRun run) => default; + + public ValueTask Deliver(IMessageConsumeContext context) { + context.LogContext = Log; + + return Handler(context); + } + } + + class TracedFirstHandler : IEventHandler { + public string DiagnosticName => nameof(TracedFirstHandler); + + public ValueTask HandleEvent(IMessageConsumeContext context) => ValueTask.FromResult(EventHandlingStatus.Success); + } + + class TracedSecondHandler : IEventHandler { + public string DiagnosticName => nameof(TracedSecondHandler); + + public ValueTask HandleEvent(IMessageConsumeContext context) => ValueTask.FromResult(EventHandlingStatus.Success); + } +} diff --git a/src/Core/test/Eventuous.Tests/TracedEventWriterTests.cs b/src/Core/test/Eventuous.Tests/TracedEventWriterTests.cs new file mode 100644 index 000000000..98f839396 --- /dev/null +++ b/src/Core/test/Eventuous.Tests/TracedEventWriterTests.cs @@ -0,0 +1,104 @@ +using System.Diagnostics; +using System.Diagnostics.Metrics; +using Eventuous.Diagnostics; +using Eventuous.Diagnostics.Tracing; + +namespace Eventuous.Tests; + +[NotInParallel] +public class TracedEventWriterTests : IDisposable { + readonly ActivityListener _listener = new() { + ShouldListenTo = source => source.Name == EventuousDiagnostics.InstrumentationName, + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded + }; + + public TracedEventWriterTests() { + ActivitySource.AddActivityListener(_listener); + Activity.Current = null; + } + + [Test] + public async Task ShouldPassASnapshotOfTheEventsEnrichedWithTracingMeta(CancellationToken cancellationToken) { + var inner = new CapturingWriter(); + var writer = TracedEventWriter.Trace(inner); + var events = new[] { new NewStreamEvent(Guid.NewGuid(), new object(), new()) }; + + await writer.AppendEvents(new("stream-1"), ExpectedStreamVersion.Any, events, cancellationToken); + + await Assert.That(inner.Received).IsNotSameReferenceAs(events); + await Assert.That(inner.Received!.Single().Id).IsEqualTo(events[0].Id); + await Assert.That(inner.Received!.Single().Metadata.GetTracingMeta().TraceId).IsNotNull(); + } + + [Test] + public async Task ShouldPassASnapshotOfTheAppendsEnrichedWithTracingMeta(CancellationToken cancellationToken) { + var inner = new CapturingWriter(); + var writer = TracedEventWriter.Trace(inner); + var events = new[] { new NewStreamEvent(Guid.NewGuid(), new object(), new()) }; + var appends = new[] { new NewStreamAppend(new("stream-1"), ExpectedStreamVersion.Any, events) }; + + await writer.AppendEvents(appends, cancellationToken); + + await Assert.That(inner.ReceivedAppends).IsNotSameReferenceAs(appends); + var received = inner.ReceivedAppends!.Single().Events; + await Assert.That(received).IsNotSameReferenceAs(events); + await Assert.That(received.Single().Id).IsEqualTo(events[0].Id); + await Assert.That(received.Single().Metadata.GetTracingMeta().TraceId).IsNotNull(); + } + + [Test] + public async Task ShouldMeasureBothAppendsWhenMetricsAreObserved(CancellationToken cancellationToken) { + var components = new List(); + + using var meterListener = new MeterListener { + InstrumentPublished = (instrument, listener) => { + if (instrument.Name == PersistenceMetrics.ProcessingRateName) listener.EnableMeasurementEvents(instrument); + } + }; + + meterListener.SetMeasurementEventCallback( + (_, _, tags, _) => { + foreach (var tag in tags) { + if (tag.Key != "component") continue; + + lock (components) components.Add(tag.Value); + } + } + ); + meterListener.Start(); + + using var metrics = new PersistenceMetrics(); + + var writer = TracedEventWriter.Trace(new CapturingWriter()); + var events = new[] { new NewStreamEvent(Guid.NewGuid(), new object(), new()) }; + + await writer.AppendEvents(new("stream-1"), ExpectedStreamVersion.Any, events, cancellationToken); + await writer.AppendEvents([new NewStreamAppend(new("stream-1"), ExpectedStreamVersion.Any, events)], cancellationToken); + + await Assert.That(components.Count(x => Equals(x, nameof(CapturingWriter)))).IsEqualTo(2); + } + + public void Dispose() => _listener.Dispose(); + + class CapturingWriter : IEventWriter { + public IReadOnlyCollection? Received { get; private set; } + public IReadOnlyCollection? ReceivedAppends { get; private set; } + + public Task AppendEvents( + StreamName stream, + ExpectedStreamVersion expectedVersion, + IReadOnlyCollection events, + CancellationToken cancellationToken + ) { + Received = events; + + return Task.FromResult(new AppendEventsResult(0, 0)); + } + + public Task AppendEvents(IReadOnlyCollection appends, CancellationToken cancellationToken) { + ReceivedAppends = appends; + + return Task.FromResult([new(0, 0)]); + } + } +} diff --git a/src/Diagnostics/test/Eventuous.Tests.Diagnostics/ActivityHelpersTests.cs b/src/Diagnostics/test/Eventuous.Tests.Diagnostics/ActivityHelpersTests.cs new file mode 100644 index 000000000..9332c2e61 --- /dev/null +++ b/src/Diagnostics/test/Eventuous.Tests.Diagnostics/ActivityHelpersTests.cs @@ -0,0 +1,50 @@ +using System.Diagnostics; +using Eventuous.Diagnostics; + +namespace Eventuous.Tests.Diagnostics; + +public class ActivityHelpersTests { + [Test] + public async Task ShouldShareTheOkStatusWithoutDescription() { + var first = ActivityStatus.Ok(); + var second = ActivityStatus.Ok(); + + await Assert.That(ReferenceEquals(first, second)).IsTrue(); + await Assert.That(first.StatusCode).IsEqualTo(ActivityStatusCode.Ok); + await Assert.That(first.Description).IsNull(); + await Assert.That(first.Exception).IsNull(); + } + + [Test] + public async Task ShouldKeepTheDescriptionOfAnOkStatus() { + var status = ActivityStatus.Ok("done"); + + await Assert.That(status.StatusCode).IsEqualTo(ActivityStatusCode.Ok); + await Assert.That(status.Description).IsEqualTo("done"); + } + + [Test] + public async Task ShouldCopyOnlyStringTagsFromTheParent() { + var previous = Activity.Current; + + try { + using var parent = new Activity("parent").Start(); + parent.SetTag("text", "value"); + parent.SetTag("number", 42); + + using var child = new Activity("child").Start(); + + child.CopyParentTag("text").CopyParentTag("number").CopyParentTag("missing"); + child.SetOrCopyParentTag("renamed", null, "text"); + child.SetOrCopyParentTag("explicit", "own", "text"); + + await Assert.That(child.GetTagItem("text")).IsEqualTo("value"); + await Assert.That(child.GetTagItem("number")).IsNull(); + await Assert.That(child.GetTagItem("missing")).IsNull(); + await Assert.That(child.GetTagItem("renamed")).IsEqualTo("value"); + await Assert.That(child.GetTagItem("explicit")).IsEqualTo("own"); + } finally { + Activity.Current = previous; + } + } +} diff --git a/src/Diagnostics/test/Eventuous.Tests.Diagnostics/TracedCommandServiceMetricsTests.cs b/src/Diagnostics/test/Eventuous.Tests.Diagnostics/TracedCommandServiceMetricsTests.cs new file mode 100644 index 000000000..954bf490c --- /dev/null +++ b/src/Diagnostics/test/Eventuous.Tests.Diagnostics/TracedCommandServiceMetricsTests.cs @@ -0,0 +1,58 @@ +using System.Diagnostics.Metrics; +using Eventuous.Diagnostics; + +namespace Eventuous.Tests.Diagnostics; + +[NotInParallel] +public class TracedCommandServiceMetricsTests { + [Test] + public async Task ShouldMeasureEveryTracedService() { + var services = new List(); + + using var meterListener = new MeterListener { + InstrumentPublished = (instrument, listener) => { + if (instrument.Name == CommandServiceMetrics.ProcessingRateName) listener.EnableMeasurementEvents(instrument); + } + }; + + meterListener.SetMeasurementEventCallback( + (_, _, tags, _) => { + foreach (var tag in tags) { + if (tag.Key != "command-service") continue; + + lock (services) services.Add(tag.Value); + } + } + ); + meterListener.Start(); + + // The metrics listener has to exist before the services: each service used to announce its own diagnostic + // listener, and the metrics listener only kept the last one it saw. + using var metrics = new CommandServiceMetrics(); + + var first = TracedCommandService.Trace(new FirstService()); + var second = TracedCommandService.Trace(new SecondService()); + + await first.Handle(new TestCommand(), CancellationToken.None); + await second.Handle(new TestCommand(), CancellationToken.None); + + await Assert.That(services).Contains(nameof(FirstService)); + await Assert.That(services).Contains(nameof(SecondService)); + } + + record TestCommand; + + record FirstState : State; + + record SecondState : State; + + class FirstService : ICommandService { + public Task> Handle(TCommand command, CancellationToken cancellationToken) where TCommand : class + => Task.FromResult(Result.FromSuccess(new(), [], 0)); + } + + class SecondService : ICommandService { + public Task> Handle(TCommand command, CancellationToken cancellationToken) where TCommand : class + => Task.FromResult(Result.FromSuccess(new(), [], 0)); + } +} diff --git a/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs b/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs index 94237e556..077a3c7a0 100644 --- a/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs +++ b/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs @@ -20,18 +20,16 @@ protected async Task ShouldMeasureSubscriptionGapCountBase() { await gapCount.CheckTag(fixture.DefaultTagKey, fixture.DefaultTagValue); } - // [Fact] - // [Trait("Category", "Diagnostics")] - // public void ShouldMeasureSubscriptionDuration() { - // Fixture.Output?.WriteLine($"Stream {Fixture.Stream}"); - // Assert.NotNull(_values); - // var duration = GetValue(_values, SubscriptionMetrics.ProcessingRateName)!; - // - // duration.Should().NotBeNull(); - // duration.CheckTag(SubscriptionMetrics.SubscriptionIdTag, Fixture.SubscriptionId); - // duration.CheckTag(Fixture.DefaultTagKey, Fixture.DefaultTagValue); - // duration.CheckTag(SubscriptionMetrics.MessageTypeTag, TestEvent.TypeName); - // } + protected async Task ShouldMeasureSubscriptionDurationBase() { + TestContext.Current?.OutputWriter.WriteLine($"Stream {fixture.Stream}"); + await Assert.That(_values).IsNotNull(); + var duration = GetValue(_values!, SubscriptionMetrics.ProcessingRateName); + + await Assert.That(duration).IsNotNull(); + await duration!.CheckTag(SubscriptionMetrics.SubscriptionIdTag, fixture.SubscriptionId); + await duration.CheckTag(fixture.DefaultTagKey, fixture.DefaultTagValue); + await duration.CheckTag(SubscriptionMetrics.MessageTypeTag, TestEvent.TypeName); + } static MetricValue? GetValue(MetricValue[] values, string metric) => values.FirstOrDefault(x => x.Name == metric); diff --git a/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Metrics/MetricsTests.cs b/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Metrics/MetricsTests.cs index 9ecc73c5e..62b532ae7 100644 --- a/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Metrics/MetricsTests.cs +++ b/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Metrics/MetricsTests.cs @@ -10,6 +10,12 @@ public class MetricsTests(MetricsFixture fixture) : MetricsTestsBase(fixture) { public async Task ShouldMeasureSubscriptionGapCountBase_Esdb() { await ShouldMeasureSubscriptionGapCountBase(); } + + [Test] + [Retry(3)] + public async Task ShouldMeasureSubscriptionDurationBase_Esdb() { + await ShouldMeasureSubscriptionDurationBase(); + } } [ClassDataSource] diff --git a/src/Postgres/test/Eventuous.Tests.Postgres/Metrics/MetricsTests.cs b/src/Postgres/test/Eventuous.Tests.Postgres/Metrics/MetricsTests.cs index 7ce59763f..0b93c79d4 100644 --- a/src/Postgres/test/Eventuous.Tests.Postgres/Metrics/MetricsTests.cs +++ b/src/Postgres/test/Eventuous.Tests.Postgres/Metrics/MetricsTests.cs @@ -10,6 +10,12 @@ public class MetricsTests(MetricsFixture fixture) : MetricsTestsBase(fixture) { public async Task ShouldMeasureSubscriptionGapCountBase_Postgres() { await ShouldMeasureSubscriptionGapCountBase(); } + + [Test] + [Retry(3)] + public async Task ShouldMeasureSubscriptionDurationBase_Postgres() { + await ShouldMeasureSubscriptionDurationBase(); + } } [ClassDataSource] diff --git a/src/SqlServer/test/Eventuous.Tests.SqlServer/Metrics/MetricsTests.cs b/src/SqlServer/test/Eventuous.Tests.SqlServer/Metrics/MetricsTests.cs index a5ef2e841..cfe21da99 100644 --- a/src/SqlServer/test/Eventuous.Tests.SqlServer/Metrics/MetricsTests.cs +++ b/src/SqlServer/test/Eventuous.Tests.SqlServer/Metrics/MetricsTests.cs @@ -10,6 +10,12 @@ public class MetricsTests(MetricsFixture fixture) : MetricsTestsBase(fixture) { public async Task ShouldMeasureSubscriptionGapCountBase_SqlServer() { await ShouldMeasureSubscriptionGapCountBase(); } + + [Test] + [Retry(3)] + public async Task ShouldMeasureSubscriptionDurationBase_SqlServer() { + await ShouldMeasureSubscriptionDurationBase(); + } } [ClassDataSource] diff --git a/src/Sqlite/test/Eventuous.Tests.Sqlite/Metrics/MetricsTests.cs b/src/Sqlite/test/Eventuous.Tests.Sqlite/Metrics/MetricsTests.cs index 4dc665e7a..b124b24d3 100644 --- a/src/Sqlite/test/Eventuous.Tests.Sqlite/Metrics/MetricsTests.cs +++ b/src/Sqlite/test/Eventuous.Tests.Sqlite/Metrics/MetricsTests.cs @@ -10,6 +10,12 @@ public class MetricsTests(MetricsFixture fixture) : MetricsTestsBase(fixture) { public async Task ShouldMeasureSubscriptionGapCountBase_Sqlite() { await ShouldMeasureSubscriptionGapCountBase(); } + + [Test] + [Retry(3)] + public async Task ShouldMeasureSubscriptionDurationBase_Sqlite() { + await ShouldMeasureSubscriptionDurationBase(); + } } [ClassDataSource] From 7cab455c9fea34ab3706fd447aba0e72750d1a3f Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 00:46:18 +0200 Subject: [PATCH 3/7] fix(application): return ThrowingCommandService results on success - ThrowingCommandService returns the inner result on success instead of always throwing - Resolver null-check messages are built only on failure - ApplicationEventSource builds error event text only when enabled Co-Authored-By: Claude Opus 5.5 (1M context) --- .../AggregateService/CommandHandlerBuilder.cs | 4 +-- .../Diagnostics/ApplicationEventSource.cs | 8 ++++-- .../ThrowingCommandService.cs | 2 +- .../ResolverNullCheckTests.cs | 23 ++++++++++++++++ .../ThrowingCommandServiceTests.cs | 26 +++++++++++++++++++ 5 files changed, 58 insertions(+), 5 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests.Application/ResolverNullCheckTests.cs create mode 100644 src/Core/test/Eventuous.Tests.Application/ThrowingCommandServiceTests.cs diff --git a/src/Core/src/Eventuous.Application/AggregateService/CommandHandlerBuilder.cs b/src/Core/src/Eventuous.Application/AggregateService/CommandHandlerBuilder.cs index 264307644..958b1060e 100644 --- a/src/Core/src/Eventuous.Application/AggregateService/CommandHandlerBuilder.cs +++ b/src/Core/src/Eventuous.Application/AggregateService/CommandHandlerBuilder.cs @@ -245,9 +245,9 @@ RegisteredHandler Build() { ); Func DefaultResolveWriter() - => _ => Ensure.NotNull(writer, $"Function to resolve event writer from {typeof(TCommand).Name} is not defined and no default writer is set"); + => _ => writer ?? throw new ArgumentNullException($"Function to resolve event writer from {typeof(TCommand).Name} is not defined and no default writer is set"); Func DefaultResolveReader() - => _ => Ensure.NotNull(reader, $"Function to resolve event reader from {typeof(TCommand).Name} is not defined and no default reader is set"); + => _ => reader ?? throw new ArgumentNullException($"Function to resolve event reader from {typeof(TCommand).Name} is not defined and no default reader is set"); } } diff --git a/src/Core/src/Eventuous.Application/Diagnostics/ApplicationEventSource.cs b/src/Core/src/Eventuous.Application/Diagnostics/ApplicationEventSource.cs index bd31de011..f9bd74461 100644 --- a/src/Core/src/Eventuous.Application/Diagnostics/ApplicationEventSource.cs +++ b/src/Core/src/Eventuous.Application/Diagnostics/ApplicationEventSource.cs @@ -18,10 +18,14 @@ class ApplicationEventSource : EventSource { const int CommandHandlerRegisteredId = 5; [NonEvent] - public void CommandHandlerNotFound() => CommandHandlerNotFound(typeof(TCommand).Name); + public void CommandHandlerNotFound() { + if (IsEnabled(EventLevel.Error, EventKeywords.All)) CommandHandlerNotFound(typeof(TCommand).Name); + } [NonEvent] - public void ErrorHandlingCommand(Exception e) => ErrorHandlingCommand(typeof(TCommand).Name, e.ToString()); + public void ErrorHandlingCommand(Exception e) { + if (IsEnabled(EventLevel.Error, EventKeywords.All)) ErrorHandlingCommand(typeof(TCommand).Name, e.ToString()); + } [NonEvent] public void CommandHandled() { diff --git a/src/Core/src/Eventuous.Application/ThrowingCommandService.cs b/src/Core/src/Eventuous.Application/ThrowingCommandService.cs index cc49091de..d508dc221 100644 --- a/src/Core/src/Eventuous.Application/ThrowingCommandService.cs +++ b/src/Core/src/Eventuous.Application/ThrowingCommandService.cs @@ -15,6 +15,6 @@ public async Task> Handle(TCommand command, Cancellatio result.ThrowIfError(); - throw new ApplicationException($"Error handling command {command}"); + return result; } } diff --git a/src/Core/test/Eventuous.Tests.Application/ResolverNullCheckTests.cs b/src/Core/test/Eventuous.Tests.Application/ResolverNullCheckTests.cs new file mode 100644 index 000000000..d855e0312 --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Application/ResolverNullCheckTests.cs @@ -0,0 +1,23 @@ +using Eventuous.Sut.Domain; + +namespace Eventuous.Tests.Application; + +public class ResolverNullCheckTests { + [Test] + public async Task ShouldFailAtCallTimeWhenNoReaderIsAvailable(CancellationToken cancellationToken) { + var service = new NoStoreService(); + + var exception = await Assert.That(async () => await service.Handle(Helpers.GetBookRoom(), cancellationToken)).Throws(); + + await Assert.That(exception!.Message).Contains("Function to resolve event reader from BookRoom is not defined and no default reader is set"); + } + + class NoStoreService : CommandService { + public NoStoreService() : base((IEventStore?)null) { + On() + .InState(ExpectedState.New) + .GetId(cmd => new(cmd.BookingId)) + .Act((_, _) => { }); + } + } +} diff --git a/src/Core/test/Eventuous.Tests.Application/ThrowingCommandServiceTests.cs b/src/Core/test/Eventuous.Tests.Application/ThrowingCommandServiceTests.cs new file mode 100644 index 000000000..f0b661128 --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Application/ThrowingCommandServiceTests.cs @@ -0,0 +1,26 @@ +using Eventuous.Sut.App; +using Eventuous.Testing; + +namespace Eventuous.Tests.Application; + +public class ThrowingCommandServiceTests { + readonly ThrowingCommandService _service; + + public ThrowingCommandServiceTests() => _service = new(new BookingService(new InMemoryEventStore())); + + [Test] + public async Task ShouldReturnResultOnSuccess(CancellationToken cancellationToken) { + var result = await _service.Handle(Helpers.GetBookRoom(), cancellationToken); + + await Assert.That(result.TryGet(out _)).IsTrue(); + } + + [Test] + public async Task ShouldThrowOnError(CancellationToken cancellationToken) { + var cmd = Helpers.GetBookRoom(); + await _service.Handle(cmd, cancellationToken); + + // The booking already exists, so the second attempt produces an error result + await Assert.That(async () => await _service.Handle(cmd, cancellationToken)).Throws(); + } +} From bd2ed1a3b9cf9fadd246689c6cf1b3a3266b5d5d Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Wed, 30 Sep 2026 23:07:31 +0200 Subject: [PATCH 4/7] perf(domain): stop allocating a view on every Changes access Co-Authored-By: Claude Sonnet 5.5 --- src/Core/src/Eventuous.Domain/Aggregate.cs | 7 +++-- .../Aggregates/ChangesViewTests.cs | 26 +++++++++++++++++++ 2 files changed, 31 insertions(+), 2 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests/Aggregates/ChangesViewTests.cs diff --git a/src/Core/src/Eventuous.Domain/Aggregate.cs b/src/Core/src/Eventuous.Domain/Aggregate.cs index 5548d5acc..7908befd2 100644 --- a/src/Core/src/Eventuous.Domain/Aggregate.cs +++ b/src/Core/src/Eventuous.Domain/Aggregate.cs @@ -13,7 +13,7 @@ namespace Eventuous; /// /// Get the list of pending changes (new events) within the scope of the current operation. /// - public IReadOnlyCollection Changes => _changes.AsReadOnly(); + public IReadOnlyCollection Changes => _changesView ??= _changes.AsReadOnly(); /// /// A collection with all the aggregate events, previously persisted and new @@ -36,10 +36,13 @@ namespace Eventuous; /// The current version is set to the original version when the aggregate is loaded from the store. /// It should increase for each state transition performed within the scope of the current operation. /// - public long CurrentVersion => OriginalVersion + Changes.Count; + public long CurrentVersion => OriginalVersion + _changes.Count; readonly List _changes = []; + // The view wraps the live list, so it can be created once and reused + IReadOnlyCollection? _changesView; + /// /// Adds an event to the list of pending changes. /// diff --git a/src/Core/test/Eventuous.Tests/Aggregates/ChangesViewTests.cs b/src/Core/test/Eventuous.Tests/Aggregates/ChangesViewTests.cs new file mode 100644 index 000000000..7f13924ae --- /dev/null +++ b/src/Core/test/Eventuous.Tests/Aggregates/ChangesViewTests.cs @@ -0,0 +1,26 @@ +namespace Eventuous.Tests.Aggregates; + +using Sut.Domain; + +public class ChangesViewTests { + [Test] + public async Task ShouldReuseTheSameViewThatReflectsNewChanges() { + var booking = new Booking(); + var view = booking.Changes; + + await Assert.That(view).HasCount(0); + await Assert.That(booking.CurrentVersion).IsEqualTo(-1); + + booking.Cancel(); + + await Assert.That(booking.Changes).IsSameReferenceAs(view); + await Assert.That(view).HasCount(1); + await Assert.That(booking.CurrentVersion).IsEqualTo(0); + + booking.ClearChanges(); + + await Assert.That(booking.Changes).IsSameReferenceAs(view); + await Assert.That(view).HasCount(0); + await Assert.That(booking.CurrentVersion).IsEqualTo(-1); + } +} From 34b6c5541f47272e19425944512ffd36df29d84d Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Wed, 30 Sep 2026 23:08:10 +0200 Subject: [PATCH 5/7] perf(serialization): reuse deserialization failure results DefaultEventSerializer and DefaultStaticEventSerializer allocated a new FailedToDeserialize per unknown type, content-type mismatch or empty payload, which is every unregistered event on $all. Records are immutable, so share one static instance per error kind. Co-Authored-By: Claude Sonnet 5.5 --- .../DefaultEventSerializer.cs | 11 ++-- .../DefaultStaticEventSerializer.cs | 11 ++-- .../DynamicSerializerFailureTests.cs | 54 +++++++++++++++++++ .../DeserializationFailureTests.cs | 53 ++++++++++++++++++ 4 files changed, 123 insertions(+), 6 deletions(-) create mode 100644 src/Core/test/Eventuous.Tests.Subscriptions/DynamicSerializerFailureTests.cs create mode 100644 src/Core/test/Eventuous.Tests/DeserializationFailureTests.cs diff --git a/src/Core/src/Eventuous.Serialization.Json.Dynamic/DefaultEventSerializer.cs b/src/Core/src/Eventuous.Serialization.Json.Dynamic/DefaultEventSerializer.cs index 4ec488624..89e7c3250 100644 --- a/src/Core/src/Eventuous.Serialization.Json.Dynamic/DefaultEventSerializer.cs +++ b/src/Core/src/Eventuous.Serialization.Json.Dynamic/DefaultEventSerializer.cs @@ -11,6 +11,11 @@ public class DefaultEventSerializer : IEventSerializer { const string DynamicSerializationMessage = "DefaultEventSerializer uses reflection-based System.Text.Json serialization. Use DefaultStaticEventSerializer with a JsonSerializerContext in trimmed or AOT applications."; + // Failure results are immutable, so one instance per error kind is shared instead of allocating per event + static readonly FailedToDeserialize UnknownType = new(DeserializationError.UnknownType); + static readonly FailedToDeserialize ContentTypeMismatch = new(DeserializationError.ContentTypeMismatch); + static readonly FailedToDeserialize PayloadEmpty = new(DeserializationError.PayloadEmpty); + readonly JsonSerializerOptions _options; readonly ITypeMapper _typeMapper; @@ -29,14 +34,14 @@ public DefaultEventSerializer(JsonSerializerOptions options, ITypeMapper? typeMa public DeserializationResult DeserializeEvent(ReadOnlySpan data, string eventType, string contentType) { var typeMapped = _typeMapper.TryGetType(eventType, out var dataType); - if (!typeMapped) return new FailedToDeserialize(DeserializationError.UnknownType); - if (contentType != ContentType) return new FailedToDeserialize(DeserializationError.ContentTypeMismatch); + if (!typeMapped) return UnknownType; + if (contentType != ContentType) return ContentTypeMismatch; var deserialized = JsonSerializer.Deserialize(data, dataType!, _options); return deserialized != null ? new SuccessfullyDeserialized(deserialized) - : new FailedToDeserialize(DeserializationError.PayloadEmpty); + : PayloadEmpty; } [UnconditionalSuppressMessage("Trimming", "IL2026", Justification = "The constructor is annotated with RequiresUnreferencedCode, so an instance only exists if the caller acknowledged the requirement")] diff --git a/src/Core/src/Eventuous.Serialization/DefaultStaticEventSerializer.cs b/src/Core/src/Eventuous.Serialization/DefaultStaticEventSerializer.cs index 31dc927a3..a675aab02 100644 --- a/src/Core/src/Eventuous.Serialization/DefaultStaticEventSerializer.cs +++ b/src/Core/src/Eventuous.Serialization/DefaultStaticEventSerializer.cs @@ -9,19 +9,24 @@ namespace Eventuous; [PublicAPI] public class DefaultStaticEventSerializer(JsonSerializerContext context, ITypeMapper? typeMapper = null) : IEventSerializer { + // Failure results are immutable, so one instance per error kind is shared instead of allocating per event + static readonly FailedToDeserialize UnknownType = new(DeserializationError.UnknownType); + static readonly FailedToDeserialize ContentTypeMismatch = new(DeserializationError.ContentTypeMismatch); + static readonly FailedToDeserialize PayloadEmpty = new(DeserializationError.PayloadEmpty); + readonly ITypeMapper _typeMapper = typeMapper ?? TypeMap.Instance; public DeserializationResult DeserializeEvent(ReadOnlySpan data, string eventType, string contentType) { var typeMapped = _typeMapper.TryGetType(eventType, out var dataType); - if (!typeMapped) return new FailedToDeserialize(DeserializationError.UnknownType); - if (contentType != ContentType) return new FailedToDeserialize(DeserializationError.ContentTypeMismatch); + if (!typeMapped) return UnknownType; + if (contentType != ContentType) return ContentTypeMismatch; var deserialized = JsonSerializer.Deserialize(data, dataType!, context); return deserialized != null ? new SuccessfullyDeserialized(deserialized) - : new FailedToDeserialize(DeserializationError.PayloadEmpty); + : PayloadEmpty; } public SerializationResult SerializeEvent(object evt) diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/DynamicSerializerFailureTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/DynamicSerializerFailureTests.cs new file mode 100644 index 000000000..28eafdb03 --- /dev/null +++ b/src/Core/test/Eventuous.Tests.Subscriptions/DynamicSerializerFailureTests.cs @@ -0,0 +1,54 @@ +using System.Text.Json; +using static Eventuous.DeserializationResult; + +namespace Eventuous.Tests.Subscriptions; + +/// +/// The reflection-based serializer shares its failure results like the static one; the subscription project already +/// references it, so it's pinned here. +/// +public class DynamicSerializerFailureTests { + static DefaultEventSerializer CreateSerializer(ITypeMapper typeMapper) => new(new(JsonSerializerDefaults.Web), typeMapper); + + static TypeMapper MapperWithTestEvent() { + var typeMapper = new TypeMapper(); + typeMapper.AddType("dynamic-failure-test-event"); + + return typeMapper; + } + + [Test] + public async Task ShouldReuseUnknownTypeFailure() { + var serializer = CreateSerializer(new TypeMapper()); + + var first = serializer.DeserializeEvent("{}"u8, "unregistered", "application/json"); + var second = serializer.DeserializeEvent("{}"u8, "other-unregistered", "application/json"); + + await Assert.That(((FailedToDeserialize)first).Error).IsEqualTo(DeserializationError.UnknownType); + await Assert.That(second).IsSameReferenceAs(first); + } + + [Test] + public async Task ShouldReuseContentTypeMismatchFailure() { + var serializer = CreateSerializer(MapperWithTestEvent()); + + var first = serializer.DeserializeEvent("{}"u8, "dynamic-failure-test-event", "application/xml"); + var second = serializer.DeserializeEvent("{}"u8, "dynamic-failure-test-event", "text/plain"); + + await Assert.That(((FailedToDeserialize)first).Error).IsEqualTo(DeserializationError.ContentTypeMismatch); + await Assert.That(second).IsSameReferenceAs(first); + } + + [Test] + public async Task ShouldReusePayloadEmptyFailure() { + var serializer = CreateSerializer(MapperWithTestEvent()); + + var first = serializer.DeserializeEvent("null"u8, "dynamic-failure-test-event", "application/json"); + var second = serializer.DeserializeEvent("null"u8, "dynamic-failure-test-event", "application/json"); + + await Assert.That(((FailedToDeserialize)first).Error).IsEqualTo(DeserializationError.PayloadEmpty); + await Assert.That(second).IsSameReferenceAs(first); + } + + record TestEvent(string Name); +} diff --git a/src/Core/test/Eventuous.Tests/DeserializationFailureTests.cs b/src/Core/test/Eventuous.Tests/DeserializationFailureTests.cs new file mode 100644 index 000000000..0873ea888 --- /dev/null +++ b/src/Core/test/Eventuous.Tests/DeserializationFailureTests.cs @@ -0,0 +1,53 @@ +using System.Text.Json.Serialization; +using static Eventuous.DeserializationResult; + +namespace Eventuous.Tests; + +public partial class DeserializationFailureTests { + static DefaultStaticEventSerializer CreateSerializer(ITypeMapper typeMapper) => new(TestJsonContext.Default, typeMapper); + + static TypeMapper MapperWithTestEvent() { + var typeMapper = new TypeMapper(); + typeMapper.AddType("test-event"); + + return typeMapper; + } + + [Test] + public async Task ShouldReuseUnknownTypeFailure() { + var serializer = CreateSerializer(new TypeMapper()); + + var first = serializer.DeserializeEvent("{}"u8, "unregistered", "application/json"); + var second = serializer.DeserializeEvent("{}"u8, "other-unregistered", "application/json"); + + await Assert.That(((FailedToDeserialize)first).Error).IsEqualTo(DeserializationError.UnknownType); + await Assert.That(second).IsSameReferenceAs(first); + } + + [Test] + public async Task ShouldReuseContentTypeMismatchFailure() { + var serializer = CreateSerializer(MapperWithTestEvent()); + + var first = serializer.DeserializeEvent("{}"u8, "test-event", "application/xml"); + var second = serializer.DeserializeEvent("{}"u8, "test-event", "text/plain"); + + await Assert.That(((FailedToDeserialize)first).Error).IsEqualTo(DeserializationError.ContentTypeMismatch); + await Assert.That(second).IsSameReferenceAs(first); + } + + [Test] + public async Task ShouldReusePayloadEmptyFailure() { + var serializer = CreateSerializer(MapperWithTestEvent()); + + var first = serializer.DeserializeEvent("null"u8, "test-event", "application/json"); + var second = serializer.DeserializeEvent("null"u8, "test-event", "application/json"); + + await Assert.That(((FailedToDeserialize)first).Error).IsEqualTo(DeserializationError.PayloadEmpty); + await Assert.That(second).IsSameReferenceAs(first); + } + + record TestEvent(string Name); + + [JsonSerializable(typeof(TestEvent))] + partial class TestJsonContext : JsonSerializerContext; +} From f71b3010d388e339a6a9f7b897f28047afd5e189 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 00:46:18 +0200 Subject: [PATCH 6/7] fix(sql): honour the metadata serializer in SQL and Redis subscriptions - Postgres, SQL Server, SQLite and Redis subscriptions deserialize metadata with the configured IMetadataSerializer through DeserializeMeta, so malformed metadata is logged (or throws DeserializationException under ThrowOnError) instead of faulting the poll loop - Metadata is skipped for events whose payload didn't deserialize - Redis $all resolves each link with a single-entry XRANGE - Postgres builds its query text once per schema instead of per access Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/Eventuous.Postgresql/Schema.cs | 26 +- .../Subscriptions/SchemaSqlTests.cs | 31 +++ .../RedisAllStreamSubscription.cs | 3 +- .../Subscriptions/RedisSubscriptionBase.cs | 7 +- .../AllStreamLinkResolutionTests.cs | 106 ++++++++ .../MetadataDeserializationTests.cs | 207 +++++++++++++++ .../Subscriptions/SqlSubscriptionBase.cs | 7 +- .../MetadataDeserializationTests.cs | 248 ++++++++++++++++++ 8 files changed, 615 insertions(+), 20 deletions(-) create mode 100644 src/Postgres/test/Eventuous.Tests.Postgres/Subscriptions/SchemaSqlTests.cs create mode 100644 src/Redis/test/Eventuous.Tests.Redis/Subscriptions/AllStreamLinkResolutionTests.cs create mode 100644 src/Redis/test/Eventuous.Tests.Redis/Subscriptions/MetadataDeserializationTests.cs create mode 100644 src/Sqlite/test/Eventuous.Tests.Sqlite/Subscriptions/MetadataDeserializationTests.cs diff --git a/src/Postgres/src/Eventuous.Postgresql/Schema.cs b/src/Postgres/src/Eventuous.Postgresql/Schema.cs index bed75ac09..5b5d63ad4 100644 --- a/src/Postgres/src/Eventuous.Postgresql/Schema.cs +++ b/src/Postgres/src/Eventuous.Postgresql/Schema.cs @@ -17,19 +17,19 @@ public class Schema(string schema = Schema.DefaultSchema) { public string Name => schema; - public string StreamMessage => GetStreamMessageTypeName(schema); - public string AppendEvents => $"select * from {schema}.append_events(@_stream_name, @_expected_version, @_created, @_messages)"; - public string ReadStreamForwards => $"select * from {schema}.read_stream_forwards(@_stream_name, @_from_position, @_count)"; - public string ReadStreamBackwards => $"select * from {schema}.read_stream_backwards(@_stream_name, @_from_position, @_count)"; - public string ReadStreamSub => $"select * from {schema}.read_stream_sub(@_stream_id, @_stream_name, @_from_position, @_count)"; - public string ReadAllForwards => $"select * from {schema}.read_all_forwards(@_from_position, @_count)"; - public string CheckStream => $"select * from {schema}.check_stream(@_stream_name, @_expected_version)"; - public string StreamExists => $"select exists (select 1 from {schema}.streams where stream_name = (@name))"; - public string TruncateStream => $"select * from {schema}.truncate_stream(@_stream_name, @_expected_version, @_position)"; - public string GetCheckpointSql => $"select position from {schema}.checkpoints where id=(@checkpointId)"; - public string AddCheckpointSql => $"insert into {schema}.checkpoints (id) values (@checkpointId)"; - public string UpdateCheckpointSql => $"update {schema}.checkpoints set position=(@position) where id=(@checkpointId)"; - public string TryInsertTombstone => $"select {schema}.try_insert_tombstone(@_gap_position, @_stream_name, @_type, @_id)"; + public string StreamMessage { get; } = GetStreamMessageTypeName(schema); + public string AppendEvents { get; } = $"select * from {schema}.append_events(@_stream_name, @_expected_version, @_created, @_messages)"; + public string ReadStreamForwards { get; } = $"select * from {schema}.read_stream_forwards(@_stream_name, @_from_position, @_count)"; + public string ReadStreamBackwards { get; } = $"select * from {schema}.read_stream_backwards(@_stream_name, @_from_position, @_count)"; + public string ReadStreamSub { get; } = $"select * from {schema}.read_stream_sub(@_stream_id, @_stream_name, @_from_position, @_count)"; + public string ReadAllForwards { get; } = $"select * from {schema}.read_all_forwards(@_from_position, @_count)"; + public string CheckStream { get; } = $"select * from {schema}.check_stream(@_stream_name, @_expected_version)"; + public string StreamExists { get; } = $"select exists (select 1 from {schema}.streams where stream_name = (@name))"; + public string TruncateStream { get; } = $"select * from {schema}.truncate_stream(@_stream_name, @_expected_version, @_position)"; + public string GetCheckpointSql { get; } = $"select position from {schema}.checkpoints where id=(@checkpointId)"; + public string AddCheckpointSql { get; } = $"insert into {schema}.checkpoints (id) values (@checkpointId)"; + public string UpdateCheckpointSql { get; } = $"update {schema}.checkpoints set position=(@position) where id=(@checkpointId)"; + public string TryInsertTombstone { get; } = $"select {schema}.try_insert_tombstone(@_gap_position, @_stream_name, @_type, @_id)"; static readonly Assembly Assembly = typeof(Schema).Assembly; diff --git a/src/Postgres/test/Eventuous.Tests.Postgres/Subscriptions/SchemaSqlTests.cs b/src/Postgres/test/Eventuous.Tests.Postgres/Subscriptions/SchemaSqlTests.cs new file mode 100644 index 000000000..4c4ff5611 --- /dev/null +++ b/src/Postgres/test/Eventuous.Tests.Postgres/Subscriptions/SchemaSqlTests.cs @@ -0,0 +1,31 @@ +using Eventuous.Postgresql; + +namespace Eventuous.Tests.Postgres.Subscriptions; + +/// +/// The queries are read on every poll, append and read, so the schema builds their text once rather than per access. +/// +public class SchemaSqlTests { + [Test] + public async Task Queries_are_built_once_for_the_schema() { + var schema = new Schema("custom_schema"); + + await Assert.That(schema.ReadAllForwards).IsEqualTo("select * from custom_schema.read_all_forwards(@_from_position, @_count)"); + await Assert.That(schema.ReadStreamSub).IsEqualTo("select * from custom_schema.read_stream_sub(@_stream_id, @_stream_name, @_from_position, @_count)"); + await Assert.That(schema.AppendEvents).IsEqualTo("select * from custom_schema.append_events(@_stream_name, @_expected_version, @_created, @_messages)"); + + await Assert.That(schema.StreamMessage).IsSameReferenceAs(schema.StreamMessage); + await Assert.That(schema.AppendEvents).IsSameReferenceAs(schema.AppendEvents); + await Assert.That(schema.ReadStreamForwards).IsSameReferenceAs(schema.ReadStreamForwards); + await Assert.That(schema.ReadStreamBackwards).IsSameReferenceAs(schema.ReadStreamBackwards); + await Assert.That(schema.ReadStreamSub).IsSameReferenceAs(schema.ReadStreamSub); + await Assert.That(schema.ReadAllForwards).IsSameReferenceAs(schema.ReadAllForwards); + await Assert.That(schema.CheckStream).IsSameReferenceAs(schema.CheckStream); + await Assert.That(schema.StreamExists).IsSameReferenceAs(schema.StreamExists); + await Assert.That(schema.TruncateStream).IsSameReferenceAs(schema.TruncateStream); + await Assert.That(schema.GetCheckpointSql).IsSameReferenceAs(schema.GetCheckpointSql); + await Assert.That(schema.AddCheckpointSql).IsSameReferenceAs(schema.AddCheckpointSql); + await Assert.That(schema.UpdateCheckpointSql).IsSameReferenceAs(schema.UpdateCheckpointSql); + await Assert.That(schema.TryInsertTombstone).IsSameReferenceAs(schema.TryInsertTombstone); + } +} diff --git a/src/Redis/src/Eventuous.Redis/Subscriptions/RedisAllStreamSubscription.cs b/src/Redis/src/Eventuous.Redis/Subscriptions/RedisAllStreamSubscription.cs index e66edc43e..bdba6e521 100644 --- a/src/Redis/src/Eventuous.Redis/Subscriptions/RedisAllStreamSubscription.cs +++ b/src/Redis/src/Eventuous.Redis/Subscriptions/RedisAllStreamSubscription.cs @@ -28,7 +28,8 @@ protected override async Task ReadEvents(IDatabase database, lo var stream = linkEvent[EventuousRedisKeys.Stream]; var streamPosition = linkEvent[Position]; - var streamEvents = await database.StreamRangeAsync(new(stream), streamPosition).NoContext(); + // The link holds the exact entry id, so read that one entry rather than the rest of the source stream + var streamEvents = await database.StreamRangeAsync(new(stream), minId: streamPosition, maxId: streamPosition, count: 1).NoContext(); var entry = streamEvents[0]; persistentEvents.Add( diff --git a/src/Redis/src/Eventuous.Redis/Subscriptions/RedisSubscriptionBase.cs b/src/Redis/src/Eventuous.Redis/Subscriptions/RedisSubscriptionBase.cs index 2d306aa94..7db387197 100644 --- a/src/Redis/src/Eventuous.Redis/Subscriptions/RedisSubscriptionBase.cs +++ b/src/Redis/src/Eventuous.Redis/Subscriptions/RedisSubscriptionBase.cs @@ -32,8 +32,6 @@ public abstract class RedisSubscriptionBase( metadataSerializer ) where T : RedisSubscriptionBaseOptions { - readonly IMetadataSerializer _metaSerializer = DefaultMetadataSerializer.Instance; - protected GetRedisDatabase GetDatabase { get; } = Ensure.NotNull(getDatabase, "Connection factory"); protected override async ValueTask Connect(SubscriptionRun run) { @@ -102,7 +100,10 @@ MessageConsumeContext ToConsumeContext(SubscriptionRun run, ReceivedEvent evt, C (ulong)evt.StreamPosition ); - var meta = (evt.JsonMetadata == null) ? new() : _metaSerializer.Deserialize(Encoding.UTF8.GetBytes(evt.JsonMetadata)); + // A payload-less context is acknowledged without entering the pipe, so its metadata would never be read + var meta = data is null ? null + : evt.JsonMetadata == null ? new() + : MetadataSerializer.DeserializeMeta(Options, Encoding.UTF8.GetBytes(evt.JsonMetadata), evt.StreamName, (ulong)evt.StreamPosition); return AsContext(run, evt, data, meta, cancellationToken); } diff --git a/src/Redis/test/Eventuous.Tests.Redis/Subscriptions/AllStreamLinkResolutionTests.cs b/src/Redis/test/Eventuous.Tests.Redis/Subscriptions/AllStreamLinkResolutionTests.cs new file mode 100644 index 000000000..e87a9acd9 --- /dev/null +++ b/src/Redis/test/Eventuous.Tests.Redis/Subscriptions/AllStreamLinkResolutionTests.cs @@ -0,0 +1,106 @@ +using System.Globalization; +using System.Reflection; +using Eventuous.Redis.Subscriptions; +using Eventuous.Subscriptions.Checkpoints; +using Eventuous.Subscriptions.Filters; +using Shouldly; +using StackExchange.Redis; +using Eventuous.Redis; +using static Eventuous.Redis.EventuousRedisKeys; + +namespace Eventuous.Tests.Redis.Subscriptions; + +/// +/// How the $all subscription resolves a link in _all to the entry in its source stream. The database is a proxy +/// that records the calls, so no server is involved. +/// +public class AllStreamLinkResolutionTests { + /// + /// Each link names one entry id. Reading from that id with no upper bound pulled the whole rest of the source + /// stream for every linked event and kept only the first entry. + /// + [Test] + public async Task Resolves_each_link_by_reading_only_the_linked_entry() { + const string stream = "booking-1"; + + var database = RecordingDatabase.Create( + [Link(stream, "1000-0"), Link(stream, "1001-0")], + [Entry("1000-0"), Entry("1001-0")] + ); + + var events = await new Probe().Read(database, 0); + + events.Select(x => x.StreamPosition).ShouldBe([10000L, 10010L]); + + ((RecordingDatabase)(object)database).Ranges + .Select(r => $"{r.Key} {r.MinId} {r.MaxId} {r.Count}") + .ShouldBe([$"{stream} 1000-0 1000-0 1", $"{stream} 1001-0 1001-0 1"], "each range must start and end at the linked entry"); + } + + static StreamEntry Link(string stream, string position) => new(RedisValue.Null, [new(EventuousRedisKeys.Stream, stream), new(Position, position)]); + + static StreamEntry Entry(string id) + => new( + id, + [ + new(MessageId, Guid.NewGuid().ToString()), + new(MessageType, "test-event"), + new(JsonData, "{}"), + new(JsonMetadata, "{}"), + new(Created, DateTime.UtcNow.ToString("o", CultureInfo.InvariantCulture)) + ] + ); + + sealed class Probe() + : RedisAllStreamSubscription( + () => null!, + new() { SubscriptionId = "redis-all-links" }, + new NoOpCheckpointStore(), + new ConsumePipe(), + null + ) { + public Task Read(IDatabase database, long position) => ReadEvents(database, position); + } + + public record Range(RedisKey Key, RedisValue? MinId, RedisValue? MaxId, int? Count); + + /// + /// Serves _all reads from the given links and range reads from the given entries by id, recording each range read. + /// + public class RecordingDatabase : DispatchProxy { + StreamEntry[] _links = []; + StreamEntry[] _entries = []; + + public List Ranges { get; } = []; + + public static IDatabase Create(StreamEntry[] links, StreamEntry[] entries) { + var proxy = Create(); + var self = (RecordingDatabase)(object)proxy; + self._links = links; + self._entries = entries; + + return proxy; + } + + protected override object? Invoke(MethodInfo? method, object?[]? args) { + switch (method?.Name) { + case nameof(IDatabaseAsync.StreamReadAsync): + return Task.FromResult(_links); + case nameof(IDatabaseAsync.StreamRangeAsync): { + var range = new Range((RedisKey)args![0]!, (RedisValue?)args[1], (RedisValue?)args[2], (int?)args[3]); + Ranges.Add(range); + + // Mimic XRANGE: entries from minId, up to maxId when given, at most count + var result = _entries + .Where(x => string.CompareOrdinal(x.Id, range.MinId ?? "-") >= 0) + .Where(x => range.MaxId is null || string.CompareOrdinal(x.Id, range.MaxId) <= 0) + .Take(range.Count ?? int.MaxValue) + .ToArray(); + + return Task.FromResult(result); + } + default: throw new NotSupportedException(method?.Name); + } + } + } +} diff --git a/src/Redis/test/Eventuous.Tests.Redis/Subscriptions/MetadataDeserializationTests.cs b/src/Redis/test/Eventuous.Tests.Redis/Subscriptions/MetadataDeserializationTests.cs new file mode 100644 index 000000000..71e22832d --- /dev/null +++ b/src/Redis/test/Eventuous.Tests.Redis/Subscriptions/MetadataDeserializationTests.cs @@ -0,0 +1,207 @@ +using System.Collections.Concurrent; +using System.Text.Json; +using Eventuous.Redis.Subscriptions; +using Eventuous.Subscriptions; +using Eventuous.Subscriptions.Checkpoints; +using Eventuous.Subscriptions.Context; +using Eventuous.Subscriptions.Filters; +using Shouldly; +using StackExchange.Redis; + +namespace Eventuous.Tests.Redis.Subscriptions; + +/// +/// How the Redis subscription turns a received entry into a consume context: which metadata serializer it uses and +/// what a metadata failure does to the poll loop. The entries come from the ReadEvents seam, so no server is involved. +/// +public class MetadataDeserializationTests { + const string KnownType = "redis-meta-test-event"; + const string UnknownType = "redis-meta-test-unknown"; + + // Multibyte characters on purpose: the payload and metadata are transcoded from strings to UTF-8 bytes. + const string Text = "héllo wörld ☕ 🌍"; + + /// + /// The serializer given to the subscription is the one that reads metadata; it used to be ignored in favour of the default. + /// + [Test] + [Timeout(30_000)] + public async Task Uses_the_configured_metadata_serializer(CancellationToken ct) { + var subscription = new ScriptedSubscription( + "redis-meta-configured", + [Event(KnownType, $$"""{"text":"{{Text}}"}""", $$"""{"note":"{{Text}}"}""")], + new MarkingMetadataSerializer() + ); + + await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 1, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("the event should have reached the handler"); + var context = subscription.Collector.Contexts.Single(); + context.Message.ShouldBeOfType().Text.ShouldBe(Text); + context.Metadata.ShouldNotBeNull(); + context.Metadata.GetString(MarkingMetadataSerializer.Key).ShouldBe(MarkingMetadataSerializer.Value, "the configured serializer should have read the metadata"); + context.Metadata.GetString("note").ShouldBe(Text); + } + + /// + /// Without ThrowOnError a malformed metadata entry is logged and skipped; it must not fault the poll loop, + /// which would otherwise fail on the same entry after every resubscribe. + /// + [Test] + [Timeout(30_000)] + public async Task Malformed_metadata_does_not_drop_the_subscription(CancellationToken ct) { + var subscription = new ScriptedSubscription( + "redis-meta-malformed", + [Event(KnownType, $$"""{"text":"{{Text}}"}""", "not json")] + ); + + var drops = 0; + await subscription.Subscribe(_ => { }, (_, _, _) => Interlocked.Increment(ref drops), ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 1, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("the event should still reach the handler"); + subscription.Collector.Contexts.Single().Metadata.ShouldBeNull(); + drops.ShouldBe(0); + } + + /// + /// Under ThrowOnError malformed metadata is still a fault, reported as a deserialization failure. + /// + [Test] + [Timeout(30_000)] + public async Task Malformed_metadata_drops_the_subscription_under_ThrowOnError(CancellationToken ct) { + var subscription = new ScriptedSubscription( + "redis-meta-malformed-throw", + [Event(KnownType, $$"""{"text":"{{Text}}"}""", "not json")], + throwOnError: true + ); + + Exception? reported = null; + await subscription.Subscribe(_ => { }, (_, _, e) => reported ??= e, ct); + var dropped = await WaitUntil(() => reported != null, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + dropped.ShouldBeTrue("malformed metadata under ThrowOnError should drop the subscription"); + reported.ShouldBeOfType(); + subscription.Collector.Contexts.ShouldBeEmpty(); + } + + /// + /// An event whose payload can't be deserialized is acknowledged without being handled, so its metadata is never read, + /// and malformed metadata on it doesn't fault the subscription even under ThrowOnError. + /// + [Test] + [Timeout(30_000)] + public async Task Metadata_of_a_payload_less_event_is_not_read(CancellationToken ct) { + var metaSerializer = new MarkingMetadataSerializer(); + + var subscription = new ScriptedSubscription( + "redis-meta-payload-less", + [ + Event(UnknownType, """{"text":"unknown"}""", "not json", position: 1), + Event(KnownType, $$"""{"text":"{{Text}}"}""", "{}", position: 2) + ], + metaSerializer, + throwOnError: true + ); + + var drops = 0; + await subscription.Subscribe(_ => { }, (_, _, _) => Interlocked.Increment(ref drops), ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 1, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("the event after the payload-less one should reach the handler"); + subscription.Collector.Contexts.Single().Message.ShouldBeOfType(); + drops.ShouldBe(0); + metaSerializer.Calls.ShouldBe(1, "only the handled event's metadata should be deserialized"); + } + + static ReceivedEvent Event(string type, string data, string? meta, long position = 1) + => new(Guid.NewGuid(), type, position, position, data, meta, DateTime.UtcNow, "redis-meta-stream"); + + static async Task WaitUntil(Func condition, TimeSpan timeout) { + var deadline = DateTime.UtcNow + timeout; + + while (DateTime.UtcNow < deadline) { + if (condition()) return true; + + await Task.Delay(20); + } + + return condition(); + } + + record TestEvent(string Text); + + record TestOptions : RedisSubscriptionBaseOptions; + + static IEventSerializer CreateSerializer() { + var typeMapper = new TypeMapper(); + typeMapper.AddType(KnownType); + + return new DefaultEventSerializer(new(JsonSerializerDefaults.Web), typeMapper); + } + + /// + /// Returns the scripted entries on the first poll and nothing afterwards. The database is never touched. + /// + sealed class ScriptedSubscription(string id, ReceivedEvent[] events, IMetadataSerializer? metaSerializer = null, bool throwOnError = false) + : RedisSubscriptionBase( + () => null!, + new() { SubscriptionId = id, ThrowOnError = throwOnError }, + new NoOpCheckpointStore(), + new ConsumePipe().AddDefaultConsumer(Handlers[id] = new()), + SubscriptionKind.All, + null, + CreateSerializer(), + metaSerializer + ) { + public CollectingHandler Collector => Handlers[id]; + + int _polls; + + protected override async Task ReadEvents(IDatabase database, long position) { + if (Interlocked.Increment(ref _polls) == 1) return events; + + await Task.Delay(20); + + return []; + } + } + + // The handler has to exist before the base constructor runs, so it is handed over through this map. + static readonly ConcurrentDictionary Handlers = new(); + + sealed class CollectingHandler : BaseEventHandler { + public ConcurrentQueue Contexts { get; } = new(); + + public override ValueTask HandleEvent(IMessageConsumeContext context) { + Contexts.Enqueue(context); + + return new(EventHandlingStatus.Success); + } + } + + /// + /// Reads metadata like the default serializer, then marks it so a test can tell which serializer ran. + /// + sealed class MarkingMetadataSerializer : IMetadataSerializer { + public const string Key = "serializer"; + public const string Value = "custom"; + + int _calls; + + public int Calls => Volatile.Read(ref _calls); + + public byte[] Serialize(Metadata evt) => DefaultMetadataSerializer.Instance.Serialize(evt); + + public Metadata? Deserialize(ReadOnlySpan bytes) { + Interlocked.Increment(ref _calls); + + return DefaultMetadataSerializer.Instance.Deserialize(bytes)?.With(Key, Value); + } + } +} diff --git a/src/Relational/src/Eventuous.Sql.Base/Subscriptions/SqlSubscriptionBase.cs b/src/Relational/src/Eventuous.Sql.Base/Subscriptions/SqlSubscriptionBase.cs index e742bbbd3..256abad28 100644 --- a/src/Relational/src/Eventuous.Sql.Base/Subscriptions/SqlSubscriptionBase.cs +++ b/src/Relational/src/Eventuous.Sql.Base/Subscriptions/SqlSubscriptionBase.cs @@ -40,8 +40,6 @@ public abstract class SqlSubscriptionBase( : EventSubscriptionWithCheckpoint(options, checkpointStore, consumePipe, concurrencyLimit, kind, loggerFactory, eventSerializer, metaSerializer), IMeasuredSubscription where TOptions : SqlSubscriptionOptionsBase where TConnection : DbConnection { - readonly IMetadataSerializer _metaSerializer = DefaultMetadataSerializer.Instance; - /// /// Create and open the SQL connection /// @@ -264,7 +262,10 @@ MessageConsumeContext ToConsumeContext(SubscriptionRun run, PersistedEvent evt, var data = DeserializeData(ContentType, evt.MessageType, Encoding.UTF8.GetBytes(evt.JsonData), evt.StreamName!, (ulong)evt.StreamPosition); - var meta = evt.JsonMetadata == null ? new() : _metaSerializer.Deserialize(Encoding.UTF8.GetBytes(evt.JsonMetadata!)); + // A payload-less context is acknowledged without entering the pipe, so its metadata would never be read + var meta = data is null ? null + : evt.JsonMetadata == null ? new() + : MetadataSerializer.DeserializeMeta(Options, Encoding.UTF8.GetBytes(evt.JsonMetadata), evt.StreamName!, (ulong)evt.StreamPosition); return AsContext(run, evt, data, meta, cancellationToken); } diff --git a/src/Sqlite/test/Eventuous.Tests.Sqlite/Subscriptions/MetadataDeserializationTests.cs b/src/Sqlite/test/Eventuous.Tests.Sqlite/Subscriptions/MetadataDeserializationTests.cs new file mode 100644 index 000000000..04ab5766c --- /dev/null +++ b/src/Sqlite/test/Eventuous.Tests.Sqlite/Subscriptions/MetadataDeserializationTests.cs @@ -0,0 +1,248 @@ +// Copyright (C) Eventuous HQ OÜ. All rights reserved +// Licensed under the Apache License, Version 2.0. + +using System.Collections.Concurrent; +using System.Text.Json; +using Eventuous.Sqlite.Subscriptions; +using Eventuous.Subscriptions; +using Eventuous.Subscriptions.Checkpoints; +using Eventuous.Subscriptions.Context; +using Eventuous.Subscriptions.Filters; +using Microsoft.Data.Sqlite; +using Shouldly; + +namespace Eventuous.Tests.Sqlite.Subscriptions; + +/// +/// How SqlSubscriptionBase turns a polled row into a consume context: which metadata serializer it uses and +/// what a metadata failure does to the poll loop. Sqlite shares that base with Postgres and SQL Server but needs no +/// container, so it is the cheapest place to pin it. +/// +/// The poll query selects literal rows from an in-memory database, so no schema or store is involved. +public class MetadataDeserializationTests { + const string KnownType = "sql-meta-test-event"; + const string UnknownType = "sql-meta-test-unknown"; + + // Multibyte characters on purpose: the payload and metadata are transcoded from strings to UTF-8 bytes. + const string Text = "héllo wörld ☕ 🌍"; + + /// + /// The serializer given to the subscription is the one that reads metadata; it used to be ignored in favour of the default. + /// + [Test] + [Timeout(30_000)] + public async Task Uses_the_configured_metadata_serializer(CancellationToken ct) { + var subscription = new ScriptedSubscription( + "sql-meta-configured", + [Row(1, KnownType, $$"""{"text":"{{Text}}"}""", $$"""{"note":"{{Text}}"}""")], + new MarkingMetadataSerializer() + ); + + await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 1, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("the event should have reached the handler"); + var context = subscription.Collector.Contexts.Single(); + context.Message.ShouldBeOfType().Text.ShouldBe(Text); + context.Metadata.ShouldNotBeNull(); + context.Metadata.GetString(MarkingMetadataSerializer.Key).ShouldBe(MarkingMetadataSerializer.Value, "the configured serializer should have read the metadata"); + context.Metadata.GetString("note").ShouldBe(Text); + } + + /// + /// Without ThrowOnError a malformed metadata row is logged and skipped; it must not fault the poll loop, + /// which would otherwise fail on the same row after every resubscribe. + /// + [Test] + [Timeout(30_000)] + public async Task Malformed_metadata_does_not_drop_the_subscription(CancellationToken ct) { + var subscription = new ScriptedSubscription("sql-meta-malformed", [Row(1, KnownType, $$"""{"text":"{{Text}}"}""", "not json")]); + + var drops = 0; + await subscription.Subscribe(_ => { }, (_, _, _) => Interlocked.Increment(ref drops), ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 1, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("the event should still reach the handler"); + subscription.Collector.Contexts.Single().Metadata.ShouldBeNull(); + drops.ShouldBe(0); + } + + /// + /// Under ThrowOnError malformed metadata is still a fault, reported as a deserialization failure. + /// + [Test] + [Timeout(30_000)] + public async Task Malformed_metadata_drops_the_subscription_under_ThrowOnError(CancellationToken ct) { + var subscription = new ScriptedSubscription( + "sql-meta-malformed-throw", + [Row(1, KnownType, $$"""{"text":"{{Text}}"}""", "not json")], + throwOnError: true + ); + + Exception? reported = null; + await subscription.Subscribe(_ => { }, (_, _, e) => reported ??= e, ct); + var dropped = await WaitUntil(() => reported != null, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + dropped.ShouldBeTrue("malformed metadata under ThrowOnError should drop the subscription"); + reported.ShouldBeOfType(); + subscription.Collector.Contexts.ShouldBeEmpty(); + } + + /// + /// An empty metadata string and the JSON literal null both mean "no metadata", not a failure, so neither + /// faults the subscription even under ThrowOnError. The empty string used to throw. + /// + [Test] + [Timeout(30_000)] + public async Task Empty_and_null_metadata_give_no_metadata(CancellationToken ct) { + var subscription = new ScriptedSubscription( + "sql-meta-empty", + [ + Row(1, KnownType, $$"""{"text":"{{Text}}"}""", ""), + Row(2, KnownType, $$"""{"text":"{{Text}}"}""", "null") + ], + throwOnError: true + ); + + var drops = 0; + await subscription.Subscribe(_ => { }, (_, _, _) => Interlocked.Increment(ref drops), ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 2, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("both events should reach the handler"); + subscription.Collector.Contexts.ShouldAllBe(x => x.Metadata == null); + drops.ShouldBe(0); + } + + /// + /// An event whose payload can't be deserialized is acknowledged without being handled, so its metadata is never read, + /// and malformed metadata on it doesn't fault the subscription even under ThrowOnError. + /// + [Test] + [Timeout(30_000)] + public async Task Metadata_of_a_payload_less_event_is_not_read(CancellationToken ct) { + var metaSerializer = new MarkingMetadataSerializer(); + + var subscription = new ScriptedSubscription( + "sql-meta-payload-less", + [ + Row(1, UnknownType, """{"text":"unknown"}""", "not json"), + Row(2, KnownType, $$"""{"text":"{{Text}}"}""", "{}") + ], + metaSerializer, + throwOnError: true + ); + + var drops = 0; + await subscription.Subscribe(_ => { }, (_, _, _) => Interlocked.Increment(ref drops), ct); + var received = await WaitUntil(() => subscription.Collector.Contexts.Count == 1, TimeSpan.FromSeconds(10)); + await subscription.Unsubscribe(_ => { }, ct); + + received.ShouldBeTrue("the event after the payload-less one should reach the handler"); + subscription.Collector.Contexts.Single().Message.ShouldBeOfType(); + drops.ShouldBe(0); + metaSerializer.Calls.ShouldBe(1, "only the handled event's metadata should be deserialized"); + } + + record EventRow(long Position, string Type, string Data, string Meta); + + static EventRow Row(long position, string type, string data, string meta) => new(position, type, data, meta); + + static async Task WaitUntil(Func condition, TimeSpan timeout) { + var deadline = DateTime.UtcNow + timeout; + + while (DateTime.UtcNow < deadline) { + if (condition()) return true; + + await Task.Delay(20); + } + + return condition(); + } + + record TestEvent(string Text); + + static IEventSerializer CreateSerializer() { + var typeMapper = new TypeMapper(); + typeMapper.AddType(KnownType); + + return new DefaultEventSerializer(new(JsonSerializerDefaults.Web), typeMapper); + } + + /// + /// Polls literal rows instead of the messages table, returning those past the subscription's position. + /// + sealed class ScriptedSubscription(string id, EventRow[] rows, IMetadataSerializer? metaSerializer = null, bool throwOnError = false) + : SqliteAllStreamSubscription( + new() { + SubscriptionId = id, + ConnectionString = "Data Source=:memory:", + ThrowOnError = throwOnError, + Polling = new() { MinIntervalMs = 1, MaxIntervalMs = 20 } + }, + new NoOpCheckpointStore(), + new ConsumePipe().AddDefaultConsumer(Collectors[id] = new()), + eventSerializer: CreateSerializer(), + metaSerializer: metaSerializer + ) { + public CollectingHandler Collector => Collectors[id]; + + protected override SqliteCommand PrepareCommand(SqliteConnection connection, long start) { + var cmd = connection.CreateCommand(); + + var values = string.Join( + " UNION ALL ", + rows.Select((_, i) => $"SELECT @id{i}, @type{i}, @pos{i}, @pos{i}, @data{i}, @meta{i}, '2026-01-01 00:00:00', 'sql-meta-stream'") + ); + + cmd.CommandText = $"WITH events(id, type, sp, gp, data, meta, created, stream) AS ({values}) SELECT * FROM events WHERE gp > @start ORDER BY gp"; + cmd.Parameters.AddWithValue("@start", start); + + for (var i = 0; i < rows.Length; i++) { + cmd.Parameters.AddWithValue($"@id{i}", Guid.NewGuid().ToString()); + cmd.Parameters.AddWithValue($"@type{i}", rows[i].Type); + cmd.Parameters.AddWithValue($"@pos{i}", rows[i].Position); + cmd.Parameters.AddWithValue($"@data{i}", rows[i].Data); + cmd.Parameters.AddWithValue($"@meta{i}", rows[i].Meta); + } + + return cmd; + } + } + + // The handler has to exist before the base constructor runs, so it is handed over through this map. + static readonly ConcurrentDictionary Collectors = new(); + + sealed class CollectingHandler : BaseEventHandler { + public ConcurrentQueue Contexts { get; } = new(); + + public override ValueTask HandleEvent(IMessageConsumeContext context) { + Contexts.Enqueue(context); + + return new(EventHandlingStatus.Success); + } + } + + /// + /// Reads metadata like the default serializer, then marks it so a test can tell which serializer ran. + /// + sealed class MarkingMetadataSerializer : IMetadataSerializer { + public const string Key = "serializer"; + public const string Value = "custom"; + + int _calls; + + public int Calls => Volatile.Read(ref _calls); + + public byte[] Serialize(Metadata evt) => DefaultMetadataSerializer.Instance.Serialize(evt); + + public Metadata? Deserialize(ReadOnlySpan bytes) { + Interlocked.Increment(ref _calls); + + return DefaultMetadataSerializer.Instance.Deserialize(bytes)?.With(Key, Value); + } + } +} From e38e910553875edbad2a6b0688cd51b2cc72ac15 Mon Sep 17 00:00:00 2001 From: Pavel Borisov Date: Thu, 1 Oct 2026 00:46:18 +0200 Subject: [PATCH 7/7] perf(subscriptions): skip metadata for payload-less events in KurrentDB and brokers - KurrentDB (all, stream, persistent), RabbitMQ, Service Bus and Pub/Sub subscriptions don't deserialize or build metadata when the payload didn't deserialize; such events are acknowledged without entering the pipe - Pub/Sub deserializes from the message memory without copying it - The Cloud Run receive log moves to Debug and logs the message id only Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Subscriptions/ServiceBusSubscription.cs | 3 +-- .../CloudRunPubSubSubscription.cs | 2 +- .../Subscriptions/GooglePubSubSubscription.cs | 4 ++-- .../Subscriptions/AllStreamSubscription.cs | 2 +- .../Subscriptions/PersistentSubscriptionBase.cs | 2 +- .../Subscriptions/StreamSubscription.cs | 14 ++++++++------ .../Subscriptions/RabbitMqSubscription.cs | 2 +- 7 files changed, 15 insertions(+), 14 deletions(-) diff --git a/src/Azure/src/Eventuous.Azure.ServiceBus/Subscriptions/ServiceBusSubscription.cs b/src/Azure/src/Eventuous.Azure.ServiceBus/Subscriptions/ServiceBusSubscription.cs index 59f50284c..ae0d801d6 100644 --- a/src/Azure/src/Eventuous.Azure.ServiceBus/Subscriptions/ServiceBusSubscription.cs +++ b/src/Azure/src/Eventuous.Azure.ServiceBus/Subscriptions/ServiceBusSubscription.cs @@ -110,7 +110,6 @@ CancellationToken ct Logger.Current = Log; var evt = DeserializeData(contentType, eventType, msg.Body, streamName); - var applicationProperties = msg.ApplicationProperties.Concat(MessageProperties(msg)); var ctx = new MessageConsumeContext( msg.MessageId, @@ -123,7 +122,7 @@ CancellationToken ct run.NextSequence(), msg.EnqueuedTime.UtcDateTime, evt, - AsMeta(applicationProperties), + evt is null ? null : AsMeta(msg.ApplicationProperties.Concat(MessageProperties(msg))), SubscriptionId, ct ); diff --git a/src/GooglePubSub/src/Eventuous.GooglePubSub.CloudRun/CloudRunPubSubSubscription.cs b/src/GooglePubSub/src/Eventuous.GooglePubSub.CloudRun/CloudRunPubSubSubscription.cs index 4fe1bf6de..5016d516d 100644 --- a/src/GooglePubSub/src/Eventuous.GooglePubSub.CloudRun/CloudRunPubSubSubscription.cs +++ b/src/GooglePubSub/src/Eventuous.GooglePubSub.CloudRun/CloudRunPubSubSubscription.cs @@ -47,7 +47,7 @@ public static void MapSubscription(WebApplication app, string path = "/") { return Results.BadRequest(); } - subscription.Log.InfoLog?.Log("Received {@Message}", envelope.Message); + subscription.Log.DebugLog?.Log("Received message {MessageId}", envelope.Message.MessageId); var data = Convert.FromBase64String(envelope.Message.Data); // ReSharper disable once ConditionIsAlwaysTrueOrFalseAccordingToNullableAPIContract diff --git a/src/GooglePubSub/src/Eventuous.GooglePubSub/Subscriptions/GooglePubSubSubscription.cs b/src/GooglePubSub/src/Eventuous.GooglePubSub/Subscriptions/GooglePubSubSubscription.cs index 377b120f2..0837e4624 100644 --- a/src/GooglePubSub/src/Eventuous.GooglePubSub/Subscriptions/GooglePubSubSubscription.cs +++ b/src/GooglePubSub/src/Eventuous.GooglePubSub/Subscriptions/GooglePubSubSubscription.cs @@ -138,7 +138,7 @@ async Task Handle(PubsubMessage msg, CancellationToken ct) { Logger.Current = Log; - var evt = DeserializeData(contentType, eventType, msg.Data.ToByteArray(), _topicName.TopicId); + var evt = DeserializeData(contentType, eventType, msg.Data.Memory, _topicName.TopicId); var ctx = new MessageConsumeContext( msg.MessageId, @@ -151,7 +151,7 @@ async Task Handle(PubsubMessage msg, CancellationToken ct) { run.NextSequence(), msg.PublishTime.ToDateTime(), evt, - AsMeta(msg.Attributes), + evt is null ? null : AsMeta(msg.Attributes), SubscriptionId, ct ); diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs index 8106abd9c..c089fc69e 100644 --- a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs +++ b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs @@ -242,7 +242,7 @@ MessageConsumeContext CreateContext(SubscriptionRun run, ResolvedEvent re, Cance run.NextSequence(), re.Event.Created, evt, - MetadataSerializer.DeserializeMeta(Options, re.Event.Metadata, re.Event.EventStreamId), + evt is null ? null : MetadataSerializer.DeserializeMeta(Options, re.Event.Metadata, re.Event.EventStreamId), SubscriptionId, cancellationToken ); diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs index b2707181b..a3c10b44e 100644 --- a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs +++ b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs @@ -216,7 +216,7 @@ MessageConsumeContext CreateContext(SubscriptionRun run, ResolvedEvent re, Cance run.NextSequence(), re.Event.Created, evt, - MetadataSerializer.DeserializeMeta(Options, re.Event.Metadata, re.Event.EventStreamId, re.Event.EventNumber), + evt is null ? null : MetadataSerializer.DeserializeMeta(Options, re.Event.Metadata, re.Event.EventStreamId, re.Event.EventNumber), SubscriptionId, cancellationToken ); diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs index d32e18b2d..8dfd106e3 100644 --- a/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs +++ b/src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs @@ -137,12 +137,14 @@ MessageConsumeContext CreateContext(SubscriptionRun run, ResolvedEvent re, Cance re.Event.EventNumber ); - var meta = MetadataSerializer.DeserializeMeta( - Options, - re.Event.Metadata, - re.Event.EventStreamId, - re.Event.EventNumber - ); + var meta = evt is null + ? null + : MetadataSerializer.DeserializeMeta( + Options, + re.Event.Metadata, + re.Event.EventStreamId, + re.Event.EventNumber + ); return new( re.Event.EventId.ToString(), diff --git a/src/RabbitMq/src/Eventuous.RabbitMq/Subscriptions/RabbitMqSubscription.cs b/src/RabbitMq/src/Eventuous.RabbitMq/Subscriptions/RabbitMqSubscription.cs index 9c13de196..e42cc264f 100644 --- a/src/RabbitMq/src/Eventuous.RabbitMq/Subscriptions/RabbitMqSubscription.cs +++ b/src/RabbitMq/src/Eventuous.RabbitMq/Subscriptions/RabbitMqSubscription.cs @@ -225,7 +225,7 @@ void LogDeliveryUndecided(Exception e) MessageConsumeContext CreateContext(BasicDeliverEventArgs received, CancellationToken cancellationToken) { var evt = DeserializeData(received.BasicProperties.ContentType!, received.BasicProperties.Type!, received.Body, received.Exchange); - var meta = received.BasicProperties.Headers != null + var meta = evt is not null && received.BasicProperties.Headers != null ? new Metadata(received.BasicProperties.Headers.ToDictionary(x => x.Key, x => x.Value)!) : null;