Skip to content

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:

csharp
// 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:

json
{
  "ConnectionStrings": {
    "DbConnectionString": "Host=localhost;Database=myapp;Username=postgres;Password=secret"
  }
}

You can pass a different name:

csharp
builder.Services.AddNagarePostgresStorage(
    builder.Configuration,
    connectionStringName: "EventStoreDb");

Event store ​

Register an event store for each event type in your system:

csharp
// 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:

csharp
builder.Services.AddPostgresEventStore<BookEvent>(tableName: "BookEvents");

Snapshot store ​

csharp
builder.Services.AddPostgresSnapshotStore<BookState>();

Default table name is Snapshots. Override it the same way:

csharp
builder.Services.AddPostgresSnapshotStore<BookState>(tableName: "BookSnapshots");

Checkpoint store ​

Each database adapter registers its own checkpoint store and lock provider:

csharp
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:

csharp
builder.Services.AddPostgresAggregate<
    BookAggregate, BookCommand, BookEvent, BookState>();

This is equivalent to:

csharp
builder.Services.AddPostgresEventStore<BookEvent>();
builder.Services.AddPostgresSnapshotStore<BookState>();
builder.Services.AddAggregate<BookAggregate, BookCommand, BookEvent, BookState>();

AggregateOptions ​

Control how aggregates load their event history:

PropertyTypeDefaultPurpose
InitializeBatchSizeint50Number 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:

csharp
builder.Services.AddPipeline<BookProjection, BookEvent>(new PipelineHostingOptions
{
    PollDelay = TimeSpan.FromMilliseconds(200),
    MaxFailureBackoff = TimeSpan.FromSeconds(60),
});
PropertyTypeDefaultPurpose
PollDelayTimeSpan100msWait between sessions: after a session ends, or while another node holds the pipeline's lock
IdlePollIntervalTimeSpan?nullJournal pipelines: longest a caught-up source waits for a wake-up before re-reading the journal. null takes the global default (below)
MaxFailureBackoffTimeSpan30sCap for the exponential backoff between consecutively failing sessions
FaultAfterConsecutiveFailuresint5Failing sessions before the readiness health check reports the pipeline as faulted
FaultClearAfterTimeSpan60sHow 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)
StandbyHolderStaleAfterTimeSpan90sWith 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
ResilienceResiliencePipeline?nullOptional Polly pipeline wrapping each processing session
DeadLetterPolicyDeadLetterPolicy?nullNull = 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:

csharp
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 100ms are 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 100ms without 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 FaultClearAfter with 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 ​

csharp
// 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:

csharp
// 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:

csharp
// 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 ​

csharp
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 ​

csharp
builder.Services.AddNagareKafkaTransport(builder.Configuration);

This reads from the "Kafka" section of IConfiguration:

json
{
  "Kafka": {
    "BootstrapServers": "localhost:9092",
    "Namespace": "my-app"
  }
}
PropertyTypeRequiredPurpose
BootstrapServersstring?One of theseStandard Kafka bootstrap servers
FullyQualifiedNamespacestring?is requiredAzure Event Hubs namespace (uses Azure Identity for auth)
NamespacestringYesPrefix for topic names and consumer groups
PropertiesDictionary<string, string>NoAdditional Kafka client properties (SASL, SSL, etc.)
RecoveryKafkaRecoveryOptionsNoWhen a client cut off from the brokers is rebuilt

For Azure Event Hubs, use FullyQualifiedNamespace instead of BootstrapServers:

json
{
  "Kafka": {
    "FullyQualifiedNamespace": "my-eventhub.servicebus.windows.net",
    "Namespace": "my-app"
  }
}

KafkaRecoveryOptions ​

json
{
  "Kafka": {
    "BootstrapServers": "broker:9092",
    "Namespace": "my-app",
    "Recovery": {
      "RebuildAfter": "00:01:30",
      "MaxRebuildInterval": "00:10:00"
    }
  }
}
PropertyTypeDefaultPurpose
RebuildAfterTimeSpan90sHow long a client must be continuously cut off before its librdkafka handle is rebuilt
MaxRebuildIntervalTimeSpan10mNominal 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 ​

csharp
builder.Services.AddMessageSubscription<
    InventoryBookHandler,
    BookIntegrationEvent,
    BookChannel>(options =>
{
    options.BatchSize = 100;
});
PropertyTypeDefaultPurpose
BatchSizeint100Acknowledged messages between broker offset commits
DeadLetterDeadLetterPolicy?nullNull = conveyor-belt default (block hard and retry); set to opt into retry-then-park
Require<T>()—none declaredDeclares 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>:

csharp
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 ​

csharp
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:

csharp
// 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: true
  • DefaultIgnoreCondition: WhenWritingNull
  • ReferenceHandler: IgnoreCycles
  • JsonStringEnumConverter for 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:

csharp
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();

流れ — flow.