Skip to content

perf(subscriptions): hot-path quick wins and bug fixes - #601

Open
Inok wants to merge 7 commits into
Eventuous:devfrom
Inok:perf/hot-path-quick-wins
Open

Inok wants to merge 7 commits into
Eventuous:devfrom
Inok:perf/hot-path-quick-wins

Conversation

@Inok

@Inok Inok commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Summary

This PR contains low-risk bug fixes and allocation cuts on the subscription consume path, plus a few from diagnostics, application and persistence. It has no redesigns and no public signature changes. Each commit covers one area and builds on its own.

Commit What
fix(subscriptions) MurmurHash3 no longer reads an unpinned string. The context converter caches are thread-safe. TracingFilter disposes only the activities it created and keeps a failed handler's error status. ConsumePipe validates filter types when the pipe is composed. Items dictionary is allocated lazily; logging scope is an array; ack/nack delegates are built once per run.
fix(diagnostics) Metrics come from every traced handler and command service, not only the last one created. Measures are skipped when nobody listens. Activity names are cached; activity helpers don't allocate per call.
fix(application) ThrowingCommandService returns the result on success. Failure-only strings are built only on failure.
perf(domain) Aggregate.Changes doesn't allocate a new view on each access.
perf(serialization) Failed-deserialization results are shared instances.
fix(sql) SQL and Redis subscriptions use the configured metadata serializer and handle bad metadata like KurrentDB does. Redis $all link resolution reads exactly one entry. Postgres builds its query text once per schema instead of on every access.
perf(subscriptions) KurrentDB, RabbitMQ, Service Bus and Pub/Sub skip metadata for events whose payload didn't deserialize. Pub/Sub reads the message body without copying it.

Observable behaviour changes

  1. Subscription spans on the synchronous path. This covers Service Bus, Pub/Sub, Cloud Run and KurrentDB persistent subscriptions. TracingFilter used to dispose the subscription activity it reused, which stopped the span before EventSubscription.Handler had finished with it. Now the handler ends the span:

    • A failed message's span is marked Recorded before it stops, so failed spans that a sampler had dropped are now exported.
    • TracingFilter no longer overwrites the error status a failed handler set (via Nack) with OK. This also corrects failed spans that were already exported, which used to show OK.
    • The span ends slightly later. It now covers the subscription's post-processing, but not the transport's ack or complete.
    • When Eventuous tracing has no listener, the filter no longer stops a foreign ambient span (the Azure SDK processing span, or the ASP.NET request span) early.

    Ignored messages are unchanged: their span was already marked not-recorded before it stopped. Checkpointed subscriptions (the async path) are unaffected, because there the filter always starts and disposes its own span.

  2. Metrics from every instance. Duration and error metrics now come from every traced handler and command service. Before, only the most recently created one was observed.

    • Expect more metric series.
    • In a process with several hosts or meter providers, each provider receives every instance's measurements.
    • Durations use Stopwatch, so they are monotonic.
    • A DiagnosticListener observer that subscribes with a predicate rejecting the measure event no longer receives it. The built-in metrics subscribe without a predicate.
  3. SQL and Redis metadata. Postgres, SQL Server and SQLite subscriptions now use the metadata serializer they are given, including one registered in DI. Before, it was injected but ignored. The concrete Redis subscriptions don't take one and keep the default. A RedisSubscriptionBase subclass that passes one now has it honoured.

    Before, malformed metadata threw out of the poll loop, whatever ThrowOnError was set to. The subscription dropped and resubscribed from the same checkpoint indefinitely. Now:

    • Without ThrowOnError, the error is logged, the event is delivered with Metadata == null, and the checkpoint moves on. Handlers that dereference metadata without a null check now run for such events.
    • With ThrowOnError, the loop behaves as before. The exception is now a DeserializationException, and an error is logged on each attempt.
    • An empty metadata string gives null metadata instead of an error. A whitespace-only string is still malformed.
  4. No metadata on payload-less events. This applies when the payload didn't deserialize: it was empty, its type isn't registered, or deserialization failed and the failure was caught. On every transport the context then has Metadata == null. Nothing reads it, because such contexts are ignored and acknowledged without entering the pipe or starting an activity. The visible effects:

    • Malformed metadata on such an event is no longer logged. It also no longer fails the subscription: under ThrowOnError on every transport, and on SQL and Redis in either mode.
    • Service Bus: suppose a message's application properties collide with the configured attribute names and its payload doesn't deserialize. It used to be abandoned, and eventually dead-lettered. It is now completed like any other ignored message.
  5. ConsumePipe checks filter types when it is built. An incompatible second filter throws InvalidContextTypeException while the pipe is composed. Before, a ConsumeFilter<,>-derived first filter already threw ArgumentException on every message for the same mismatch. The only composition that used to work and is now rejected: a first filter that implements IConsumeFilter<,> directly (skipping that per-message check), declares a TOut the second filter can't consume, yet passes it compatible contexts at runtime.

  6. ThrowingCommandService returns the result on success. It used to always throw ApplicationException, a regression from d0499d4 (2024). This restores the earlier behaviour.

  7. Redis $all resolves each link with a single-entry XRANGE, which picks the same entry as before whenever it exists. A linked entry can only go missing through an external XDEL or XTRIM.

    • Now a missing entry throws. The subscription then drops and keeps resubscribing at that page.
    • Before, it delivered the next entry of the source stream instead, which meant a wrong event plus a duplicate. If there was no next entry, it threw the same way.
  8. Cloud Run Pub/Sub receive log moved from Info to Debug. It logs the message id only, and no longer the payload and attributes.

  9. Shared instances. ActivityStatus.Ok() and FailedToDeserialize results are now shared instances, and Aggregate.Changes returns the same read-only view on every access. This is only visible through reference equality.

  10. DefaultConsumer logging scope state is a KeyValuePair array instead of a dictionary, with the same keys and values.

Public API

  • No signature changes.
  • BaseTracer.StartMeasure went from private to private protected.

Testing

On net10.0, locally, these suites pass:

Suite Passed
Core 36
Subscriptions 163
Application 24
Diagnostics 13
Gateway 10
DI 3
SQLite 53
KurrentDB 83
Postgres 56
Redis 34
RabbitMQ 5
MongoDB 9
Pub/Sub 3
Service Bus 44 of 50 (6 skipped by design)

Not verified yet:

  • net8.0 and net9.0: compiled only, tests not run.
  • The SQL Server suite: excluded on macOS.
  • The span-export change: reasoned from source, not observed in a collector.
  • No benchmark numbers yet. I'll add before/after numbers from src/Benchmarks before marking this ready.

🤖 Generated with Claude Code

@github-actions

github-actions Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Test Results

   48 files  ±  0     48 suites  ±0   16m 5s ⏱️ +50s
  663 tests + 50    663 ✅ + 50  0 💤 ±0  0 ❌ ±0 
1 358 runs  +132  1 358 ✅ +132  0 💤 ±0  0 ❌ ±0 

Results for commit e38e910. ± Comparison against base commit ed6cfd1.

This pull request removes 9 and adds 59 tests. Note that renamed tests count towards both.
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(09/21/2026 13:40:13 +00:00)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(09/21/2026 13:40:13)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(9a8c03f6-9bdf-4072-a43a-772bdae7bb21)
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 4, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 8, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 4, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 8, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 4, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 8, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 })
Eventuous.Tests.Aggregates.ChangesViewTests ‑ ShouldReuseTheSameViewThatReflectsNewChanges
Eventuous.Tests.Application.ResolverNullCheckTests ‑ ShouldFailAtCallTimeWhenNoReaderIsAvailable
Eventuous.Tests.Application.ThrowingCommandServiceTests ‑ ShouldReturnResultOnSuccess
Eventuous.Tests.Application.ThrowingCommandServiceTests ‑ ShouldThrowOnError
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(10/01/2026 10:17:54 +00:00)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(10/01/2026 10:17:54)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(126757bc-8dcf-46dd-96cf-865ae4f38afb)
Eventuous.Tests.DeserializationFailureTests ‑ ShouldReuseContentTypeMismatchFailure
Eventuous.Tests.DeserializationFailureTests ‑ ShouldReusePayloadEmptyFailure
Eventuous.Tests.DeserializationFailureTests ‑ ShouldReuseUnknownTypeFailure
…

♻️ This comment has been updated with latest results.

@Inok
Inok force-pushed the perf/hot-path-quick-wins branch 3 times, most recently from 51f5c55 to 327fd91 Compare September 30, 2026 23:43
@Inok
Inok marked this pull request as ready for review October 1, 2026 09:57
@Inok
Inok force-pushed the perf/hot-path-quick-wins branch from 327fd91 to 7207869 Compare October 1, 2026 09:58
@qodo-free-for-open-source-projects

Copy link
Copy Markdown
Contributor

PR Summary by Qodo

Fix subscription consume-path bugs and reduce hot-path allocations

🐞 Bug fix ✨ Enhancement 🧪 Tests 🕐 40+ Minutes

Grey Divider

AI Description

• Fix tracing ownership, metric coverage, command results, and metadata error handling.
• Skip unused metadata and reduce allocations across subscriptions, serializers, and persistence.
• Add regression tests for consume behavior, diagnostics, serialization, SQL, and Redis.
Diagram

graph TD
  T["Transport subscriptions"] --> D["Deserialize payload"] --> V{"Payload available?"} -->|yes| M["Deserialize metadata"] --> C["Consume context"] --> P["Consume pipe"] --> H["Handlers and metrics"]
  V -->|no| I["Ignore and acknowledge"]
Loading
High-Level Assessment

Keep the targeted fixes in existing subscription and diagnostics abstractions rather than introducing a new pipeline. Review the observable behavior changes—especially malformed metadata, span export, and unsnapshotted event collections—before relying on the allocation improvements.

Files changed (57) +1615 / -150

Enhancement (20) +130 / -66
CommandHandlerBuilder.csDefer missing-store error message construction +2/-2

Defer missing-store error message construction

• Reader and writer resolvers construct their ArgumentNullException messages only when the dependency is missing.

src/Core/src/Eventuous.Application/AggregateService/CommandHandlerBuilder.cs

ApplicationEventSource.csAvoid formatting disabled command error events +6/-2

Avoid formatting disabled command error events

• Command names and exception text are produced only when the corresponding error event is enabled.

src/Core/src/Eventuous.Application/Diagnostics/ApplicationEventSource.cs

ActivityExtensions.csRead parent tags without enumeration +2/-1

Read parent tags without enumeration

• Parent-tag lookup uses GetTagItem while retaining the prior string-only behavior.

src/Core/src/Eventuous.Diagnostics/ActivityExtensions.cs

ActivityStatus.csShare plain successful activity status +4/-1

Share plain successful activity status

• Ok() returns a shared status when no description is supplied.

src/Core/src/Eventuous.Diagnostics/ActivityStatus.cs

Aggregate.csReuse the aggregate changes view +5/-2

Reuse the aggregate changes view

• Changes lazily retains one live read-only view, while CurrentVersion reads the backing list count directly.

src/Core/src/Eventuous.Domain/Aggregate.cs

BaseTracer.csSkip unobserved persistence measures +12/-6

Skip unobserved persistence measures

• Tracing paths create measure contexts only when their event is observed. A private protected helper centralizes the check.

src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/BaseTracer.cs

TracedEventWriter.csPass append collections through without copying +15/-11

Pass append collections through without copying

• Tracing enriches event metadata in place and passes the original event or append collection to the writer. It also skips unobserved measures.

src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs

DefaultEventSerializer.csReuse dynamic deserialization failures +8/-3

Reuse dynamic deserialization failures

• Unknown-type, content-type-mismatch, and empty-payload results are shared rather than allocated for each failure.

src/Core/src/Eventuous.Serialization.Json.Dynamic/DefaultEventSerializer.cs

DefaultStaticEventSerializer.csReuse static deserialization failures +8/-3

Reuse static deserialization failures

• The static serializer shares immutable results for its three routine failure kinds.

src/Core/src/Eventuous.Serialization/DefaultStaticEventSerializer.cs

DefaultConsumer.csUse an array for the consumer logging scope +4/-4

Use an array for the consumer logging scope

• The logging scope retains its keys and values without allocating a dictionary.

src/Core/src/Eventuous.Subscriptions/Consumers/DefaultConsumer.cs

MessageConsumeContext.csAllocate context items on first access +3/-1

Allocate context items on first access

• Contexts without item usage no longer allocate ContextItems; repeated access returns the same instance.

src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs

EventSubscription.csCache subscription activity names +11/-1

Cache subscription activity names

• Activity names are built once per message type and subscription instead of on every handled message.

src/Core/src/Eventuous.Subscriptions/EventSubscription.cs

EventSubscriptionWithCheckpoint.csReuse acknowledgement delegates per run +23/-4

Reuse acknowledgement delegates per run

• Each checkpointed run owns its ack and nack delegates rather than constructing closures for every message.

src/Core/src/Eventuous.Subscriptions/EventSubscriptionWithCheckpoint.cs

CloudRunPubSubSubscription.csLimit Cloud Run receive logging +1/-1

Limit Cloud Run receive logging

• The receive log moves to Debug and includes only the message ID instead of the full envelope.

src/GooglePubSub/src/Eventuous.GooglePubSub.CloudRun/CloudRunPubSubSubscription.cs

GooglePubSubSubscription.csAvoid Pub/Sub body copy and unused metadata +2/-2

Avoid Pub/Sub body copy and unused metadata

• Payload deserialization reads message memory directly; attributes become metadata only for deserialized events.

src/GooglePubSub/src/Eventuous.GooglePubSub/Subscriptions/GooglePubSubSubscription.cs

AllStreamSubscription.csSkip metadata for unusable all-stream events +1/-1

Skip metadata for unusable all-stream events

• Metadata is not deserialized when the event payload did not deserialize.

src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/AllStreamSubscription.cs

PersistentSubscriptionBase.csSkip metadata for unusable persistent events +1/-1

Skip metadata for unusable persistent events

• Persistent-subscription contexts have null metadata when their payload is absent.

src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/PersistentSubscriptionBase.cs

StreamSubscription.csSkip metadata for unusable stream events +8/-6

Skip metadata for unusable stream events

• Stream subscriptions deserialize metadata only after payload deserialization succeeds.

src/KurrentDB/src/Eventuous.KurrentDB/Subscriptions/StreamSubscription.cs

Schema.csBuild PostgreSQL query strings once per schema +13/-13

Build PostgreSQL query strings once per schema

• Query properties retain their strings instead of interpolating SQL on every access.

src/Postgres/src/Eventuous.Postgresql/Schema.cs

RabbitMqSubscription.csSkip RabbitMQ headers for unusable events +1/-1

Skip RabbitMQ headers for unusable events

• Headers are converted to metadata only if the payload deserializes.

src/RabbitMq/src/Eventuous.RabbitMq/Subscriptions/RabbitMqSubscription.cs

Bug fix (13) +75 / -71
ServiceBusSubscription.csSkip properties for payload-less Service Bus messages +1/-2

Skip properties for payload-less Service Bus messages

• Application properties become metadata only when the payload deserializes, so ignored messages are not failed by unusable metadata.

src/Azure/src/Eventuous.Azure.ServiceBus/Subscriptions/ServiceBusSubscription.cs

CommandServiceActivity.csShare command metrics source across services +10/-4

Share command metrics source across services

• All traced command services publish through one listener. Unobserved measures are skipped, with nullable handling on failures.

src/Core/src/Eventuous.Application/Diagnostics/CommandServiceActivity.cs

TracedCommandService.csRemove per-service metrics listener +1/-5

Remove per-service metrics listener

• Command tracing delegates metrics publication to the shared source in CommandServiceActivity.

src/Core/src/Eventuous.Application/Diagnostics/TracedCommandService.cs

ThrowingCommandService.csReturn successful command results +1/-1

Return successful command results

• The wrapper still throws for error results but returns the result when handling succeeds.

src/Core/src/Eventuous.Application/ThrowingCommandService.cs

Measure.csTime measures with a monotonic clock +2/-3

Time measures with a monotonic clock

• Duration measurement switches from wall-clock timestamps to Stopwatch.

src/Core/src/Eventuous.Diagnostics/Metrics/Measure.cs

MessageConsumeContextConverter.csMake context-conversion caches concurrency-safe +18/-25

Make context-conversion caches concurrency-safe

• Conversions use a concurrent cache, and registered converters are read from copy-on-write snapshots during concurrent registration.

src/Core/src/Eventuous.Subscriptions/Consumers/MessageConsumeContextConverter.cs

ConsumePipe.csValidate the second filter at composition +1/-1

Validate the second filter at composition

• A second filter whose input cannot accept the first filter's output now fails when added, before messages arrive.

src/Core/src/Eventuous.Subscriptions/Filters/ConsumePipe.cs

MurmurHash3.csHash partition keys without an unpinned pointer +3/-14

Hash partition keys without an unpinned pointer

• The hash reads bytes from the string's span, preserving hash values while removing unsafe pointer access.

src/Core/src/Eventuous.Subscriptions/Filters/Partitioning/MurmurHash3.cs

TracingFilter.csPreserve subscription span ownership and error status +12/-5

Preserve subscription span ownership and error status

• The filter disposes only spans it starts and does not replace a nacked handler's error status with OK.

src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs

TracedEventHandler.csObserve every traced handler and cache span names +16/-4

Observe every traced handler and cache span names

• Handlers share a metrics source, omit unobserved measures, and reuse activity names by message type.

src/Core/src/Eventuous.Subscriptions/Handlers/TracedEventHandler.cs

RedisAllStreamSubscription.csResolve Redis links with exact-entry ranges +2/-1

Resolve Redis links with exact-entry ranges

• Each $all link requests only its specified source entry rather than the rest of the stream.

src/Redis/src/Eventuous.Redis/Subscriptions/RedisAllStreamSubscription.cs

RedisSubscriptionBase.csHonor Redis metadata serializer and failure policy +4/-3

Honor Redis metadata serializer and failure policy

• Redis uses the configured serializer and shared metadata error handling, while skipping metadata for payload-less events.

src/Redis/src/Eventuous.Redis/Subscriptions/RedisSubscriptionBase.cs

SqlSubscriptionBase.csHonor SQL metadata serializer and failure policy +4/-3

Honor SQL metadata serializer and failure policy

• SQL subscriptions use the configured serializer and shared error handling, and do not read metadata for payload-less rows.

src/Relational/src/Eventuous.Sql.Base/Subscriptions/SqlSubscriptionBase.cs

Tests (23) +1410 / -12
ResolverNullCheckTests.csTest missing reader failure timing +23/-0

Test missing reader failure timing

• Verifies that a service without a configured reader fails on command handling with the expected error message.

src/Core/test/Eventuous.Tests.Application/ResolverNullCheckTests.cs

ThrowingCommandServiceTests.csTest command wrapper success and failure +26/-0

Test command wrapper success and failure

• Covers returning a successful result and throwing for an error result.

src/Core/test/Eventuous.Tests.Application/ThrowingCommandServiceTests.cs

ConsumePipeTests.csTest early filter incompatibility rejection +14/-0

Test early filter incompatibility rejection

• Confirms that adding an incompatible second filter throws without modifying the pipe.

src/Core/test/Eventuous.Tests.Subscriptions/ConsumePipeTests.cs

ContextItemsTests.csTest lazy context-items identity +23/-0

Test lazy context-items identity

• Checks repeated access and preservation of items when a context is wrapped.

src/Core/test/Eventuous.Tests.Subscriptions/ContextItemsTests.cs

DynamicSerializerFailureTests.csTest shared dynamic serializer failures +54/-0

Test shared dynamic serializer failures

• Checks result identity and error kinds for unknown types, mismatched content types, and empty payloads.

src/Core/test/Eventuous.Tests.Subscriptions/DynamicSerializerFailureTests.cs

MessageConsumeContextConverterTests.csExercise concurrent context conversion and registration +90/-0

Exercise concurrent context conversion and registration

• Runs conversions across many types while registering converters to smoke-test cache and snapshot safety.

src/Core/test/Eventuous.Tests.Subscriptions/MessageConsumeContextConverterTests.cs

MurmurHash3Tests.csPin partition-hash compatibility +23/-0

Pin partition-hash compatibility

• Verifies known hashes, including Unicode and empty keys, and checks null input handling.

src/Core/test/Eventuous.Tests.Subscriptions/MurmurHash3Tests.cs

TracedEventHandlerTests.csTest handler metrics and cached activity names +143/-0

Test handler metrics and cached activity names

• Confirms metrics arrive from multiple traced handlers and that handler and subscription span names remain unchanged.

src/Core/test/Eventuous.Tests.Subscriptions/TracedEventHandlerTests.cs

TracingFilterTests.csTest tracing-filter span ownership +101/-0

Test tracing-filter span ownership

• Checks that reused spans remain open, newly started spans stop, and nacked messages retain error status.

src/Core/test/Eventuous.Tests.Subscriptions/TracingFilterTests.cs

ChangesViewTests.csTest live, reusable aggregate changes view +26/-0

Test live, reusable aggregate changes view

• Checks view identity, reflected changes, and version calculations before and after clearing events.

src/Core/test/Eventuous.Tests/Aggregates/ChangesViewTests.cs

DeserializationFailureTests.csTest shared static serializer failures +53/-0

Test shared static serializer failures

• Checks identity and error kinds for the static serializer's routine failure results.

src/Core/test/Eventuous.Tests/DeserializationFailureTests.cs

TracedEventWriterTests.csTest append identity, tracing metadata, and metrics +100/-0

Test append identity, tracing metadata, and metrics

• Verifies that both append overloads preserve caller collections, enrich metadata, and produce observed metrics.

src/Core/test/Eventuous.Tests/TracedEventWriterTests.cs

ActivityHelpersTests.csTest shared status and parent-tag behavior +50/-0

Test shared status and parent-tag behavior

• Checks plain OK status identity and confirms that only string parent tags are copied.

src/Diagnostics/test/Eventuous.Tests.Diagnostics/ActivityHelpersTests.cs

TracedCommandServiceMetricsTests.csTest metrics from multiple command services +58/-0

Test metrics from multiple command services

• Verifies that services with different state types both contribute measurements.

src/Diagnostics/test/Eventuous.Tests.Diagnostics/TracedCommandServiceMetricsTests.cs

MetricsTests.csAdd reusable subscription-duration assertion +10/-12

Add reusable subscription-duration assertion

• Replaces a commented-out duration check with an executable assertion for subscription and message tags.

src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs

MetricsTests.csCheck KurrentDB subscription duration metric +6/-0

Check KurrentDB subscription duration metric

• Adds the shared duration-and-tags assertion to the KurrentDB metrics suite.

src/KurrentDB/test/Eventuous.Tests.KurrentDB/Metrics/MetricsTests.cs

MetricsTests.csCheck Postgres subscription duration metric +6/-0

Check Postgres subscription duration metric

• Adds the shared duration-and-tags assertion to the Postgres metrics suite.

src/Postgres/test/Eventuous.Tests.Postgres/Metrics/MetricsTests.cs

SchemaSqlTests.csTest cached PostgreSQL query text +31/-0

Test cached PostgreSQL query text

• Checks representative SQL values and reference identity across repeated property reads.

src/Postgres/test/Eventuous.Tests.Postgres/Subscriptions/SchemaSqlTests.cs

AllStreamLinkResolutionTests.csTest exact Redis link resolution +106/-0

Test exact Redis link resolution

• A recording database verifies that each link requests only its named entry and that the resolved events match.

src/Redis/test/Eventuous.Tests.Redis/Subscriptions/AllStreamLinkResolutionTests.cs

MetadataDeserializationTests.csTest Redis metadata handling through the poll loop +207/-0

Test Redis metadata handling through the poll loop

• Covers custom serializers, malformed metadata under both error policies, and metadata skipped for ignored payloads.

src/Redis/test/Eventuous.Tests.Redis/Subscriptions/MetadataDeserializationTests.cs

MetricsTests.csCheck SQL Server subscription duration metric +6/-0

Check SQL Server subscription duration metric

• Adds the shared duration-and-tags assertion to the SQL Server metrics suite.

src/SqlServer/test/Eventuous.Tests.SqlServer/Metrics/MetricsTests.cs

MetricsTests.csCheck SQLite subscription duration metric +6/-0

Check SQLite subscription duration metric

• Adds the shared duration-and-tags assertion to the SQLite metrics suite.

src/Sqlite/test/Eventuous.Tests.Sqlite/Metrics/MetricsTests.cs

MetadataDeserializationTests.csTest shared SQL metadata handling with SQLite +248/-0

Test shared SQL metadata handling with SQLite

• In-memory polling tests custom serializers, malformed and empty metadata, error-policy behavior, and skipped metadata for ignored events.

src/Sqlite/test/Eventuous.Tests.Sqlite/Subscriptions/MetadataDeserializationTests.cs

Other (1) +0 / -1
Eventuous.Subscriptions.csprojRemove unsafe-code build setting +0/-1

Remove unsafe-code build setting

• Unsafe blocks are no longer required after the partition hash switches to a span byte view.

src/Core/src/Eventuous.Subscriptions/Eventuous.Subscriptions.csproj

@qodo-free-for-open-source-projects

qodo-free-for-open-source-projects Bot commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Code Review by Qodo

🐞 Bugs (0) 📘 Rule violations (0) 📎 Requirement gaps (0) 🎨 UX issues (0) 🔗 Cross-repo conflicts (0) 📜 Skill insights (0)

Grey Divider


Remediation recommended

1. Batch appends can use changed inputs ✓ Resolved
Description
TracedEventWriter.AppendEvents now forwards the caller's appends collection and its event
collections rather than passing a materialized snapshot. If a caller changes either collection while
an append is in flight, the KurrentDB writer can enumerate changed request contents or map responses
against a different stream order after its network await.
Code

src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs[58]

+            var results = await writer.AppendEvents(appends, cancellationToken).NoContext();
Evidence
The writer interface accepts IReadOnlyCollection, which can have a mutable backing collection. The
traced writer now passes that collection through; the KurrentDB implementation builds deferred
requests from it, awaits the network append, and then enumerates it again to associate responses
with streams. The previous traced-writer path materialized both the outer collection and each event
collection.

src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs[51-58]
src/Core/src/Eventuous.Persistence/EventStore/IEventWriter.cs[29-39]
src/KurrentDB/src/Eventuous.KurrentDB/KurrentDBEventStore.cs[161-183]

Agent prompt
The issue below was found during a code review. Follow the provided context and guidance below and implement a solution

## Issue description
The traced writer no longer snapshots batch appends, so downstream asynchronous writers can observe changes to caller-owned collections during a write.
## Fix Focus Areas
- src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs[51-58]
## Recommended Fix
Materialize the outer append collection and each event collection before calling the inner writer, while avoiding unnecessary copies of individual event objects. Preserve tracing enrichment and the response order.

ⓘ Copy this prompt and use it to remediate the issue with your preferred AI generation tools


2. Parallel access can lose context items ✓ Resolved
Description
MessageConsumeContext.Items initializes _items with an unsynchronized ??=, allowing two
first-time callers to receive different ContextItems instances. If one caller adds an item while
the other replaces _items, later accesses through the context or its wrappers cannot see the item
added to the displaced instance.
Code

src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs[49]

+    public ContextItems      Items             => _items ??= new();
Evidence
The new getter can allocate and return an instance before another caller replaces the field. Wrapped
contexts forward item access to the same underlying message context, so they do not protect callers
from this race; AddItem writes to the particular instance it receives.

src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs[26-49]
src/Core/src/Eventuous.Subscriptions/Context/WrappedConsumeContext.cs[25-25]
src/Core/src/Eventuous.Subscriptions/Context/ContextItems.cs[18-21]

Agent prompt
The issue below was found during a code review. Follow the provided context and guidance below and implement a solution

## Issue description
Concurrent first accesses to a message context can publish different item bags, losing items written to the displaced bag.
## Fix Focus Areas
- src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs[26-49]
## Recommended Fix
Keep lazy allocation, but publish the newly created `ContextItems` with `Interlocked.CompareExchange` and always return the published instance. Add a concurrent first-access test.

ⓘ Copy this prompt and use it to remediate the issue with your preferred AI generation tools


Grey Divider

Tip of the day
💡 Did you know, you can keep summaries lean with Findings visible per group, which tucks the rest behind a View link

More tips ↗ | Customize Qodo ↗ | Qodo docs ↗

Grey Divider

Qodo Logo

Comment thread src/Core/src/Eventuous.Persistence/Diagnostics/Tracing/TracedEventWriter.cs Outdated
Comment thread src/Core/src/Eventuous.Subscriptions/Context/MessageConsumeContext.cs Outdated
Inok and others added 7 commits October 1, 2026 12:10
…tions

- 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) <noreply@anthropic.com>
…easures

- 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) <noreply@anthropic.com>
- 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) <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
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 <noreply@anthropic.com>
- 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) <noreply@anthropic.com>
…DB 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) <noreply@anthropic.com>
@Inok
Inok force-pushed the perf/hot-path-quick-wins branch from 7207869 to e38e910 Compare October 1, 2026 10:12
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant