Observability
Nagare instruments itself with System.Diagnostics.Activity and System.Diagnostics.Metrics — the same primitives ASP.NET Core and EF Core use. Any OpenTelemetry collector picks them up.
There is no Nagare.OpenTelemetry package, and there isn't going to be one. The ActivitySource and Meter are part of the core library so a single .AddSource("Nagare") and .AddMeter("Nagare") is enough — same pattern as Marten, Wolverine, and MassTransit.
Tracing setup
Register Nagare's activity source and meter with your OpenTelemetry configuration:
builder.Services.AddOpenTelemetry()
.WithTracing(tracing => tracing
.AddSource("Nagare")
.AddOtlpExporter())
.WithMetrics(metrics => metrics
.AddMeter("Nagare")
.AddOtlpExporter());That captures every span and counter the framework emits. The source name is also exposed as NagareActivity.SourceName if you prefer not to hard-code the string.
Traced operations
Nagare creates spans for six operations:
| Span name | Kind | When |
|---|---|---|
nagare.aggregate.ask | internal | A command is sent to an aggregate |
nagare.eventstore.append | producer | Events are written to the store |
nagare.eventstore.read | internal | Events are read from the store |
nagare.subscription.handle | consumer | A subscription processes an event |
nagare.subscription.checkpoint | internal | A subscription saves its position |
nagare.outbox.dispatch | consumer | The outbox runner delivers a side-effect |
A typical command flow produces three nested spans: aggregate.ask wraps eventstore.read (loading the aggregate) and eventstore.append (writing new events). Subscription processing produces subscription.handle with periodic subscription.checkpoint spans. Outbox delivery produces outbox.dispatch per row, linked back to the original eventstore.append via the stored traceparent.
Tags
Each span carries attributes that let you filter and group traces:
| Tag | Value |
|---|---|
nagare.aggregate.id | The aggregate instance ID |
nagare.aggregate.type | The aggregate's type name |
nagare.event.type | The event stream's type name |
nagare.event.count | Number of events in this operation |
nagare.command.type | The command's type name |
nagare.subscription.id | The subscription's identifier |
nagare.position | Global event store position |
nagare.version | Aggregate version number |
nagare.outbox.dispatch_id | Stable id of the outbox row being delivered |
nagare.outbox.target | Logical name of the destination (handler/sink) |
nagare.outbox.attempt | Attempt counter — 1 on first try |
nagare.outbox.outcome | dispatched, transient, or dead |
nagare.outbox.source_process | Process aggregate that produced the dispatch, if any |
nagare.replay | true when the handler was driven by a catch-up replay |
You can find every command handled by a given aggregate, trace from the HTTP request through the command to the events it produced, and follow those events through subscriptions and outbox deliveries.
Snapshot write failures
Snapshot writes happen after the events are committed, so a failed snapshot cannot fail the command — rolling back would make the caller retry and double-apply events that are already durable. Instead, Nagare records a nagare.snapshot.write_failed event on the current span and leaves the last-snapshot marker unchanged, so the next command re-attempts the snapshot. Snapshots are only a read-time optimisation; the events are the source of truth.
The event carries these tags:
| Tag | Value |
|---|---|
aggregate.id | The aggregate instance ID |
snapshot.version | The version the snapshot was being written at |
exception.type | Full name of the exception thrown by the snapshot store |
exception.message | The exception message |
This is the only observability surface for snapshot failures — Nagare deliberately takes no ILogger dependency, so there is no log entry. The event appears on the span that was current when the command was persisted (usually nagare.aggregate.ask), whenever an ActivitySource("Nagare") listener is active. Alert on it like any other span event.
One exception to the swallow-and-continue rule: cancellation. An OperationCanceledException from the snapshot write propagates to the caller — a cancelled command must not report success.
Metrics
Counters and a histogram are exposed under the same Nagare meter:
| Instrument | Type | Description |
|---|---|---|
nagare.outbox.dispatched | counter | Rows successfully dispatched |
nagare.outbox.dead | counter | Rows that exhausted retries and were marked Dead |
nagare.outbox.attempts | counter | Dispatch attempts, tagged by outcome |
nagare.outbox.dispatch.duration_ms | histogram | Per-row dispatch latency |
nagare.outbox.retention.deleted | counter | Rows removed by the retention sweep |
Tag the dashboards by nagare.outbox.target to break the totals down per sink.
Linking subscription spans to writes
Events written today might be processed by a subscription seconds later — or replayed by a projection rebuild years later. Naive parent-child propagation breaks the second case: the original write trace is long gone from the backend, leaving handler spans pointing at nothing.
Because of that, handler spans are deliberately not parented to the writing span. This matches the OpenTelemetry messaging convention, stays deterministic across replays, and degrades gracefully when the writer trace has been sampled out or aged out of the backend.
The write side is captured at the append site: the framework reads Activity.Current when Append is called, formats a W3C traceparent, and stores it in event metadata. (The writing process must have a listener active for the Nagare source at that moment; if no Activity is current, no traceparent is stored.) That stored traceparent is what powers outbox dispatch linking — each dispatch span attaches to the producing span across the async hop.
Pipeline handler spans are emitted as stream.handle on the Nagare.Streams source, tagged with nagare.stream_id, nagare.position, and nagare.event_type. They start as their own roots; to correlate a handler span back to the original write, follow the traceparent stored in the event's metadata.
For most users, the link default is the right call. Turn parent-child on when you specifically want the demo-friendly view and your write traces are reliably retained.
Event metadata
Beyond tracing spans, you can attach metadata to individual events. This metadata is persisted in the event store and available everywhere the event is read.
var metadata = new EventMetadata(
CorrelationId: requestId,
CausationId: $"http:{Request.Path}",
ActorId: currentUser.Id,
Headers: new Dictionary<string, string> { ["tenant"] = currentUser.TenantId },
Timestamp: DateTimeOffset.UtcNow);
await aggregate.Ask(new BorrowBook("user-42"), metadata);In projections, the metadata is available on the envelope:
public async Task Handle(EventEnvelope<BookEvent> envelope)
{
var actorId = envelope.Metadata?.ActorId;
var correlationId = envelope.Metadata?.CorrelationId;
var tenant = envelope.Metadata?.Headers?.GetValueOrDefault("tenant");
// ...
}What to put in metadata
| Field | Purpose | Example |
|---|---|---|
CorrelationId | Trace a chain of events back to the original request | HTTP request ID |
CausationId | Identify what caused this event | The command or event that triggered it |
ActorId | Stable id of whoever asked. Pairs with CommandSource. | user-42 for HTTP, a job name for the scheduler, a migration tag for system tasks. Prefer internal ids over PII — events are immutable. |
CommandType / CommandSource | Auto-attribution | Filled by Aggregate.Ask and ProcessGrain; you usually don't set these. |
Headers | Free-form IReadOnlyDictionary<string, string> for ride-along context | Tenant id, branch id, feature-flag bucket. Surfaces in dashboards and projections; isn't interpreted by the framework. |
Timestamp | Custom timestamp | Override the store's default timestamp |
TraceParent | W3C trace context captured at append time | Auto-populated; do not set by hand |
TraceState | Vendor-specific trace propagation companion to TraceParent | Auto-populated; do not set by hand |
TraceParent and TraceState are written by the framework at append time when an Activity is in scope — you don't set them. They drive the linking decision described above. Custom IEventMetadata types are passed through untouched, so spans-linking only works when you use the built-in EventMetadata record.
A middleware is a good place to attach metadata automatically:
public class CorrelationMiddleware(IHttpContextAccessor http) : ICommandMiddleware
{
public async Task<IReply> InvokeAsync(AskContext context, AskDelegate next)
{
var requestId = http.HttpContext?.TraceIdentifier;
var actorId = http.HttpContext?.User.FindFirst("sub")?.Value;
var enriched = context with
{
Metadata = new EventMetadata(
CorrelationId: requestId,
ActorId: actorId)
};
return await next(enriched);
}
}Register it once and every command carries correlation data.
Health checks
Nagare registers these health checks to report whether the system is ready to serve traffic.
Event store readiness
EventStoreReadyHealthCheck reports healthy once the event store's database table has been created and verified. It reports unhealthy during startup while the initialization service runs CREATE TABLE IF NOT EXISTS.
Subscription readiness
SubscriptionsReadyHealthCheck tracks each subscription individually. It reports healthy only when every registered subscription has completed its initial catch-up (replayed historical events up to the current position). During startup, it lists which subscriptions are still initializing.
This is useful for Kubernetes readiness probes. A service shouldn't receive traffic until its projections have caught up. Otherwise, queries against read models return stale or empty results.
Stream readiness
StreamsReadyHealthCheck covers every hosted pipeline (projections, message subscriptions, process-manager routes). Each stream is in one of four states, and separately may be faulted:
- Not ready: the pipeline hasn't started pumping for a reason other than another node holding its lock.
Preparehasn't finished, or no lock attempt has resolved yet (for example, the lock store has been unreachable since startup). - Standby: another node holds the pipeline's lock. This is how a multi-replica deployment runs, so a standby stream is healthy. It becomes ready when this node acquires the lock.
- Standby behind an unresponsive holder: another node holds the lock, but with the lease registry on, the holder has stopped renewing its lease: its
renewed_atis older thanStandbyHolderStaleAfter(default 90s), or no node has recorded a lease for the lock for that long. The stream is listed underholderUnresponsivewith the reason (for examplelock held by web-2/1, whose lease was last renewed 95s ago) and the check isDegraded. A PostgreSQL session lock outlives a holder whose machine or network died without closing its connection until the server notices, which with default TCP keepalives can take hours; without this check every surviving replica would report a healthy standby while nothing pumps. It is never reported as not ready: restarting a standby replica cannot free a lock the dead holder's connection still owns, and the same signal also fires for holders that are fine but not renewing (see below). Without the lease registry, or when the lock provider does not record leases, there is nothing to judge the holder by, and standby is always healthy. - Ready: this node holds the lock and is pumping.
- Faulted:
FaultAfterConsecutiveFailuressessions in a row have failed. The pipeline keeps retrying with backoff, and the fault clears once it makes progress again (see pipeline hosting). A fault describes this node's own attempts. On standby it clears when the checkpoint moves past the element it failed on, or, when there was no failed element or the checkpoint cannot show progress here, afterFaultClearAfterof unbroken standby. A standby node that can read an unmoved checkpoint keeps its element fault, because the pipeline is visibly still stuck there. The node holding the lock reports its own fault.
The check is Unhealthy if any stream is not ready. The description starts Streams not yet ready: and also lists faults and unresponsive holders when there are any. If no stream is not ready but some are faulted, it is Unhealthy with Streams faulted: …, followed by any unresponsive holders. If the only problem is standby streams behind an unresponsive holder, it is Degraded with Streams on standby behind an unresponsive lock holder: …. Otherwise, with every stream ready or on standby, it is Healthy. A fault never hides a stream that didn't start.
Base decisions on HealthCheckResult.Data rather than parsing the description. Matching the Streams faulted prefix misses faults whenever another stream is also not ready, because the description then starts Streams not yet ready:.
| Key | Type | Contents |
|---|---|---|
notReady (StreamsReadyHealthCheck.NotReadyKey) | string[] | Ids of streams that are not pumping and not on standby, sorted |
standby (StreamsReadyHealthCheck.StandbyKey) | string[] | Ids of streams whose lock another node holds, and whose holder looks alive, sorted |
faulted (StreamsReadyHealthCheck.FaultedKey) | IReadOnlyDictionary<string, string> | Faulted stream id → failure reason, sorted |
holderUnresponsive (StreamsReadyHealthCheck.HolderUnresponsiveKey) | IReadOnlyDictionary<string, string> | Standby stream id → why its lock holder looks unresponsive, sorted |
All four keys are always present, and empty when there is nothing to report. A stream is in at most one of notReady, standby and holderUnresponsive (a ready stream is in none), and may also be in faulted.
Alert on holderUnresponsive; do not restart for it. It means nothing may be pumping the stream anywhere, which needs a person: find the holder named in the reason, and if its node is gone, terminate its database session (on PostgreSQL, pg_terminate_backend on the session holding the advisory lock) so a standby replica can take the lock. Restarting the standby replicas does not help, and would restart all of them at once. Likewise ignore standby: treating it as fatal would restart every standby replica of a multi-replica service indefinitely. Whether a liveness probe treats notReady as fatal is a separate choice; a stream is also not ready while it warms up and while the lock store is unreachable.
Turn the lease registry on for every node of a service before relying on the holder check, and register it after the storage backend (see cluster); a pipeline whose lock provider does not record leases logs a warning at start and skips the check. A holder that does not write leases looks like a holder that stopped renewing, so its standby peers report holderUnresponsive after StandbyHolderStaleAfter. In particular, during a rolling upgrade from a Nagare version before 0.19.0 (which did not renew leases) standby streams on the new nodes may report holderUnresponsive (Degraded) until the old nodes are gone: an old holder never renews, and a new node's boot sweep deletes its row after 2 minutes. There is no way to tell such a holder from a dead one. Set StandbyHolderStaleAfter to Timeout.InfiniteTimeSpan to turn the check off for the rollout. The check compares the holder's renewed_at, stamped with the holder's clock, against this node's clock; keep clock skew between nodes well under the margin between the 30s renewal and StandbyHolderStaleAfter.
Repository storage readiness
RepositoryStorageReadyHealthCheck reports healthy once all document store tables have been created. Like the event store check, it transitions from unhealthy to healthy during startup.
Using health checks
The health checks are registered automatically when you add event stores, subscriptions, or repository stores. Wire them into ASP.NET Core's health check endpoint:
app.MapHealthChecks("/health/ready", new HealthCheckOptions
{
Predicate = _ => true
});In Kubernetes, point your readiness probe at this endpoint:
readinessProbe:
httpGet:
path: /health/ready
port: 8080
initialDelaySeconds: 5
periodSeconds: 10The service stays out of the load balancer until the event store is initialized, all subscriptions have caught up, and all document store tables exist.
Connecting the pieces
A production observability setup ties these together:
- Tracing shows you what happened: which command was issued, what events it produced, how long each step took
- Metadata shows you why: who issued the command, what request triggered it, what earlier event caused it
- Health checks show you readiness: is the system caught up and safe to serve traffic
The tracing spans and metadata flow into your existing observability stack (Datadog, Jaeger, Grafana Tempo, Azure Monitor). The health checks integrate with your existing orchestrator. Nagare doesn't impose its own monitoring layer. It fits into whatever you already run.