Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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<TEvent> 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<TEvent>` 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<TEvent>`) 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<TEvent, SequentialEventDispatcher>()` (or `UseDefaultEventDispatcher<SequentialEventDispatcher>()` 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.
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,14 @@ public static class EntityFrameworkExtensions
/// <para><strong>Note:</strong></para>
/// The DbContext must already be registered in the service collection.
/// This method does not register the DbContext itself.
/// <para><strong>Shared DbContext:</strong></para>
/// The outbox writes through the scoped <typeparamref name="TContext"/> of the publishing caller, the same
/// instance that scoped event handlers receive. <c>PublishAsync</c> 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
/// <c>UseDefaultEventDispatcher&lt;SequentialEventDispatcher&gt;()</c>, because the default parallel
/// dispatcher runs them concurrently. Storing a message calls <c>SaveChangesAsync</c> on the shared
/// context, which also saves any changes the caller has made but not yet saved.
/// </remarks>
/// <example>
/// <code>
Expand Down
7 changes: 7 additions & 0 deletions src/NetEvolve.Pulse.EntityFramework/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<MyEvent, SequentialEventDispatcher>()` for that event (or `UseDefaultEventDispatcher<SequentialEventDispatcher>()` 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
Expand Down
3 changes: 2 additions & 1 deletion src/NetEvolve.Pulse.Extensibility/IEvent.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@
/// Events are immutable notifications where multiple subscribers may react independently.
/// </summary>
/// <remarks>
/// ⚠️ 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 <c>AddOutbox()</c> always runs first and on its own.
/// Use past-tense names (OrderCreated, PaymentProcessed, UserRegistered).
/// </remarks>
/// <example>
Expand Down
2 changes: 2 additions & 0 deletions src/NetEvolve.Pulse.Extensibility/IEventDispatcher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
/// <item><description><c>PrioritizedEventDispatcher</c>: Orders handlers by <see cref="IPrioritizedEventHandler{TEvent}.Priority"/> before execution</description></item>
/// <item><description><c>TransactionalEventDispatcher</c>: Stores events in <see cref="IEventOutbox"/> for reliable delivery</description></item>
/// </list>
/// The outbox handler registered by <c>AddOutbox()</c> 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.
/// <para><strong>Custom Implementations:</strong></para>
/// Implement this interface for advanced scenarios such as:
/// <list type="bullet">
Expand Down
8 changes: 8 additions & 0 deletions src/NetEvolve.Pulse/Dispatchers/ParallelEventDispatcher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,14 @@
/// <para><strong>⚠️ Caution:</strong></para>
/// Not suitable when handler execution order matters or when handlers access shared mutable state.
/// Consider <see cref="SequentialEventDispatcher"/> 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 <c>DbContext</c> or database connection: they run concurrently on it, which EF Core and
/// most ADO.NET providers do not support. Register
/// <c>UseEventDispatcherFor&lt;TEvent, SequentialEventDispatcher&gt;()</c> for such events instead, or
/// <c>UseDefaultEventDispatcher&lt;SequentialEventDispatcher&gt;()</c> to make every event sequential.
/// The outbox handler registered by <c>AddOutbox()</c> is not affected: the mediator always runs it
/// first and on its own, and passes only the remaining handlers to this dispatcher.
/// </remarks>
/// <example>
/// <code>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 <see cref="ParallelEventDispatcher"/> for independent handlers.
/// <para><strong>Outbox:</strong></para>
/// The outbox handler registered by <c>AddOutbox()</c> is never passed to this dispatcher: the mediator
/// always runs it first and on its own, and passes only the remaining handlers here.
/// </remarks>
/// <example>
/// <code>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,9 @@
/// <item><description>Consider using scoped lifetime when concurrency varies per request</description></item>
/// <item><description>Monitor queue depth in high-throughput scenarios</description></item>
/// </list>
/// <para><strong>Outbox:</strong></para>
/// The outbox handler registered by <c>AddOutbox()</c> is never passed to this dispatcher: the mediator
/// always runs it first and on its own, and passes only the remaining handlers here.
/// </remarks>
/// <example>
/// <code>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@
/// <para><strong>Execution Behavior:</strong></para>
/// 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 <c>AddOutbox()</c> 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.
/// <para><strong>Error Handling:</strong></para>
/// Individual handler failures do not prevent subsequent handlers from executing.
/// All handlers are executed regardless of failures. If any handlers fail, an
Expand Down
85 changes: 83 additions & 2 deletions src/NetEvolve.Pulse/Internals/PulseMediator.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
using Microsoft.Extensions.Logging;
using NetEvolve.Pulse.Dispatchers;
using NetEvolve.Pulse.Extensibility;
using NetEvolve.Pulse.Outbox;

/// <summary>
/// Internal implementation of <see cref="IMediator"/> that coordinates dispatching requests and events to their handlers.
Expand Down Expand Up @@ -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 <c>AddOutbox()</c> 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 <c>DbContext</c> or connection
/// with each other need <see cref="SequentialEventDispatcher"/>, because the default
/// <see cref="ParallelEventDispatcher"/> runs them concurrently.
/// </remarks>
public Task PublishAsync<TEvent>([NotNull] TEvent message, CancellationToken cancellationToken = default)
where TEvent : IEvent
Expand Down Expand Up @@ -179,7 +185,7 @@ public Task<TResponse> SendAsync<TCommand, TResponse>(
/// <returns>A task representing the asynchronous execution of the event through the interceptor pipeline.</returns>
private Task ExecuteAsync<TEvent>(
TEvent msg,
IEnumerable<IEventHandler<TEvent>> handlers,
IEventHandler<TEvent>[] handlers,
IServiceProvider serviceProvider,
CancellationToken cancellationToken
)
Expand All @@ -190,9 +196,20 @@ CancellationToken cancellationToken
// Resolve dispatcher: keyed by event type first, then global, then default
var dispatcher = serviceProvider.GetKeyedService<IEventDispatcher>(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<TEvent>);
var otherHandlers =
outboxHandlers.Length == 0
? handlers
: Array.FindAll(handlers, static h => h is not OutboxEventHandler<TEvent>);

// 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;
Expand Down Expand Up @@ -338,6 +355,70 @@ is NativeAotInterceptorExtensions.Marker marker
return [.. _serviceProvider.GetServices<TInterceptor>()];
}

/// <summary>
/// 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.
/// </summary>
/// <typeparam name="TEvent">The type of event being processed.</typeparam>
/// <param name="message">The event to process.</param>
/// <param name="outboxHandlers">The outbox handlers, invoked first and one after another.</param>
/// <param name="otherHandlers">The remaining handlers, passed to <paramref name="dispatcher"/>.</param>
/// <param name="dispatcher">The dispatcher for the remaining handlers.</param>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>A task representing the asynchronous dispatch.</returns>
/// <exception cref="AggregateException">
/// 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 <see cref="AggregateException"/>.
/// </exception>
private async Task DispatchOutboxFirstAsync<TEvent>(
TEvent message,
IEventHandler<TEvent>[] outboxHandlers,
IEventHandler<TEvent>[] otherHandlers,
IEventDispatcher dispatcher,
CancellationToken cancellationToken
)
where TEvent : IEvent
{
cancellationToken.ThrowIfCancellationRequested();

var exceptions = new List<Exception>();

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);
}
}

/// <summary>
/// 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.
Expand Down
7 changes: 7 additions & 0 deletions src/NetEvolve.Pulse/OutboxExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,13 @@ public static class OutboxExtensions
/// <para><strong>Usage:</strong></para>
/// An <see cref="IOutboxRepository"/> implementation must be registered separately by calling
/// <c>AddSqlServerOutbox</c>, <c>AddEntityFrameworkOutbox</c>, or a custom implementation.
/// <para><strong>Dispatch:</strong></para>
/// The outbox handler writes through the caller's scoped services (for example the same <c>DbContext</c>).
/// <c>PublishAsync</c> therefore always invokes it first
/// and on its own. Only the remaining handlers go to the configured <see cref="IEventDispatcher"/>, so the
/// outbox write never runs concurrently with another handler, even under the default
/// <see cref="Dispatchers.ParallelEventDispatcher"/>. User handlers that share a scoped <c>DbContext</c> or
/// connection with each other still need <c>UseDefaultEventDispatcher&lt;SequentialEventDispatcher&gt;()</c>.
/// </remarks>
public static IMediatorBuilder AddOutbox(
this IMediatorBuilder configurator,
Expand Down
Loading
Loading