Configuration
This page lists every option you can configure in Nagare, organized by package.
Database connection
Before registering event stores or projections, set up a database connection. Each adapter has its own registration method:
// PostgreSQL
builder.Services.AddNagarePostgresStorage(builder.Configuration);
// SQL Server
builder.Services.AddNagareSqlServerStorage(builder.Configuration);
// SQLite
builder.Services.AddNagareSqliteStorage(builder.Configuration);
// MySQL
builder.Services.AddNagareMySqlStorage(builder.Configuration);All three read a connection string from IConfiguration. By default, they look for "DbConnectionString" in the connection strings section:
{
"ConnectionStrings": {
"DbConnectionString": "Host=localhost;Database=myapp;Username=postgres;Password=secret"
}
}You can pass a different name:
builder.Services.AddNagarePostgresStorage(
builder.Configuration,
connectionStringName: "EventStoreDb");Event store
Register an event store for each event type in your system:
// PostgreSQL
builder.Services.AddPostgresEventStore<BookEvent>();
// SQL Server
builder.Services.AddSqlServerEventStore<BookEvent>();
// SQLite
builder.Services.AddSqliteEventStore<BookEvent>();
// MySQL
builder.Services.AddMySqlEventStore<BookEvent>();By default, events are stored in a table called Events. You can specify a different name:
builder.Services.AddPostgresEventStore<BookEvent>(tableName: "BookEvents");Snapshot store
builder.Services.AddPostgresSnapshotStore<BookState>();Default table name is Snapshots. Override it the same way:
builder.Services.AddPostgresSnapshotStore<BookState>(tableName: "BookSnapshots");Checkpoint store
Each database adapter registers its own checkpoint store and lock provider:
builder.Services.AddPostgresCheckpointStore();
// or
builder.Services.AddSqlServerCheckpointStore();
// or
builder.Services.AddSqliteCheckpointStore();
// or
builder.Services.AddMySqlCheckpointStore();The PostgreSQL, SQL Server, and MySQL checkpoint stores also register a distributed lock provider (PostgresLockProvider, MsSqlLockProvider, or MySqlLockProvider). SQLite registers NoopLockProvider since it's single-process.
Aggregates
Full registration shorthand
Each adapter provides a method that registers the event store, snapshot store, and aggregate in one call:
builder.Services.AddPostgresAggregate<
BookAggregate, BookCommand, BookEvent, BookState>();This is equivalent to:
builder.Services.AddPostgresEventStore<BookEvent>();
builder.Services.AddPostgresSnapshotStore<BookState>();
builder.Services.AddAggregate<BookAggregate, BookCommand, BookEvent, BookState>();AggregateOptions
Control how aggregates load their event history:
| Property | Type | Default | Purpose |
|---|---|---|---|
InitializeBatchSize | int | 50 | Number of events to read per batch when loading an aggregate |
There's no explicit options registration. The batch size is a property on the aggregate options record. Configure it when you have aggregates with long event histories and want to tune the page size for loading.
Subscriptions
PipelineHostingOptions
Every subscription, projection, message subscription, and process-manager event route runs on the same pipeline runtime, configured through PipelineHostingOptions:
builder.Services.AddPipeline<BookProjection, BookEvent>(new PipelineHostingOptions
{
PollDelay = TimeSpan.FromMilliseconds(200),
MaxFailureBackoff = TimeSpan.FromSeconds(60),
});| Property | Type | Default | Purpose |
|---|---|---|---|
PollDelay | TimeSpan | 100ms | Wait between sessions: after a session ends, or while another node holds the pipeline's lock |
IdlePollInterval | TimeSpan? | null | Journal pipelines: longest a caught-up source waits for a wake-up before re-reading the journal. null takes the global default (below) |
MaxFailureBackoff | TimeSpan | 30s | Cap for the exponential backoff between consecutively failing sessions |
FaultAfterConsecutiveFailures | int | 5 | Failing sessions before the readiness health check reports the pipeline as faulted |
FaultClearAfter | TimeSpan | 60s | How long a session must run with nothing in flight before a fault clears when no element has been handled (an idle live source that recovered), and how long a node must stay on standby before it drops a fault it cannot otherwise resolve. Timeout.InfiniteTimeSpan disables both. Must not exceed ~49.7 days (System.Threading.Timer's maximum) |
StandbyHolderStaleAfter | TimeSpan | 90s | With the lease registry on: how stale the lock holder's lease may be before a standby stream is listed under holderUnresponsive and the check reports Degraded. Holders renew every 30s and this node re-reads the lease every third of this value, so keep it above 45s plus clock skew. Timeout.InfiniteTimeSpan disables the check. See stream readiness |
Resilience | ResiliencePipeline? | null | Optional Polly pipeline wrapping each processing session |
DeadLetterPolicy | DeadLetterPolicy? | null | Null = conveyor-belt default (block hard and retry); set to opt into retry-then-park |
Idle poll interval
A caught-up journal pipeline waits on ISubscriptionWakeUp before it reads its journal again. IdlePollInterval caps that wait. It applies to every AddPipeline<TEvent> / AddPipeline<TSubscription, TEvent> pipeline, every journal live tap, and everything built on them (projections, journal message producers, process-manager routes). Set the default for all of them once; PipelineHostingOptions.IdlePollInterval and LiveTapHostingOptions.IdlePollInterval override it per registration:
builder.Services.AddNagareStreams(o => o.IdlePollInterval = TimeSpan.FromSeconds(5));The default is 100ms. How far to raise it depends on the wake-up:
- PostgreSQL (LISTEN/NOTIFY). A NOTIFY ends the wait at once, including one that arrives between the read and the wait, and a (re)established LISTEN forces a re-read. The interval only bounds latency when a NOTIFY is genuinely lost, so seconds are safe. It also sets the idle load: every caught-up pipeline issues one journal query per interval, so 100 pipelines at
100msare about 1,000 queries per second doing nothing. - MySQL, SQL Server, SQLite (no push wake-up). The interval is the polling interval, so it is added latency for every new event. Keep it low. Nagare logs a warning when a pipeline is configured above
100mswithout a push wake-up.
PollDelay does not affect the idle poll.
Failure handling is the conveyor-belt default: when a handler throws, the pipeline blocks hard at that element. The checkpoint never advances past it, nothing behind it is processed, and the same element is retried with exponential backoff (up to MaxFailureBackoff) until it succeeds — no attempt budget, no skipping. The block is per-pipeline and loud, not fatal: other pipelines keep running, and after FaultAfterConsecutiveFailures failing sessions the readiness health check flips so the stall is visible. Skipping an element is only possible by explicitly setting a DeadLetterPolicy — see Dead-letter handling.
A fault clears, and the backoff resets, as soon as the pipeline shows it has recovered. It does not wait for the session to end: a live journal or Kafka session only returns when cancelled. Recovery means either:
- an element is handled (or parked under a
DeadLetterPolicy) and its offset saved, or - the session runs for
FaultClearAfterwith no element in flight: a live source that recovered and is waiting at the head has nothing to handle.
When the last failure was on a specific element, only that element counts: the fault clears when the element at that same offset is handled or parked. Elements from other partitions of a multi-partition source, or messages a transport replays on re-subscribe, do not clear it, and neither does the idle timer. When the failure was not on an element (the source, the offset store or the connection failed), the first element handled clears it.
A fault on an element also clears when the checkpoint moves past it without this node handling it: another node took the lock and handled or parked the element, or an operator moved the checkpoint (for example from the dashboard). A node on standby clears the fault as soon as it sees the checkpoint has moved. A node that loads a moved checkpoint at the start of a session stops waiting for the failed element, and the rules above apply to that session. Offsets are compared by value: by Equals, and failing that by their System.Text.Json serialisation, so an offset class or a record holding a collection works even though the JSON offset stores return a new instance on every load. An offset that cannot be serialised is never treated as moved, and the fault stays.
A node on standby also drops a fault it cannot resolve that way once it has been on standby, without a failed lock attempt or checkpoint read, for FaultClearAfter. That covers a fault with no failed element (the lock store, source or offset store was unreachable, and the cluster has since recovered with another node holding the lock) and a failed element on an offset store that always loads null (Kafka group commits and TransportOffsetStore, where the broker tracks the position, so this node cannot see whether the checkpoint moved). If the problem persists, the node holding the lock reports its own fault. A standby node that reads an unmoved checkpoint keeps its fault. Clearing on standby never resets the consecutive-failure count: only progress on this node does. Two nodes that take turns failing on the same element therefore still reach FaultAfterConsecutiveFailures, and a node that takes the lock again and fails re-faults at once. A node on standby polls for the lock every PollDelay, whatever its failure count, so it takes over promptly when the holder releases.
A poison element never clears the fault. It fails before its offset is saved, and while it is in flight (including dead-letter retries) the idle timer does not count. The idle clear is provisional until a second FaultClearAfter passes. If the session fails first, the forgiven failures come back, so a source that blocks past FaultClearAfter and then fails re-faults immediately rather than flapping to healthy. Keep FaultClearAfter above the longest a source read can block before it surfaces an error, which is the connect timeout plus the command timeout.
Database-specific subscription registration
// Registers the subscription with the correct checkpoint store for the database
builder.Services.AddPostgresSubscription<BookProjection, BookEvent>();
builder.Services.AddSqlServerSubscription<BookProjection, BookEvent>();
builder.Services.AddSqliteSubscription<BookProjection, BookEvent>();Repository stores
Register read model repositories for dependency injection and health checks:
// Interface + implementation
builder.Services.AddRepositoryStore<IBookRepository, BookRepository>();
// Or just the implementation
builder.Services.AddRepositoryStore<BookRepository>();Both overloads also register the RepositoryStorageReadyHealthCheck.
Event upcasters
Register upcasters individually or as a chain:
// Individual registration (rebuilds the chain on each call)
builder.Services.AddEventUpcaster<BookAddedV1Upcaster>();
builder.Services.AddEventUpcaster<BookAddedV2Upcaster>();
// Or register the full chain at once
builder.Services.AddEventUpcasters(
new BookAddedV1Upcaster(),
new BookAddedV2Upcaster());See Event Versioning for details on writing upcasters.
Command middleware
builder.Services.AddCommandMiddleware<AuditMiddleware>();
builder.Services.AddCommandMiddleware<AuthorizationMiddleware>();Middleware executes in registration order. The first registered middleware is the outermost wrapper.
See Middleware for patterns and examples.
Messaging
Kafka transport
builder.Services.AddNagareKafkaTransport(builder.Configuration);This reads from the "Kafka" section of IConfiguration:
{
"Kafka": {
"BootstrapServers": "localhost:9092",
"Namespace": "my-app"
}
}| Property | Type | Required | Purpose |
|---|---|---|---|
BootstrapServers | string? | One of these | Standard Kafka bootstrap servers |
FullyQualifiedNamespace | string? | is required | Azure Event Hubs namespace (uses Azure Identity for auth) |
Namespace | string | Yes | Prefix for topic names and consumer groups |
Properties | Dictionary<string, string> | No | Additional Kafka client properties (SASL, SSL, etc.) |
Recovery | KafkaRecoveryOptions | No | When a client cut off from the brokers is rebuilt |
For Azure Event Hubs, use FullyQualifiedNamespace instead of BootstrapServers:
{
"Kafka": {
"FullyQualifiedNamespace": "my-eventhub.servicebus.windows.net",
"Namespace": "my-app"
}
}KafkaRecoveryOptions
{
"Kafka": {
"BootstrapServers": "broker:9092",
"Namespace": "my-app",
"Recovery": {
"RebuildAfter": "00:01:30",
"MaxRebuildInterval": "00:10:00"
}
}
}| Property | Type | Default | Purpose |
|---|---|---|---|
RebuildAfter | TimeSpan | 90s | How long a client must be continuously cut off before its librdkafka handle is rebuilt |
MaxRebuildInterval | TimeSpan | 10m | Nominal ceiling on the gap between rebuilds while an outage lasts; the gap doubles up to it |
Both are then spread by up to 25% either way, so a fleet does not rebuild in lockstep against a returning broker — which means an actual gap can sit 25% above MaxRebuildInterval.
Keep RebuildAfter comfortably above session.timeout.ms (30s as the transport configures it). A replaced consumer handle is destroyed without leaving its group, so the broker evicts it after that timeout; rebuilding faster than it can be evicted leaves the group in permanent rebalance. See broker outages.
The transport sets statistics.interval.ms to 15s on the producer, which is what lets an idle publisher prove it is connected: a sample showing a learned broker up is the only positive evidence a client with no traffic can give. Override it to 0 through Properties and a publisher that was cut off stays cut off until its next send tells it otherwise.
MessageSubscriptionOptions
builder.Services.AddMessageSubscription<
InventoryBookHandler,
BookIntegrationEvent,
BookChannel>(options =>
{
options.BatchSize = 100;
});| Property | Type | Default | Purpose |
|---|---|---|---|
BatchSize | int | 100 | Acknowledged messages between broker offset commits |
DeadLetter | DeadLetterPolicy? | null | Null = conveyor-belt default (block hard and retry); set to opt into retry-then-park |
Require<T>() | — | none declared | Declares the union members this subscription reads; undeclared members are skipped as routing, not failure |
Failure handling follows the same conveyor-belt default as every other pipeline — see Messaging → Failure handling.
The consumer group is declared on the handler itself, via the GroupId property on IMessageHandler<T>:
public class InventoryBookHandler : IMessageHandler<BookIntegrationEvent>
{
public string GroupId => "inventory.book-availability";
public Task Handle(BookIntegrationEvent message, MessageContext context, CancellationToken ct)
{
// ...
}
}Use a capability-based prefix (e.g. "inventory.book-availability", "reference-data.carer") rather than a service-based name. If the handler is later moved to a different service during a split, its group ID — and therefore its consumed position — transfers naturally. See Messaging → The GroupId property for the full rationale.
Message producer registration
builder.Services.AddMessageProducer<
BookMessageMapper, // IMessageMapper implementation
BookEvent, // Source event type (from event store)
BookChannel, // Channel definition
BookIntegrationEvent // Target message type (published to transport)
>();See Messaging for the full producer/consumer pattern.
Lock providers
For single-instance deployments, the default NoopLockProvider is fine. For multi-instance deployments, register a distributed lock provider:
// PostgreSQL (registered automatically by AddPostgresCheckpointStore)
services.AddSingleton<ILockProvider>(sp =>
new PostgresLockProvider(connectionString));
// SQL Server (registered automatically by AddSqlServerCheckpointStore)
services.AddSingleton<ILockProvider>(sp =>
new MsSqlLockProvider(connectionString));Lock providers prevent multiple instances of the same service from processing the same subscription events simultaneously.
JSON serialization
Nagare uses System.Text.Json with these defaults:
PropertyNameCaseInsensitive: trueDefaultIgnoreCondition: WhenWritingNullReferenceHandler: IgnoreCyclesJsonStringEnumConverterfor enum serialization
These settings are applied to event serialization and document store serialization. If you need custom converters, register them through NagareJsonSettings before adding your event stores.
Full example
A complete Program.cs for a service using PostgreSQL:
var builder = WebApplication.CreateBuilder(args);
// Database
builder.Services.AddNagarePostgresStorage(builder.Configuration);
builder.Services.AddPostgresCheckpointStore();
// Aggregate (includes event store + snapshot store)
builder.Services.AddPostgresAggregate<
BookAggregate, BookCommand, BookEvent, BookState>();
// Subscriptions
builder.Services.AddPostgresSubscription<BookProjection, BookEvent>();
builder.Services.AddPipeline<BookProjection, BookEvent>(options =>
{
options.BatchSize = 25;
options.PollDelay = TimeSpan.FromMilliseconds(50);
});
// Read models
builder.Services.AddRepositoryStore<IBookRepository, BookRepository>();
// Middleware
builder.Services.AddCommandMiddleware<CorrelationMiddleware>();
// Upcasters
builder.Services.AddEventUpcaster<BookAddedV1Upcaster>();
// Messaging (Kafka)
builder.Services.AddNagareKafkaTransport(builder.Configuration);
builder.Services.AddMessageProducer<
BookMessageMapper, BookEvent, BookChannel, BookIntegrationEvent>();
// Observability
builder.Services.AddOpenTelemetry()
.WithTracing(tracing => tracing
.AddSource("Nagare")
.AddOtlpExporter());
var app = builder.Build();
app.MapHealthChecks("/health/ready");
app.MapGet("/books/{id}", async (string id, IBookRepository repo) =>
{
var book = await repo.GetById(id);
return book is not null ? Results.Ok(book) : Results.NotFound();
});
app.Run();