Skip to main content

Delivering domain events

A raised domain event is a promise that something will happen. An order was placed, so a confirmation goes out, stock is reserved and a van is booked. Raising the event only puts it on the aggregate; nothing has happened yet. The question this page answers is when the promised thing does happen: inside the save, in the same transaction as the aggregate, or after it, from a record that survives a crash.

DDDToolkit.EntityFramework offers one mode for each answer. In-process dispatch runs your handlers while SaveChanges runs. The outbox writes the events to a table in the same transaction and delivers them afterwards. Configure exactly one for production. If you configure neither, an aggregate that raised events refuses to save rather than lose them; see When no delivery mode is configured.

Both modes keep the event inside this process. To send it somewhere else, add a sink to the outbox: that is a third destination and a separate page, Integration events.

This page assumes a context wired up as in Entity Framework.

Choosing​

In process, the handlers run inside the save. Whatever they change on the same context is written with the aggregate, in one transaction, and a handler that throws stops the save:

Show the code: dispatching in process

Register the toolkit with Mediator as the dispatcher, and add the toolkit to the context:

builder.Services.AddMediator(options => options.ServiceLifetime = ServiceLifetime.Scoped);
builder.Services.AddDDDToolkitEntityFramework(options => options.DispatchWithMediator());

builder.Services.AddDbContext<OrderingContext>((services, options) => options
.UseNpgsql(connectionString)
.UseDDDToolkit(services));

The event is a Mediator notification, and a handler is an ordinary Mediator handler. It runs in the saving context's scope, so what it changes on that context is saved with the order:

[DomainEventName("ordering.order-cancelled")]
public sealed record OrderCancelled(OrderId OrderId, string Reason) : DomainEvent, INotification;

public sealed class OrderLog(ILogger<OrderLog> logger) : INotificationHandler<OrderCancelled>
{
public ValueTask Handle(OrderCancelled notification, CancellationToken cancellationToken)
{
logger.LogInformation("Order {OrderId} was cancelled: {Reason}", notification.OrderId, notification.Reason);
return default;
}
}

Ordering/Application/Orders/DomainEvents/OrderLog.cs

See In-process dispatch.

Through the outbox, the save only writes the events down, next to the aggregate. Delivering them is a separate step that comes after the commit, and keeps trying until it succeeds:

Show the code: the outbox

Map the outbox table in the context:

protected override void OnModelCreating(ModelBuilder modelBuilder)
=> modelBuilder.AddDomainEventOutbox(Database);

Turn the outbox on, and run the processor that delivers from it:

builder.Services.AddDDDToolkitEntityFramework(options =>
{
options.DispatchWithMediator(); // where the processor delivers to
options.UseOutbox(outbox => outbox.RegisterEventsFromAssemblyContaining<Program>());
});

builder.Services.AddOutboxBackgroundService<OrderingContext>(TimeSpan.FromSeconds(2));

The handlers are the same as in process. They must be idempotent, because a crash between delivering and marking the row delivers it again. See The outbox.

In-processOutbox
When handlers runInside SaveChanges, before the writeAfter the commit, on the processor
Survives a crashNoYes
Handler failureAborts the saveRecorded on the row, retried
DeliveryBest-effort, at most onceAt-least-once
Handlers must be idempotentNoYes
Extra tableNoYes
Extra moving partNoA processor or background service

Use in-process dispatch for side effects inside the same database, where the transaction is the guarantee you want. Use the outbox for anything that leaves the process, and give it a sink to leave through.

Nothing stops you from mapping the outbox table from the start and switching later. The example context does exactly that, so moving from one mode to the other is a change in the registration with no schema change.

In-process dispatch​

With DispatchInProcess and no outbox, handlers run inside SaveChanges, before the database is written. Whichever way you hand the events to your handlers, that gives you three things:

  • Anything a handler changes on the same DbContext rides the same save, and therefore the same transaction. The interceptor calls DetectChanges after each round so those changes are seen.
  • A throwing handler aborts the save. Nothing is written.
  • New events raised by handlers on tracked aggregates are dispatched in a further round.

What you do not get is durability. Delivery is best-effort. The events are dequeued from the aggregate before dispatch, so once a handler has them they exist only in memory: if the save then fails, the events are gone. Nothing survives a process crash. And a handler that talks to the outside world has already sent the mail or made the HTTP call by the time a later failure rolls the transaction back.

Do not call SaveChanges from a handler in this mode. The save is already in progress and it will pick your changes up.

The events reach your handlers through a delegate. DDDToolkit.Mediator writes it for you; you can also write it yourself.

The short way: DispatchWithMediator()​

DDDToolkit.Mediator writes the delegate for you, against Mediator:

dotnet add package Temp.DDDToolkit.Mediator
using DDDToolkit.Mediator;

builder.Services.AddMediator(options => options.ServiceLifetime = ServiceLifetime.Scoped);
builder.Services.AddDDDToolkitEntityFramework(options => options.DispatchWithMediator());

It resolves IPublisher from the scope that owns the saving DbContext and publishes each event in the order it was raised, awaiting one before starting the next. That is the delegate below plus a check that each event is publishable at all, so it serves both delivery modes: add UseOutbox and the processor delivers through the same call.

Three things are worth knowing before you reach for it.

Your events have to implement Mediator.INotification. Mediator cannot publish anything else. One marker interface for the whole solution is the usual way to say it once:

public interface IOrderingEvent : IDomainEvent, INotification;

An event that does not implement it makes the dispatch throw, naming the event type. It is not skipped. By the time the delegate runs the interceptor has already dequeued the event from the aggregate, so skipping would destroy it with no row, no log and nothing to retry.

Register Mediator as scoped when handlers touch the DbContext. Mediator registers every handler as a singleton by default, and a singleton cannot depend on a scoped service. The lifetime is read off your AddMediator call at compile time, so it has to be written there; setting it any other way throws at start-up.

Keep Mediator.SourceGenerator in your composition root. Mediator generates its implementation and AddMediator into whichever assembly the generator runs in, so reference the generator from the project that builds the container and from nowhere else. Handlers in other projects are still discovered, as long as those projects are referenced. DDDToolkit.Mediator itself references only Mediator.Abstractions, so it adds no generator to your domain projects.

Mediator also reports MSG0005 at build time for a notification that no handler handles. That is usually the mistake it looks like, but an event you deliberately leave unhandled needs the warning suppressed.

The delegate​

DispatchWithMediator() is a convenience. The toolkit core has no mediator dependency and is not getting one: in-process delivery is a delegate, and you can write it against any library or none.

Both modes deliver through the same delegate, so handlers are written once:

Func<IServiceProvider, IReadOnlyList<IDomainEvent>, CancellationToken, Task>

The provider is the scope that owns the saving DbContext, and the events arrive in the order they were raised. Written against Mediator's IPublisher, it looks like this:

builder.Services.AddDDDToolkitEntityFramework(options =>
{
options.DispatchInProcess(async (services, events, cancellationToken) =>
{
var publisher = services.GetRequiredService<IPublisher>();
foreach (var domainEvent in events)
{
await publisher.Publish(domainEvent, cancellationToken);
}
});
});

There is one delegate per process. A second DispatchInProcess or DispatchWithMediator() throws rather than replace the first.

The outbox​

With UseOutbox, nothing is dispatched at save time. Instead one row per event is added to the saving context, so the events commit atomically with the aggregate, and a separate processor delivers them afterwards through the same dispatch delegate.

builder.Services.AddDDDToolkitEntityFramework(options =>
{
options.DispatchWithMediator(); // or your own DispatchInProcess(...) delegate
options.UseOutbox(outbox => outbox.RegisterEventsFromAssemblyContaining<Program>());
});

builder.Services.AddOutboxBackgroundService<OrderingContext>(TimeSpan.FromSeconds(2));

The rows need a table in the context's model: modelBuilder.AddDomainEventOutbox(Database) in OnModelCreating, described under The table. A save with an outbox configured and no table mapped throws, naming the context.

An event is never lost and never published for a transaction that rolled back. Rolling back the aggregate rolls back its events, because they are rows in the same transaction.

The cost is at-least-once delivery. A message is marked processed only after its handlers returned, so a crash in between redelivers it. Handlers must be idempotent, keyed on EventId.

Written this way the outbox is durable but still in-process: the processor hands the event back to the same delegate. Add a sink and the row leaves the process instead:

options.UseOutbox(outbox =>
{
outbox.RegisterEventsFromAssemblyContaining<Program>();
outbox.SendTo<ServiceBusSink>();
});

outbox.AlsoDispatchInProcess = true asks for both; see Which delivery wins. Sinks, the published contract that is not your domain event, and the inbox that makes at-least-once delivery safe to consume are all on Integration events.

The outbox in detail​

The table​

Map it in OnModelCreating:

protected override void OnModelCreating(ModelBuilder modelBuilder)
{
modelBuilder.AddDomainEventOutbox(Database);
}

Database is the context's own property. The method reads one thing from it, the provider name, which is what lets the timestamp columns be the right shape for the database you are actually on. See Timestamps below.

The table is ddd.OutboxMessages by default. The method takes tableName and schema if you want something else, and schema: null puts it in the provider's default schema. SQLite has no schemas and ignores the argument, so the table is plain OutboxMessages there. The mapping sets a primary key on Id, an index on ProcessedAt, and the lengths below.

ColumnTypeMeaning
IdGuid, key, never generatedThe event's EventId. This is the idempotency key
EventNamestring, required, 256The stable name, from [DomainEventName] or the convention, ordering.order-placed
Payloadstring, requiredThe event serialized with System.Text.Json
VersionintThe shape the payload was written in, 1 unless the event type says otherwise; see The outbox reading its own old rows
OccurredAtDateTimeOffsetTaken from the event
AggregateTypestring?, 512CLR type name of the aggregate that raised it
AggregateIdstring?, 256The aggregate's key as text, parts joined with |
CreatedAtDateTimeOffsetWhen the row was written
ProcessedAtDateTimeOffset?When the handlers succeeded, null while pending
AttemptsintHow often delivery was attempted
NextAttemptAtDateTimeOffset?When a failed message may be tried again, null when it is due now; see Failures, retries and poison messages
LastErrorstring?, 4000Type and message of the last failure

CreatedAt comes from options.TimeProvider, which defaults to TimeProvider.System and can be replaced in tests. So does the current time the processor compares NextAttemptAt with.

The index is on ProcessedAt alone, because that is what separates the few pending rows from the many delivered ones. NextAttemptAt is not in it. It only divides the pending rows into due and waiting, and "null or not after now" is not one range an index could seek, so adding it would make the index bigger without making the poll faster.

Payloads are written with System.Text.Json through outbox.JsonOptions, which by default are case-insensitive on read and carry the toolkit's SingleValueObjectConverterFactory, so identifiers and single value objects are stored as their raw values rather than as objects. The same options read the payload back, so change them with care once messages exist.

Timestamps​

OccurredAt, CreatedAt, ProcessedAt and NextAttemptAt are DateTimeOffset in the model. What they become in the database depends on the provider, and you can say otherwise:

// The provider's own instant type.
modelBuilder.AddDomainEventOutbox(Database);

// A UTC DateTime column, on every provider.
modelBuilder.AddDomainEventOutbox(Database, timestamps: DomainEventTimestamps.UtcDateTime);
ProviderProviderDefaultUtcDateTime
PostgreSQLtimestamp with time zonetimestamp with time zone
SQL Serverdatetimeoffsetdatetime2
SQLiteUTC DateTime as textUTC DateTime as text

The one thing the outbox needs from these columns is that they sort as instants, because the processor reads pending messages oldest first and compares NextAttemptAt with the current time. SQLite stores a DateTimeOffset as text and refuses to order by one or compare one, so on SQLite the toolkit converts to a UTC DateTime whatever you ask for. Every other provider has an instant type that orders correctly, and gets it.

Whichever column it lands in, the value is normalized to UTC on the way in. Write a DateTimeOffset that carries +02:00 and the row holds the same instant with an offset of zero, and reads back that way. That makes the two shapes interchangeable in meaning, and it is also what keeps Npgsql happy: PostgreSQL refuses a DateTimeOffset whose offset is not zero.

AddDomainEventInbox takes the same argument for its own ProcessedAt. Pass both the same thing.

If you write migrations by hand, CreateDomainEventOutbox and CreateDomainEventInbox take timestamps too and default to the same ProviderDefault. They read MigrationBuilder.ActiveProvider, so the hand-written table and the scaffolded one agree.

A database created by an earlier 3.0 build may have the older column type on SQL Server; see From an earlier 3.0 build before upgrading one.

Registering event types​

The row stores a name, so the processor needs a name-to-type map:

options.UseOutbox(outbox =>
{
outbox.RegisterEventsFromAssemblyContaining<Program>();
outbox.RegisterEventsFromAssembly(typeof(OrderPlaced).Assembly);
outbox.RegisterEvent<OrderCancelled>();
});

The three calls are additive, so use whichever suits; most applications need only the first. The assembly scans register every concrete type implementing IDomainEvent. Registering two types under the same stable name throws an ArgumentException naming both, which is the failure you want at start-up rather than at delivery time.

This is where a stable name earns its keep. The name is written into the row and read back later, so it must not change while rows are waiting. The conventional name, the module and the class name in kebab case, does not change when the class moves; when the class is renamed, [DomainEventName] keeps the old one. See Stable names.

The compiler can write this registration for you. With the Entity Framework package referenced, every assembly gets an Add{Module}IntegrationEvents() that registers each of its domain events under the name it is stored as, worked out at compile time, so nothing is scanned at start-up:

[assembly: Module("Ordering")]

[DomainEventName("ordering.order-received")] // renamed from OrderReceived; rows keep the old name
public sealed record OrderPlaced(OrderId Order) : DomainEvent;

public sealed record OrderCancelled(OrderId Order) : DomainEvent;
IntegrationEventExtensions.g.cs, shortened
public static OutboxOptions AddOrderingIntegrationEvents(this OutboxOptions outbox)
{
ArgumentNullException.ThrowIfNull(outbox);
outbox.RegisterEvent<OrderCancelled>("ordering.order-cancelled", 1);
outbox.RegisterEvent<OrderPlaced>("ordering.order-received", 1);
return outbox;
}

OrderPlaced goes into the row under its pinned name and OrderCancelled under the conventional one. Call outbox.AddOrderingIntegrationEvents() in place of the assembly scan. The same method registers what a module publishes, which Integration events covers.

An event name nobody registered is not fatal. The processor records the failure on the row, with a LastError that names the missing event and the registration call, increments Attempts, leaves ProcessedAt null and carries on with the rest of the batch. Register the type and the next attempt delivers the row.

Running the processor​

For a hosted application, the background service:

builder.Services.AddOutboxBackgroundService<OrderingContext>(
pollingInterval: TimeSpan.FromSeconds(2),
batchSize: 100);

It registers the processor, the polling options and the hosted service. On each tick it drains the outbox batch after batch until a batch delivers nothing, each batch in its own service scope so the processor and the handlers get a fresh context. A failure in a tick is logged and the next tick tries again.

To drive it yourself, from a job scheduler or a test, register the processor alone:

builder.Services.AddOutboxProcessor<OrderingContext>();
var processor = scope.ServiceProvider.GetRequiredService<OutboxProcessor<OrderingContext>>();
var delivered = await processor.ProcessPendingAsync(batchSize: 100, cancellationToken);

ProcessPendingAsync loads the pending messages that are due, oldest first, by CreatedAt, then OccurredAt, then Id, delivers each one on its own, and returns how many were delivered successfully. The processor throws at construction if the outbox is not enabled, or if it has neither a sink nor a dispatch delegate to deliver through.

Each message is saved through the same scoped context, so a handler that resolves that context and changes an aggregate commits its change together with the processed mark. Events raised by that change become new outbox rows.

Failures, retries and poison messages​

A handler that throws does not stop the batch. The exception type and message are written to LastError, truncated at 4000 characters, Attempts is incremented, ProcessedAt stays null, and the processor moves to the next message. The message is retried later, and on success LastError is cleared.

A failed message is not tried again straight away. The processor records in NextAttemptAt when it may be, and loads only the messages that are due, still oldest first. The wait grows with every failure:

After attempt1234 and later
The message waits5 s25 s2 min 5 s10 min

The wait does two things. A message that keeps failing does not hold back the messages written after it: the next poll reaches past it. And a sink that is down for a few seconds costs a message one attempt, not all of them, because a busy drain does not load it again before its time.

RetryDelay sets the schedule. It is given the number of attempts made so far, 1 after the first failure:

options.UseOutbox(outbox => outbox.RetryDelay = attempts => TimeSpan.FromSeconds(30 * attempts));

TimeSpan.Zero retries on the next poll. OutboxOptions.DefaultRetryDelay is the default schedule, for a policy that builds on it.

MaxAttempts defaults to 10. Messages that reached it are no longer loaded:

options.UseOutbox(outbox => outbox.MaxAttempts = 5);

With both defaults, a message that keeps failing is given up on about an hour after its first attempt. It stays in the table with its last error for you to inspect, and with no NextAttemptAt, so resetting Attempts retries it on the next poll. To retry a message that is still waiting, set its NextAttemptAt to null.

Delivered rows stay too, until something deletes them. services.AddDomainEventRetention<TContext>(...) deletes them once they are older than a window you choose, and never touches a row that was not delivered. See Keeping the tables small.

A batch that delivers nothing ends the background service's drain for that tick, whether its messages failed or none was due. A sink that is down is asked about one batch a tick, not about every message in a loop. The next tick carries on: the messages behind the ones that failed, and those that failed once their time has come.

An outbox per context​

UseOutbox(...) configures one outbox shared by every context. With a context per module, give each producing module an outbox of its own, registered by that module:

services.AddDDDToolkitEntityFramework(options => options.UseOutbox<OrderingContext>(outbox =>
{
outbox.RegisterEventsFromAssemblyContaining<Order>();
outbox.SendToModules();
}));

services.AddOutboxBackgroundService<OrderingContext>(TimeSpan.FromSeconds(2));

That works because AddDDDToolkitEntityFramework can be called as often as you like. The first call registers everything, and every call configures the one options object the process has. That is what lets each module of a modular monolith register its own part next to its own context: its outbox with options.UseOutbox<TContext>(...), the contracts it reads with options.MapIntegrationEvents(...). The host keeps what is process-wide, such as DispatchWithMediator(). The dispatch delegate can only be set once; a second call throws instead of quietly handing one module's events to another module's publisher. See Integration events for the module registration end to end.

A context uses its own outbox when it has one, the shared one when it has not, and none at all when neither exists: its events are then dispatched in process at save time, as if the outbox were not there. options.OutboxFor(typeof(TContext)) says which of the three applies. A processor for a context without an outbox refuses to start and names the context.

Honest limits​

  • Ordering is best-effort. Messages are loaded oldest first, but a failed message waits while messages written after it are delivered, so delivery order is not the order of writing.
  • A waiting message waits out its time. Once a sink is back, a message that failed during the outage is not sent until its NextAttemptAt, up to 10 minutes later with the default schedule. Nothing tells the processor the sink recovered. Set NextAttemptAt to null to send the waiting messages at once.
  • A slow failure still takes its time. The wait keeps a failing message from being tried too often, but each attempt still lasts as long as the sink takes to fail, and a batch is delivered one message at a time. A batch of 100 messages that each wait out a 30 second timeout holds the processor for 50 minutes. Keep sink timeouts short, or the batch small.
  • Concurrent processors can double-deliver. The processor takes no lock, so two instances polling the same table may both pick up the same row. Combined with the crash window before ProcessedAt is written, that is the at-least-once guarantee: handlers must be idempotent, keyed on EventId, which is OutboxMessage.Id.
  • Delivery is not immediate. It happens on the next poll, not at commit.
  • Two sinks share one row. When a message goes to more than one sink and one of them refuses, the message as a whole is retried and the sinks that accepted it see it again. See Integration events.

Tests/DDDToolkit.EntityFramework.Tests/OutboxTests.cs exercises the transactional write, the retries, MaxAttempts, unknown event names and the background service, and OutboxRetryTests.cs next to it the wait between attempts.

The dispatch loop​

Both modes run through the same loop in PublishDomainEventsInterceptor, once per SaveChanges:

  1. Dequeue the pending events of every tracked aggregate. If there are none, stop.
  2. In outbox mode, add a row per event and go back to step 1.
  3. Otherwise call the dispatch delegate with the whole batch, then call DetectChanges, then go back to step 1.

The loop exists because a handler may change a tracked aggregate, which may raise a further event. MaxDispatchRounds caps it at 10 by default. If a further batch of events appears after that many rounds, SaveChanges throws an InvalidOperationException naming the events still pending and suggesting either a handler that triggers itself or a higher limit:

options.MaxDispatchRounds = 20;

When no delivery mode is configured​

If an aggregate has pending events and neither DispatchInProcess nor UseOutbox was called, SaveChanges throws rather than dropping the events. The message names the aggregate types involved and both calls that would fix it. A save with no pending events needs no delivery mode, so a context that only writes aggregates which raise nothing works out of the box.

Prefer SaveChangesAsync​

The dispatch delegate is asynchronous. Synchronous SaveChanges therefore blocks on it. That is safe in console applications and in ASP.NET Core, which have no synchronization context, but it can deadlock under a UI or legacy ASP.NET synchronization context. Both overloads are otherwise identical: events are dispatched before the write either way.

Why Mediator and not MediatR​

MediatR did this job for the first two major versions of the toolkit and does it well. From version 13 it is commercially licensed. This repository prefers dependencies its users can take for free, so the examples and DDDToolkit.Mediator target Mediator, which is MIT and source generated rather than reflection based. Nothing here stops you using MediatR: write the delegate and it publishes through IPublisher exactly as it always did.

Where to look next​