diff --git a/decisions/2026-09-29-outbox-handler-runs-before-event-dispatcher.md b/decisions/2026-09-29-outbox-handler-runs-before-event-dispatcher.md new file mode 100644 index 00000000..d86e2d80 --- /dev/null +++ b/decisions/2026-09-29-outbox-handler-runs-before-event-dispatcher.md @@ -0,0 +1,58 @@ +--- +authors: + - Martin Stühmer + +applyTo: + - "src/NetEvolve.Pulse/Internals/PulseMediator.cs" + - "src/NetEvolve.Pulse/Dispatchers/*.cs" + - "src/NetEvolve.Pulse/Outbox/OutboxEventHandler*.cs" + +created: 2026-09-29 + +lastModified: 2026-09-29 + +state: proposed + +instructions: | + MUST invoke OutboxEventHandler in PulseMediator sequentially and before the configured IEventDispatcher, inside the event interceptor chain, and MUST pass only the remaining handlers to the dispatcher. + MUST keep the handler error contract: every handler runs, and failures are thrown afterwards as one flat AggregateException. + MUST keep ParallelEventDispatcher as the default and document that user handlers sharing a scoped DbContext or connection need SequentialEventDispatcher. +--- + +# Decision: The Outbox Handler Runs Before the Event Dispatcher + +`PulseMediator.PublishAsync` runs the framework's `OutboxEventHandler` on its own and in order, before the configured `IEventDispatcher` gets the other handlers. The outbox write never overlaps another handler in the caller's scope. + +## Context + +Since #599, event handlers are resolved from the caller's scope, so scoped handlers share the caller's `DbContext` and transaction. `AddOutbox()` registers `OutboxEventHandler<>` as an open generic for every event. With the Entity Framework outbox it calls `AddAsync` and `SaveChangesAsync` on the scoped `TContext`. + +The default `ParallelEventDispatcher` ran all resolved handlers through `Parallel.ForEachAsync`. The outbox handler always adds a second handler, so an event with a single user handler never took the single-handler fast path. A user handler that used the same `DbContext` therefore ran concurrently with the outbox write. EF Core does not support that ([Avoiding DbContext threading issues](https://learn.microsoft.com/ef/core/dbcontext-configuration/#avoiding-dbcontext-threading-issues)). It fails intermittently with "A second operation was started on this context instance", or it corrupts the unit of work without any error (#812). + +## Decision + +- `PulseMediator` splits the resolved handlers into the outbox handlers (`is OutboxEventHandler`) and the rest. +- Inside the innermost step of the event interceptor chain, it invokes the outbox handlers one after another. It then passes only the remaining handlers to the resolved dispatcher: a keyed per-event dispatcher, the global one, or the parallel default. Interceptors such as event filters therefore still wrap and can suppress the outbox write. +- The error contract does not change. An outbox failure does not stop the other handlers. All failures are thrown together as one flat `AggregateException` after every handler has run. +- Events without an outbox handler take the unchanged path. +- `ParallelEventDispatcher` remains the default. Its documentation and the READMEs now state that user handlers sharing a scoped `DbContext` or connection must use `UseEventDispatcherFor()` (or `UseDefaultEventDispatcher()` for all events). + +## Consequences + +- With the default configuration, the outbox write never runs concurrently with a user handler, for every dispatcher. +- An outbox plus exactly one user handler now takes the parallel dispatcher's single-handler fast path. +- The outbox write always happens first, including with `SequentialEventDispatcher`. Before, the order followed registration order. Nothing documented relied on the old order. +- Several user handlers that share one `DbContext` can still collide under the parallel default. That is the documented contract of `ParallelEventDispatcher`, and the documented workaround is `SequentialEventDispatcher`. +- The dispatcher decides nothing for the outbox handler, so a custom dispatcher (for example rate limiting or priorities) no longer sees it. +- A secondary effect is not changed by this decision. `EntityFrameworkOutboxRepository.AddAsync` calls `SaveChangesAsync` on the shared context, so it also commits changes the caller has made but not yet saved. Running the outbox first at least keeps it from flushing half-built changes of other handlers. The Entity Framework README documents this. + +## Alternatives Considered + +- **Make `SequentialEventDispatcher` the default.** This is safe for all shared scoped services, but it is a breaking behavioral change for every consumer, and it removes the throughput of independent handlers. It was rejected in favour of the smaller, targeted fix. Consumers can still opt in. +- **Fix the single-handler case inside `ParallelEventDispatcher`.** This only covers one dispatcher. Keyed, rate-limited, prioritized and custom dispatchers would still run the outbox concurrently. +- **Run the outbox handler after the other handlers.** Its `SaveChangesAsync` would then commit whatever the user handlers left pending. Running it first keeps the flush limited to the caller's state. +- **Add a marker interface for scope-sharing framework handlers.** The outbox handler is currently the only one, so that would be a speculative abstraction. + +## Related Decisions (Optional) + +- [Extensibility Interface Evolution Pre-1.0](./2026-09-24-extensibility-interface-evolution-pre-1-0.md) - No public interface changes; this is a behavioral fix inside the mediator. diff --git a/src/NetEvolve.Pulse.EntityFramework/EntityFrameworkExtensions.cs b/src/NetEvolve.Pulse.EntityFramework/EntityFrameworkExtensions.cs index b115b03c..29c91525 100644 --- a/src/NetEvolve.Pulse.EntityFramework/EntityFrameworkExtensions.cs +++ b/src/NetEvolve.Pulse.EntityFramework/EntityFrameworkExtensions.cs @@ -36,6 +36,14 @@ public static class EntityFrameworkExtensions /// Note: /// The DbContext must already be registered in the service collection. /// This method does not register the DbContext itself. + /// Shared DbContext: + /// The outbox writes through the scoped of the publishing caller, the same + /// instance that scoped event handlers receive. PublishAsync always runs the outbox handler first + /// and on its own, so it never uses the context concurrently with another handler. If two or more of your + /// own handlers for the same event use the context, register + /// UseDefaultEventDispatcher<SequentialEventDispatcher>(), because the default parallel + /// dispatcher runs them concurrently. Storing a message calls SaveChangesAsync on the shared + /// context, which also saves any changes the caller has made but not yet saved. /// /// /// diff --git a/src/NetEvolve.Pulse.EntityFramework/README.md b/src/NetEvolve.Pulse.EntityFramework/README.md index 276d9789..4ec6ffe1 100644 --- a/src/NetEvolve.Pulse.EntityFramework/README.md +++ b/src/NetEvolve.Pulse.EntityFramework/README.md @@ -154,6 +154,13 @@ public class OrderService } ``` +### Event Handlers Sharing the DbContext + +When you publish through `IMediator.PublishAsync`, the outbox writes through the caller's scoped `DbContext`. That is the same instance your scoped event handlers get. The mediator always runs the outbox handler first and on its own, so it never uses the context concurrently with one of your handlers. + +- If two or more of your own handlers for the same event use the context, register `UseEventDispatcherFor()` for that event (or `UseDefaultEventDispatcher()` for all events). The default `ParallelEventDispatcher` runs them concurrently, which EF Core does not support. +- Storing an outbox message calls `SaveChangesAsync` on the shared context. This also saves any changes the caller has made but not yet saved. Wrap the business changes and `PublishAsync` in one transaction (see above) if they must commit or roll back together. + ## Multi-Provider Support `OutboxMessageConfigurationFactory.Create(this)` automatically picks the right column types diff --git a/src/NetEvolve.Pulse.Extensibility/IEvent.cs b/src/NetEvolve.Pulse.Extensibility/IEvent.cs index 2d83d9e9..494e6625 100644 --- a/src/NetEvolve.Pulse.Extensibility/IEvent.cs +++ b/src/NetEvolve.Pulse.Extensibility/IEvent.cs @@ -5,7 +5,8 @@ /// Events are immutable notifications where multiple subscribers may react independently. /// /// -/// ⚠️ Event handlers execute in parallel and should be idempotent. They should not depend on execution order. +/// ⚠️ With the default dispatcher, event handlers execute in parallel and should be idempotent. They should not +/// depend on execution order. The outbox handler registered by AddOutbox() always runs first and on its own. /// Use past-tense names (OrderCreated, PaymentProcessed, UserRegistered). /// /// diff --git a/src/NetEvolve.Pulse.Extensibility/IEventDispatcher.cs b/src/NetEvolve.Pulse.Extensibility/IEventDispatcher.cs index 0850294c..5d17db4f 100644 --- a/src/NetEvolve.Pulse.Extensibility/IEventDispatcher.cs +++ b/src/NetEvolve.Pulse.Extensibility/IEventDispatcher.cs @@ -18,6 +18,8 @@ /// PrioritizedEventDispatcher: Orders handlers by before execution /// TransactionalEventDispatcher: Stores events in for reliable delivery /// +/// The outbox handler registered by AddOutbox() is never passed to a dispatcher: the mediator always runs it +/// first and on its own, and passes only the remaining handlers to the dispatcher. /// Custom Implementations: /// Implement this interface for advanced scenarios such as: /// diff --git a/src/NetEvolve.Pulse/Dispatchers/ParallelEventDispatcher.cs b/src/NetEvolve.Pulse/Dispatchers/ParallelEventDispatcher.cs index f03a0fc2..80731ffc 100644 --- a/src/NetEvolve.Pulse/Dispatchers/ParallelEventDispatcher.cs +++ b/src/NetEvolve.Pulse/Dispatchers/ParallelEventDispatcher.cs @@ -25,6 +25,14 @@ /// ⚠️ Caution: /// Not suitable when handler execution order matters or when handlers access shared mutable state. /// Consider for such scenarios. +/// Event handlers are resolved from the caller's scope, so scoped handlers share the caller's scoped +/// services. Do not combine this dispatcher with two or more handlers of the same event that use the +/// same scoped DbContext or database connection: they run concurrently on it, which EF Core and +/// most ADO.NET providers do not support. Register +/// UseEventDispatcherFor<TEvent, SequentialEventDispatcher>() for such events instead, or +/// UseDefaultEventDispatcher<SequentialEventDispatcher>() to make every event sequential. +/// The outbox handler registered by AddOutbox() is not affected: the mediator always runs it +/// first and on its own, and passes only the remaining handlers to this dispatcher. /// /// /// diff --git a/src/NetEvolve.Pulse/Dispatchers/PrioritizedEventDispatcher.cs b/src/NetEvolve.Pulse/Dispatchers/PrioritizedEventDispatcher.cs index aa36c7bf..6b8977b3 100644 --- a/src/NetEvolve.Pulse/Dispatchers/PrioritizedEventDispatcher.cs +++ b/src/NetEvolve.Pulse/Dispatchers/PrioritizedEventDispatcher.cs @@ -34,6 +34,9 @@ /// Sequential group execution impacts overall throughput compared to fully parallel execution. /// Use only when handler ordering across groups is critical. /// Consider for independent handlers. +/// Outbox: +/// The outbox handler registered by AddOutbox() is never passed to this dispatcher: the mediator +/// always runs it first and on its own, and passes only the remaining handlers here. /// /// /// diff --git a/src/NetEvolve.Pulse/Dispatchers/RateLimitedEventDispatcher.cs b/src/NetEvolve.Pulse/Dispatchers/RateLimitedEventDispatcher.cs index fde26fc2..693540d9 100644 --- a/src/NetEvolve.Pulse/Dispatchers/RateLimitedEventDispatcher.cs +++ b/src/NetEvolve.Pulse/Dispatchers/RateLimitedEventDispatcher.cs @@ -33,6 +33,9 @@ /// Consider using scoped lifetime when concurrency varies per request /// Monitor queue depth in high-throughput scenarios /// +/// Outbox: +/// The outbox handler registered by AddOutbox() is never passed to this dispatcher: the mediator +/// always runs it first and on its own, and passes only the remaining handlers here. /// /// /// diff --git a/src/NetEvolve.Pulse/Dispatchers/SequentialEventDispatcher.cs b/src/NetEvolve.Pulse/Dispatchers/SequentialEventDispatcher.cs index 04d83687..6cada876 100644 --- a/src/NetEvolve.Pulse/Dispatchers/SequentialEventDispatcher.cs +++ b/src/NetEvolve.Pulse/Dispatchers/SequentialEventDispatcher.cs @@ -10,6 +10,8 @@ /// Execution Behavior: /// Handlers execute one at a time in the order they were registered in the DI container. /// Each handler completes before the next one starts. +/// The outbox handler registered by AddOutbox() is never passed to this dispatcher: the mediator +/// always runs it first and on its own, so registration order applies only to the remaining handlers. /// Error Handling: /// Individual handler failures do not prevent subsequent handlers from executing. /// All handlers are executed regardless of failures. If any handlers fail, an diff --git a/src/NetEvolve.Pulse/Internals/PulseMediator.cs b/src/NetEvolve.Pulse/Internals/PulseMediator.cs index c23ff993..cf61ff71 100644 --- a/src/NetEvolve.Pulse/Internals/PulseMediator.cs +++ b/src/NetEvolve.Pulse/Internals/PulseMediator.cs @@ -8,6 +8,7 @@ using Microsoft.Extensions.Logging; using NetEvolve.Pulse.Dispatchers; using NetEvolve.Pulse.Extensibility; +using NetEvolve.Pulse.Outbox; /// /// Internal implementation of that coordinates dispatching requests and events to their handlers. @@ -74,6 +75,11 @@ public PulseMediator( /// handlers share the caller's scoped services (e.g. the same DbContext) and can participate in the /// caller's transaction — a prerequisite for atomic outbox writes. When the mediator is resolved from /// the root provider, scoped handler dependencies live in the root scope for the provider's lifetime. + /// Because the outbox handler registered by AddOutbox() writes through those shared services, it is + /// invoked first and on its own. Only the remaining handlers go to the dispatcher, so the outbox write never + /// runs concurrently with another handler. User handlers that share a scoped DbContext or connection + /// with each other need , because the default + /// runs them concurrently. /// public Task PublishAsync([NotNull] TEvent message, CancellationToken cancellationToken = default) where TEvent : IEvent @@ -179,7 +185,7 @@ public Task SendAsync( /// A task representing the asynchronous execution of the event through the interceptor pipeline. private Task ExecuteAsync( TEvent msg, - IEnumerable> handlers, + IEventHandler[] handlers, IServiceProvider serviceProvider, CancellationToken cancellationToken ) @@ -190,9 +196,20 @@ CancellationToken cancellationToken // Resolve dispatcher: keyed by event type first, then global, then default var dispatcher = serviceProvider.GetKeyedService(typeof(TEvent)) ?? _eventDispatcher; + // The outbox handler writes through the caller's scoped services (e.g. the same DbContext), so it + // must never run concurrently with other handlers. It runs first and sequentially; only the + // remaining handlers go through the dispatcher. + var outboxHandlers = Array.FindAll(handlers, static h => h is OutboxEventHandler); + var otherHandlers = + outboxHandlers.Length == 0 + ? handlers + : Array.FindAll(handlers, static h => h is not OutboxEventHandler); + // Create the dispatch action that uses the resolved dispatcher Task DispatchAsync(TEvent message, CancellationToken token) => - dispatcher.DispatchAsync(message, handlers, InvokeHandlerAsync, token); + outboxHandlers.Length == 0 + ? dispatcher.DispatchAsync(message, handlers, InvokeHandlerAsync, token) + : DispatchOutboxFirstAsync(message, outboxHandlers, otherHandlers, dispatcher, token); // Build the interceptor chain from innermost (dispatcher) to outermost (first interceptor) var next = DispatchAsync; @@ -338,6 +355,70 @@ is NativeAotInterceptorExtensions.Marker marker return [.. _serviceProvider.GetServices()]; } + /// + /// Invokes the outbox handlers sequentially and then dispatches the remaining handlers, so the outbox + /// write never overlaps another handler that shares the caller's scoped services. + /// + /// The type of event being processed. + /// The event to process. + /// The outbox handlers, invoked first and one after another. + /// The remaining handlers, passed to . + /// The dispatcher for the remaining handlers. + /// A token to monitor for cancellation requests. + /// A task representing the asynchronous dispatch. + /// + /// Thrown after all handlers have run when an outbox handler failed; it also carries the failures of + /// the remaining handlers, flattened from the dispatcher's . + /// + private async Task DispatchOutboxFirstAsync( + TEvent message, + IEventHandler[] outboxHandlers, + IEventHandler[] otherHandlers, + IEventDispatcher dispatcher, + CancellationToken cancellationToken + ) + where TEvent : IEvent + { + cancellationToken.ThrowIfCancellationRequested(); + + var exceptions = new List(); + + foreach (var handler in outboxHandlers) + { + try + { + await InvokeHandlerAsync(handler, message, cancellationToken).ConfigureAwait(false); + } + catch (Exception ex) when (!cancellationToken.IsCancellationRequested) + { + exceptions.Add(ex); + } + } + + if (otherHandlers.Length > 0) + { + try + { + await dispatcher + .DispatchAsync(message, otherHandlers, InvokeHandlerAsync, cancellationToken) + .ConfigureAwait(false); + } + catch (AggregateException ex) when (exceptions.Count > 0) + { + exceptions.AddRange(ex.InnerExceptions); + } + catch (Exception ex) when (exceptions.Count > 0 && !cancellationToken.IsCancellationRequested) + { + exceptions.Add(ex); + } + } + + if (exceptions.Count > 0) + { + throw new AggregateException("One or more event handlers failed.", exceptions); + } + } + /// /// Invokes a single event handler and logs any exception that occurs during execution. /// This wrapper ensures errors are captured before being re-thrown to the dispatcher. diff --git a/src/NetEvolve.Pulse/OutboxExtensions.cs b/src/NetEvolve.Pulse/OutboxExtensions.cs index 40d84013..f1b1b847 100644 --- a/src/NetEvolve.Pulse/OutboxExtensions.cs +++ b/src/NetEvolve.Pulse/OutboxExtensions.cs @@ -37,6 +37,13 @@ public static class OutboxExtensions /// Usage: /// An implementation must be registered separately by calling /// AddSqlServerOutbox, AddEntityFrameworkOutbox, or a custom implementation. + /// Dispatch: + /// The outbox handler writes through the caller's scoped services (for example the same DbContext). + /// PublishAsync therefore always invokes it first + /// and on its own. Only the remaining handlers go to the configured , so the + /// outbox write never runs concurrently with another handler, even under the default + /// . User handlers that share a scoped DbContext or + /// connection with each other still need UseDefaultEventDispatcher<SequentialEventDispatcher>(). /// public static IMediatorBuilder AddOutbox( this IMediatorBuilder configurator, diff --git a/src/NetEvolve.Pulse/README.md b/src/NetEvolve.Pulse/README.md index 0b773f39..c1b0dd3d 100644 --- a/src/NetEvolve.Pulse/README.md +++ b/src/NetEvolve.Pulse/README.md @@ -281,6 +281,21 @@ processorOptions.EventTypeOverrides[typeof(BulkEvent)] = new OutboxEventTypeOpti See [NetEvolve.Pulse.EntityFramework](https://www.nuget.org/packages/NetEvolve.Pulse.EntityFramework/) or [NetEvolve.Pulse.SqlServer](https://www.nuget.org/packages/NetEvolve.Pulse.SqlServer/) for persistence provider setup. +#### Outbox and Handlers Sharing the Caller's Scope + +`PublishAsync` resolves event handlers from the caller's scope, so scoped handlers get the same `DbContext` or connection as the publishing code and the outbox. Since EF Core does not support concurrent operations on one `DbContext`, the mediator always runs the outbox handler first and on its own. Only the remaining handlers go to the configured dispatcher, so the outbox write never overlaps another handler, even under the default `ParallelEventDispatcher`. + +If two or more of your own handlers for the same event use the same scoped `DbContext` or connection, the parallel default still runs them concurrently. Switch to sequential dispatch for these events: + +```csharp +services.AddPulse(config => config + .AddOutbox() + .UseEventDispatcherFor() +); +``` + +To make every event sequential instead, use `UseDefaultEventDispatcher()`. + ### Payload Serialization Pulse uses `IPayloadSerializer` (from `NetEvolve.Pulse.Extensibility`) for all internal serialization needs, including outbox message payloads, distributed cache entries, and audit trail data. A default implementation based on System.Text.Json is registered automatically when you call `AddPulse()`. diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxSharedContextTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxSharedContextTestsBase.cs new file mode 100644 index 00000000..c1835064 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxSharedContextTestsBase.cs @@ -0,0 +1,123 @@ +namespace NetEvolve.Pulse.Tests.Integration.Outbox; + +using System.Diagnostics.CodeAnalysis; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.Extensions.DependencyInjection; +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Extensibility; +using NetEvolve.Pulse.Extensibility.Outbox; +using NetEvolve.Pulse.Outbox; +using NetEvolve.Pulse.Tests.Integration.Internals; +using NetEvolve.Pulse.Tests.Integration.Internals.Outbox; + +/// +/// Verifies that publishing an event under the default dispatcher never uses the caller's scoped +/// concurrently from the Entity Framework outbox and a user handler (#812). +/// +/// +/// The user handler occupies the shared context for a while through EF Core's own +/// , exactly like a slow query would, so any outbox write that runs +/// at the same time fails with EF Core's "A second operation was started on this context instance" +/// . This keeps the test deterministic on fast providers such as +/// SQLite, where a real query would finish before the overlap could be observed. +/// +[TestGroup("Outbox")] +[Timeout(300_000)] +public abstract class EntityFrameworkOutboxSharedContextTestsBase( + IServiceFixture databaseServiceFixture, + IServiceInitializer databaseInitializer +) : PulseTestsBase(databaseServiceFixture, databaseInitializer) +{ + private const int PublishCount = 3; + + private const string HandlerPayload = "handler"; + + /// How long the user handler occupies the shared context. + private static readonly TimeSpan HoldDuration = TimeSpan.FromMilliseconds(250); + + // Test names double as table names and must stay within MySQL's 64-character identifier limit. + [Test] + public async Task PublishAsync_WithSharedDbContextHandler_PersistsBoth(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var context = services.GetRequiredService(); + + var mediator = services.GetRequiredService(); + + for (var i = 0; i < PublishCount; i++) + { + await mediator + .PublishAsync(new SharedContextEvent { Id = $"Test{i:D3}" }, token) + .ConfigureAwait(false); + } + + var outbox = services.GetRequiredService(); + var pending = await outbox.GetPendingCountAsync(token).ConfigureAwait(false); + var handled = await context + .OutboxMessages.AsNoTracking() + .CountAsync(m => m.Payload == HandlerPayload, token) + .ConfigureAwait(false); + + _ = await Assert.That(pending).IsEqualTo(PublishCount); + _ = await Assert.That(handled).IsEqualTo(PublishCount); + }, + cancellationToken, + configureServices: services => + services + .AddScoped, SharedContextHandler>() + .Configure(options => options.DisableProcessing = true) + ) + .ConfigureAwait(false); + + [SuppressMessage( + "Major Code Smell", + "S1144:Unused private types or members should be removed", + Justification = "Resolved through dependency injection." + )] + private sealed class SharedContextHandler( + EntityFrameworkOutboxInitializer.TestDbContext context, + TimeProvider timeProvider + ) : IEventHandler + { + public async Task HandleAsync(SharedContextEvent message, CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + // Occupy the shared context the way an in-flight query does. + using (context.GetService().EnterCriticalSection()) + { + await Task.Delay(HoldDuration, cancellationToken).ConfigureAwait(false); + } + + // The handler's own change: a row it writes through the same context. + var now = timeProvider.GetUtcNow(); + _ = await context + .OutboxMessages.AddAsync( + new OutboxMessage + { + Id = Guid.NewGuid(), + EventType = typeof(SharedContextEvent), + Payload = HandlerPayload, + CreatedAt = now, + UpdatedAt = now, + Status = OutboxMessageStatus.Completed, + }, + cancellationToken + ) + .ConfigureAwait(false); + _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); + } + } + + private sealed class SharedContextEvent : IEvent + { + public string? CausationId { get; set; } + public string? CorrelationId { get; set; } + + public required string Id { get; init; } + + public DateTimeOffset? PublishedAt { get; set; } + } +} diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/InMemoryEntityFrameworkOutboxSharedContextTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/InMemoryEntityFrameworkOutboxSharedContextTests.cs new file mode 100644 index 00000000..909d13e1 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/InMemoryEntityFrameworkOutboxSharedContextTests.cs @@ -0,0 +1,17 @@ +namespace NetEvolve.Pulse.Tests.Integration.Outbox; + +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Tests.Integration.Internals; +using NetEvolve.Pulse.Tests.Integration.Internals.Outbox; +using NetEvolve.Pulse.Tests.Integration.Internals.Services; + +[ClassDataSource( + Shared = [SharedType.None, SharedType.PerTestSession] +)] +[TestGroup("InMemory")] +[TestGroup("EntityFramework")] +[InheritsTests] +public class InMemoryEntityFrameworkOutboxSharedContextTests( + IServiceFixture databaseServiceFixture, + IServiceInitializer databaseInitializer +) : EntityFrameworkOutboxSharedContextTestsBase(databaseServiceFixture, databaseInitializer); diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/MySqlEntityFrameworkOutboxSharedContextTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/MySqlEntityFrameworkOutboxSharedContextTests.cs new file mode 100644 index 00000000..6cf3e759 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/MySqlEntityFrameworkOutboxSharedContextTests.cs @@ -0,0 +1,17 @@ +namespace NetEvolve.Pulse.Tests.Integration.Outbox; + +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Tests.Integration.Internals; +using NetEvolve.Pulse.Tests.Integration.Internals.Outbox; +using NetEvolve.Pulse.Tests.Integration.Internals.Services; + +[ClassDataSource( + Shared = [SharedType.None, SharedType.None] +)] +[TestGroup("MySql")] +[TestGroup("EntityFramework")] +[InheritsTests] +public class MySqlEntityFrameworkOutboxSharedContextTests( + IServiceFixture databaseServiceFixture, + IServiceInitializer databaseInitializer +) : EntityFrameworkOutboxSharedContextTestsBase(databaseServiceFixture, databaseInitializer); diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/PostgreSqlEntityFrameworkOutboxSharedContextTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/PostgreSqlEntityFrameworkOutboxSharedContextTests.cs new file mode 100644 index 00000000..d6523534 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/PostgreSqlEntityFrameworkOutboxSharedContextTests.cs @@ -0,0 +1,17 @@ +namespace NetEvolve.Pulse.Tests.Integration.Outbox; + +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Tests.Integration.Internals; +using NetEvolve.Pulse.Tests.Integration.Internals.Outbox; +using NetEvolve.Pulse.Tests.Integration.Internals.Services; + +[ClassDataSource( + Shared = [SharedType.None, SharedType.None] +)] +[TestGroup("PostgreSql")] +[TestGroup("EntityFramework")] +[InheritsTests] +public class PostgreSqlEntityFrameworkOutboxSharedContextTests( + IServiceFixture databaseServiceFixture, + IServiceInitializer databaseInitializer +) : EntityFrameworkOutboxSharedContextTestsBase(databaseServiceFixture, databaseInitializer); diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/SQLiteEntityFrameworkOutboxSharedContextTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/SQLiteEntityFrameworkOutboxSharedContextTests.cs new file mode 100644 index 00000000..b2ac2efe --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/SQLiteEntityFrameworkOutboxSharedContextTests.cs @@ -0,0 +1,17 @@ +namespace NetEvolve.Pulse.Tests.Integration.Outbox; + +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Tests.Integration.Internals; +using NetEvolve.Pulse.Tests.Integration.Internals.Outbox; +using NetEvolve.Pulse.Tests.Integration.Internals.Services; + +[ClassDataSource( + Shared = [SharedType.None, SharedType.None] +)] +[TestGroup("SQLite")] +[TestGroup("EntityFramework")] +[InheritsTests] +public class SQLiteEntityFrameworkOutboxSharedContextTests( + IServiceFixture databaseServiceFixture, + IServiceInitializer databaseInitializer +) : EntityFrameworkOutboxSharedContextTestsBase(databaseServiceFixture, databaseInitializer); diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/SqlServerEntityFrameworkOutboxSharedContextTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/SqlServerEntityFrameworkOutboxSharedContextTests.cs new file mode 100644 index 00000000..ec4256d5 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/SqlServerEntityFrameworkOutboxSharedContextTests.cs @@ -0,0 +1,17 @@ +namespace NetEvolve.Pulse.Tests.Integration.Outbox; + +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Tests.Integration.Internals; +using NetEvolve.Pulse.Tests.Integration.Internals.Outbox; +using NetEvolve.Pulse.Tests.Integration.Internals.Services; + +[ClassDataSource( + Shared = [SharedType.None, SharedType.None] +)] +[TestGroup("SqlServer")] +[TestGroup("EntityFramework")] +[InheritsTests] +public class SqlServerEntityFrameworkOutboxSharedContextTests( + IServiceFixture databaseServiceFixture, + IServiceInitializer databaseInitializer +) : EntityFrameworkOutboxSharedContextTestsBase(databaseServiceFixture, databaseInitializer); diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Internals/PulseMediatorOutboxDispatchTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Internals/PulseMediatorOutboxDispatchTests.cs new file mode 100644 index 00000000..5df9f4d2 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Unit/Internals/PulseMediatorOutboxDispatchTests.cs @@ -0,0 +1,313 @@ +namespace NetEvolve.Pulse.Tests.Unit.Internals; + +using System.Collections.Concurrent; +using System.Diagnostics.CodeAnalysis; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Dispatchers; +using NetEvolve.Pulse.Extensibility; +using NetEvolve.Pulse.Extensibility.Outbox; +using NetEvolve.Pulse.Internals; +using NetEvolve.Pulse.Outbox; +using TUnit.Core; + +/// +/// Verifies how dispatches the framework's +/// next to user handlers that share the caller's scoped services (e.g. the same DbContext). +/// +[TestGroup("Internals")] +public class PulseMediatorOutboxDispatchTests +{ + [Test] + public async Task PublishAsync_WithOutboxAndScopedHandler_DefaultDispatcher_NeverUsesSharedDependencyConcurrently( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var services = CreateServices(); + _ = services.AddScoped, SharedDependencyHandler>(); + var serviceProvider = services.BuildServiceProvider(); + + var scope = serviceProvider.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var mediator = CreateMediator(scope.ServiceProvider); + var shared = scope.ServiceProvider.GetRequiredService(); + + await mediator.PublishAsync(new TestEvent(), cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(string.Join(",", shared.Calls)).IsEqualTo("outbox,handler"); + _ = await Assert.That(shared.MaxConcurrency).IsEqualTo(1); + } + } + } + + [Test] + public async Task PublishAsync_WithOutboxAndSingleHandler_PassesOnlyUserHandlerToDispatcher( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var dispatcher = new RecordingDispatcher(); + var services = CreateServices(); + _ = services.AddScoped, SharedDependencyHandler>(); + var serviceProvider = services.BuildServiceProvider(); + + var scope = serviceProvider.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var mediator = CreateMediator(scope.ServiceProvider, dispatcher); + var shared = scope.ServiceProvider.GetRequiredService(); + + await mediator.PublishAsync(new TestEvent(), cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(dispatcher.DispatchedHandlers).HasSingleItem(); + _ = await Assert.That(dispatcher.DispatchedHandlers[0]).IsTypeOf(); + _ = await Assert.That(string.Join(",", shared.Calls)).IsEqualTo("outbox,handler"); + } + } + } + + [Test] + public async Task PublishAsync_WithOnlyOutboxHandler_StoresEventWithoutDispatcher( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var dispatcher = new RecordingDispatcher(); + var serviceProvider = CreateServices().BuildServiceProvider(); + + var scope = serviceProvider.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var mediator = CreateMediator(scope.ServiceProvider, dispatcher); + var shared = scope.ServiceProvider.GetRequiredService(); + + await mediator.PublishAsync(new TestEvent(), cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(string.Join(",", shared.Calls)).IsEqualTo("outbox"); + _ = await Assert.That(dispatcher.DispatchedHandlers).IsEmpty(); + } + } + } + + [Test] + public async Task PublishAsync_WithFailingOutbox_StillRunsOtherHandlersAndThrowsFlatAggregate( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var services = CreateServices(); + _ = services.AddScoped, SharedDependencyHandler>(); + _ = services.AddScoped, ThrowingHandler>(); + var serviceProvider = services.BuildServiceProvider(); + + var scope = serviceProvider.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var mediator = CreateMediator(scope.ServiceProvider, new SequentialEventDispatcher()); + var shared = scope.ServiceProvider.GetRequiredService(); + shared.FailOutbox = true; + + var exception = await Assert.ThrowsAsync(async () => + await mediator.PublishAsync(new TestEvent(), cancellationToken).ConfigureAwait(false) + ); + + using (Assert.Multiple()) + { + _ = await Assert.That(exception!.InnerExceptions).Count().IsEqualTo(2); + _ = await Assert.That(exception.InnerExceptions).All(x => x is InvalidOperationException); + _ = await Assert.That(shared.Calls).Contains("handler"); + } + } + } + + [Test] + public async Task PublishAsync_WithSuppressingInterceptor_DoesNotStoreOutboxMessage( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var services = CreateServices(); + _ = services.AddScoped, SharedDependencyHandler>(); + _ = services.AddSingleton, SuppressingInterceptor>(); + var serviceProvider = services.BuildServiceProvider(); + + var scope = serviceProvider.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var mediator = CreateMediator(scope.ServiceProvider); + var shared = scope.ServiceProvider.GetRequiredService(); + + await mediator.PublishAsync(new TestEvent(), cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(shared.Calls).IsEmpty(); + } + } + + private static ServiceCollection CreateServices() + { + var services = new ServiceCollection(); + _ = services.AddLogging(); + _ = services.AddScoped(); + _ = services.AddScoped(); + _ = services.AddScoped, OutboxEventHandler>(); + return services; + } + + private static PulseMediator CreateMediator( + IServiceProvider serviceProvider, + IEventDispatcher? dispatcher = null + ) => + new( + serviceProvider.GetRequiredService>(), + serviceProvider, + TimeProvider.System, + dispatcher + ); + + /// + /// Stands in for a scoped DbContext: records who used it and how many callers used it at once. + /// + private sealed class SharedDependency + { + private int _active; + private int _maxConcurrency; + + public ConcurrentQueue Calls { get; } = new(); + + public int MaxConcurrency => Volatile.Read(ref _maxConcurrency); + + public bool FailOutbox { get; set; } + + public async Task UseAsync(string caller, CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var active = Interlocked.Increment(ref _active); + try + { + int current; + do + { + current = Volatile.Read(ref _maxConcurrency); + } while ( + active > current && Interlocked.CompareExchange(ref _maxConcurrency, active, current) != current + ); + + Calls.Enqueue(caller); + + // Hold the "context" long enough for an overlapping caller to observe it. + await Task.Delay(100, cancellationToken).ConfigureAwait(false); + } + finally + { + _ = Interlocked.Decrement(ref _active); + } + } + } + + [SuppressMessage( + "Major Code Smell", + "S1144:Unused private types or members should be removed", + Justification = "Resolved through dependency injection." + )] + private sealed class SharedDependencyOutbox(SharedDependency shared) : IEventOutbox + { + public async Task StoreAsync(TEvent message, CancellationToken cancellationToken = default) + where TEvent : IEvent + { + cancellationToken.ThrowIfCancellationRequested(); + + await shared.UseAsync("outbox", cancellationToken).ConfigureAwait(false); + + if (shared.FailOutbox) + { + throw new InvalidOperationException("Outbox failure"); + } + } + } + + [SuppressMessage( + "Major Code Smell", + "S1144:Unused private types or members should be removed", + Justification = "Resolved through dependency injection." + )] + private sealed class SharedDependencyHandler(SharedDependency shared) : IEventHandler + { + public Task HandleAsync(TestEvent message, CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + return shared.UseAsync("handler", cancellationToken); + } + } + + private sealed class ThrowingHandler : IEventHandler + { + public Task HandleAsync(TestEvent message, CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + throw new InvalidOperationException("Handler failure"); + } + } + + private sealed class SuppressingInterceptor : IEventInterceptor + { + public Task HandleAsync( + TestEvent message, + Func handler, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + return Task.CompletedTask; + } + } + + private sealed class RecordingDispatcher : IEventDispatcher + { + public List DispatchedHandlers { get; } = []; + + public async Task DispatchAsync( + TEvent message, + IEnumerable> handlers, + Func, TEvent, CancellationToken, Task> invoker, + CancellationToken cancellationToken + ) + where TEvent : IEvent + { + cancellationToken.ThrowIfCancellationRequested(); + + foreach (var handler in handlers) + { + DispatchedHandlers.Add(handler); + await invoker(handler, message, cancellationToken).ConfigureAwait(false); + } + } + } + + private sealed class TestEvent : IEvent + { + public string Id { get; init; } = Guid.NewGuid().ToString(); + + public string? CausationId { get; set; } + public string? CorrelationId { get; set; } + + public DateTimeOffset? PublishedAt { get; set; } + } +}