This guide documents the optional SharpCoreDB.CQRS package and how to use it with SharpCoreDB Event Sourcing.
SharpCoreDB.CQRS adds a lightweight command-side layer on top of SharpCoreDB so applications can model explicit commands, dispatch them through handlers, collect domain events on aggregates, and publish integration messages through an outbox.
At a package level, SharpCoreDB.CQRS is responsible for:
- command contracts and handlers
- command dispatching
- aggregate-side pending domain event collection
- outbox storage and dispatch orchestration
- retry, dead-letter, and hosted worker support
- bridging aggregate events into reliable outbox messages
SharpCoreDB.CQRS does not persist event streams and it does not execute projections. Those responsibilities stay in SharpCoreDB.EventSourcing and SharpCoreDB.Projections.
The 1.9.5 synchronized release keeps the CQRS docs aligned with the shipped package surface: persistent outbox storage, retry and dead-letter metadata, the hosted outbox worker, and the aggregate-to-outbox bridge are all included in the documented baseline.
For a fast setup, see:
docs/cqrs/QUICKSTART.md
SharpCoreDB.CQRS provides:
- Command contracts and handlers
- In-memory and DI-based command dispatchers
AggregateRootbase class for pending domain events- Outbox model and dispatch service
- In-memory outbox store
- Persistent SharpCoreDB-backed outbox store
- Retry-aware outbox failure recording (
attempt_count,last_error,next_attempt_utc) - Hosted outbox worker (
PeriodicTimerpolling)
Non-goals:
- No required dependency on MediatR
- No required transport (Kafka/Rabbit/etc. are integrated via
IOutboxPublisher)
dotnet add package SharpCoreDB.CQRS --version 2.0.0using SharpCoreDB.CQRS;
public readonly record struct PlaceOrderCommand(string OrderId, decimal Amount) : ICommand;using SharpCoreDB.CQRS;
public sealed class PlaceOrderCommandHandler : ICommandHandler<PlaceOrderCommand>
{
public Task<CommandDispatchResult> HandleAsync(PlaceOrderCommand command, CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
if (string.IsNullOrWhiteSpace(command.OrderId))
{
return Task.FromResult(CommandDispatchResult.Failure("OrderId is required."));
}
if (command.Amount <= 0)
{
return Task.FromResult(CommandDispatchResult.Failure("Amount must be > 0."));
}
return Task.FromResult(CommandDispatchResult.Success());
}
}var dispatcher = new InMemoryCommandDispatcher();
dispatcher.RegisterHandler(new PlaceOrderCommandHandler());
var result = await dispatcher.DispatchAsync(new PlaceOrderCommand("order-1", 199.99m), ct);
Console.WriteLine(result.Success);using Microsoft.Extensions.DependencyInjection;
using SharpCoreDB.CQRS;
var services = new ServiceCollection();
services.AddSharpCoreDBCqrs();
services.AddCommandHandler<PlaceOrderCommand, PlaceOrderCommandHandler>();
var provider = services.BuildServiceProvider();
var dispatcher = provider.GetRequiredService<ICommandDispatcher>();
var result = await dispatcher.DispatchAsync(new PlaceOrderCommand("order-2", 249.50m), ct);Use AggregateRoot to capture domain events and dequeue them for persistence/outbox publication.
using SharpCoreDB.CQRS;
public sealed class OrderAggregate : AggregateRoot
{
public void Place(string orderId, decimal amount)
{
RaiseEvent(new OrderPlaced(orderId, amount));
}
}
public readonly record struct OrderPlaced(string OrderId, decimal Amount);
var aggregate = new OrderAggregate();
aggregate.Place("order-3", 50m);
var pendingEvents = aggregate.DequeuePendingEvents();OutboxMessage: payload + routing metadataIOutboxStore: add/read/publish/failure-record lifecycleIOutboxPublisher: pluggable external publisher contractOutboxDispatchService: orchestrates read -> publish -> publish/failure outcome
using SharpCoreDB.CQRS;
public sealed class ConsoleOutboxPublisher : IOutboxPublisher
{
public Task PublishAsync(OutboxMessage message, CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
Console.WriteLine($"Published {message.MessageId} ({message.MessageType})");
return Task.CompletedTask;
}
}var store = new InMemoryOutboxStore();
var publisher = new ConsoleOutboxPublisher();
var dispatch = new OutboxDispatchService(store, publisher);
await store.AddAsync(new OutboxMessage(
MessageId: "msg-1",
AggregateId: "order-1",
MessageType: "OrderPlaced",
Payload: "{}"u8.ToArray(),
CreatedAtUtc: DateTimeOffset.UtcNow,
IsPublished: false),
ct);
var published = await dispatch.DispatchUnpublishedAsync(100, ct);using Microsoft.Extensions.DependencyInjection;
using SharpCoreDB;
using SharpCoreDB.CQRS;
var services = new ServiceCollection();
services.AddSharpCoreDB();
services.AddSharpCoreDBCqrs();
services.AddPersistentOutbox(); // default table: scdb_outbox
services.AddSingleton<IOutboxPublisher, ConsoleOutboxPublisher>();
var provider = services.BuildServiceProvider();var dispatch = provider.GetRequiredService<OutboxDispatchService>();
var published = await dispatch.DispatchUnpublishedAsync(50, ct);For modern schema tables, the store keeps retry metadata columns:
attempt_count(INTEGER)last_error(TEXT)next_attempt_utc(TEXT, round-trip timestamp)
For legacy tables without retry metadata, store operations remain compatible.
When a publisher throws for one message:
- Dispatch does not abort the entire batch.
- It calls
IOutboxStore.RecordFailureAsync(messageId, error, ct). - Store increments attempt count and schedules the next attempt using configured backoff.
- When attempts reach
MaxAttempts, the message is moved to dead-letter storage.
services.AddSharpCoreDBCqrs();
services.AddOutboxRetryPolicy(options =>
{
options.BaseDelay = TimeSpan.FromSeconds(2);
options.MaxDelay = TimeSpan.FromMinutes(2);
options.MaxAttempts = 5;
options.DeadLetterTableName = "scdb_outbox_deadletter";
});var outboxStore = provider.GetRequiredService<IOutboxStore>();
var deadLetters = await outboxStore.GetDeadLettersAsync(100, ct);When a dead-lettered message is ready to retry (after a root cause fix, for example), move it back to the outbox with a reset attempt counter:
var outboxStore = provider.GetRequiredService<IOutboxStore>();
// Inspect dead letters and decide which to requeue
var deadLetters = await outboxStore.GetDeadLettersAsync(100, ct);
foreach (var message in deadLetters)
{
Console.WriteLine($"Dead letter: {message.MessageId} ({message.MessageType})");
}
// Move a specific message back to the live outbox
await outboxStore.RequeueDeadLetterAsync("message-id-to-requeue", ct);The requeue operation:
- Removes the message from the dead-letter table
- Inserts it back into the outbox with
attempt_count = 0andnext_attempt_utcset to now - Makes it immediately eligible for the next dispatch cycle
If the message ID does not exist in the dead-letter table, the call is a no-op.
Effect: at-least-once delivery with bounded retries and deterministic dead-lettering.
Run dispatch continuously with AddOutboxWorker.
services.AddOutboxWorker(options =>
{
options.BatchSize = 100;
options.PollInterval = TimeSpan.FromSeconds(5);
// options.MaxIterations = 10; // Optional for bounded runs/tests
});Typical production setup:
services.AddSharpCoreDB();
services.AddSharpCoreDBCqrs();
services.AddPersistentOutbox();
services.AddSingleton<IOutboxPublisher, ConsoleOutboxPublisher>();
services.AddOutboxWorker(opts =>
{
opts.BatchSize = 50;
opts.PollInterval = TimeSpan.FromSeconds(10);
});- Handle command.
- Persist aggregate/events.
- Write integration event to outbox.
- Let dispatcher/worker publish externally.
- On success: mark published.
- On failure: record retry metadata and reschedule.
This keeps command-side transactions isolated from external transport failures.
Recommended tests for CQRS integrations:
- Handler behavior tests (
Success/ validation failures) - Dispatcher tests (registered/unregistered handlers)
- Outbox tests:
- add/read lifecycle
- mark published behavior
- failure recording updates retry metadata
- scheduled messages excluded until due
- legacy schema compatibility
- Background worker tests:
- periodic dispatch
- cancellation propagation
- exception logging/continuation
The repository contains xUnit v3 coverage for these paths in tests/SharpCoreDB.CQRS.Tests.