Messaging
Within a bounded context, projections subscribe directly to the event store. No broker needed. But when events need to cross context boundaries, you need durable delivery. If the process crashes between emitting and handling, an in-memory approach loses the event.
Nagare's messaging layer provides transport-agnostic abstractions for publishing and consuming integration events. You write handlers and mappers once; the transport is a single registration line.
The two packages
| Package | Purpose |
|---|---|
Nagare.Messaging | Channel, publisher, handler, and mapper abstractions |
Nagare.Messaging.Kafka | Apache Kafka transport using the Confluent client |
dotnet add package Nagare.Messaging
dotnet add package Nagare.Messaging.KafkaDefining a channel
A channel names the conduit that messages flow through. In Kafka this maps to a topic. In other transports it maps to whatever the equivalent is — a queue, an exchange, a subscription.
public class BookChannel : IMessageChannel<BookIntegrationEvent>
{
public string ChannelName => "book-integration-events";
}The generic parameter is the message type that flows through this channel.
Publishing events
The typical pattern is a mapper that reads domain events from the event store and transforms them into integration messages. Define the mapper:
public class BookMessageMapper : IMessageMapper<BookEvent, BookIntegrationEvent>
{
public BookIntegrationEvent? Map(BookEvent @event, EventMapperContext context)
{
return @event switch
{
BookBorrowed e => new BookBorrowedIntegration(
context.AggregateId, e.BorrowerId),
BookReturned _ => new BookReturnedIntegration(context.AggregateId),
_ => null // filter out events we don't publish
};
}
}Returning null filters an event out — it won't be published.
The EventMapperContext gives you the aggregate ID, position, timestamp, and correlation metadata from the source event, without coupling the mapper to the full EventEnvelope type.
Registration
builder.Services.AddMessageProducer<
BookMessageMapper, // the mapper
BookEvent, // source (domain) event type
BookChannel, // the channel
BookIntegrationEvent // target (integration) message type
>();This registers a subscription on the BookEvent stream. Each event passes through the mapper; non-null results are published to the channel.
Consuming messages
On the other side, a bounded context subscribes to the channel and reacts:
public class InventoryBookHandler : IMessageHandler<BookIntegrationEvent>
{
private readonly IAggregateRepository<InventoryAggregate,
InventoryCommand, InventoryEvent, InventoryState> _repo;
public InventoryBookHandler(IAggregateRepository<InventoryAggregate,
InventoryCommand, InventoryEvent, InventoryState> repo) => _repo = repo;
public string GroupId => "inventory.book-availability";
public async Task Handle(
BookIntegrationEvent message, MessageContext context, CancellationToken ct)
{
if (message is BookBorrowedIntegration borrowed)
{
var aggregate = await _repo.Load(new AggregateId(borrowed.BookId));
await aggregate.Ask(new MarkAsUnavailable());
}
}
}The handler receives the message directly — no envelope wrapping. The MessageContext carries the channel name, message ID, timestamp, and any metadata.
The GroupId property
Every handler declares its own GroupId. This is the competing-consumer group the transport uses (Kafka consumer group, Service Bus subscription, etc.). Multiple instances of the same service share the same group, so each message is processed once.
The group ID is a property of the handler class itself — not a registration-time option. This matters when services get split:
- If you name the group after a service (e.g.
"inventory-service"), moving the handler into a new service forces you to either keep the old name (misleading) or change it (reset offsets, replay everything). - If you name the group after the capability (e.g.
"inventory.book-availability","reference-data.carer","availability-reactions.absence"), the group travels with the handler. When the handler moves to a new service, its consumed position moves with it.
Pick capability-based names. Your future self will thank you when a service splits.
Registration
builder.Services.AddMessageSubscription<
InventoryBookHandler,
BookIntegrationEvent,
BookChannel
>(options =>
{
options.BatchSize = 100;
});Failure handling: the conveyor belt
A subscription is a conveyor belt, and a conveyor belt stops when something on it breaks. When a handler throws, the default is block hard and retry: the message is not acknowledged, the subscription halts at it, and the same message is retried with exponential backoff until it succeeds or someone fixes the cause. No message is ever skipped, and the stall is surfaced through the readiness health check and the dashboard rather than buried in a log line. This matches every other pipeline in the framework — projections and process-manager event routes behave identically.
Skipping is never the default anywhere; it is always an explicit opt-in. The one message-level exception is routing, not failure: a message whose discriminator sits outside the subscription's declared subtype set (Require<T>()) is not for this consumer, so it is skipped with its offset advancing.
Where losing a message is acceptable and a poison message must not hold up the ones behind it, opt into dead-lettering:
builder.Services.AddMessageSubscription<
InventoryBookHandler,
BookIntegrationEvent,
BookChannel
>(options =>
{
options.DeadLetter = new DeadLetterPolicy
{
MaxRetries = 5,
// 1s, 5s, 30s, 2m, 10m by default; the last entry repeats.
BackoffSchedule = DeadLetterPolicy.DefaultSchedule,
};
});Now a failing message is retried with the policy's backoff schedule. Once retries are exhausted, the message is parked in the registered IStreamDeadLetterStore — type, payload, exception, and attempt count intact — and the belt keeps moving.
Where parked messages land depends on your storage registration. Register a relational dead-letter store (shipped with all four backends — PostgreSQL, MySQL, SQL Server, SQLite) and parked messages are durable: they survive restarts and show up in the dashboard's stream detail view. Without one, they park in process memory. One exception either way: transport offsets (Kafka partition/offset tuples) aren't numeric, so broker-sourced messages park in the store's in-memory fallback even when a relational store is registered — the records stay visible and removable through IStreamDeadLetterStore, just not across restarts.
Redriving a parked message
Redrive is a human decision: fix the root cause first, then replay.
- Remove the parked record — through the dashboard's dead-letter view or
IStreamDeadLetterStore.RemoveAsync. - Rewind the subscription so the message is consumed again. Broker offsets are a broker-side operation, so reset the consumer group (the checkpoint) from outside the process — e.g.
kafka-consumer-groups --reset-offsetsfor Kafka — or republish the message.
TransportOffsetStore.ResetAsync is deliberately a no-op: it keeps the IOffsetStore<T> contract without pretending it can rewind a broker it doesn't control.
Kafka transport
builder.Services.AddNagareKafkaTransport(builder.Configuration);In appsettings.json:
{
"Kafka": {
"BootstrapServers": "localhost:9092",
"Namespace": "my-app"
}
}For Azure Event Hubs, use the fully qualified namespace:
{
"Kafka": {
"FullyQualifiedNamespace": "my-eventhub.servicebus.windows.net",
"Namespace": "my-app"
}
}Additional Kafka properties (SASL, SSL) go in the Properties dictionary:
{
"Kafka": {
"BootstrapServers": "broker:9092",
"Namespace": "my-app",
"Properties": {
"security.protocol": "SASL_SSL",
"sasl.mechanism": "PLAIN",
"sasl.username": "$ConnectionString",
"sasl.password": "your-connection-string"
}
}
}Broker outages and recovery
Confluent's client only registers librdkafka's error callback when you give the builder an error handler. Nagare gives it one, so authentication failures, transport failures, connection-setup timeouts and _ALL_BROKERS_DOWN reach your ILogger instead of stderr. Without that they are invisible to the process: Consume returns null during a total outage, exactly as it does on a quiet topic.
A client that stays cut off gets a new librdkafka handle, in place. The ITransportConsumer you were handed does not change identity, so nothing that holds it — the DI container, the offset store, the pipeline source, the dashboard's telemetry view — needs to know. Nothing restarts.
- Cut off starts at the first broker-level error and is only ever cleared by positive evidence: a consumer group assignment, a record consumed, a record confirmed persisted, or a producer statistics sample showing a learned broker up. librdkafka suppresses an identical error for 30 s and raises
_ALL_BROKERS_DOWNonce per down-transition, so the error stream goes quiet in the middle of an outage — "nothing failed lately" is not recovery. - The rebuild happens after
Recovery.RebuildAfter(default 90 s), on the thread that owns the handle, then backs off towardsRecovery.MaxRebuildIntervalwhile the outage lasts. The replaced handle is destroyed on its own thread, because destroying one whose broker is gone takes minutes. - A fatal error rebuilds at once. librdkafka marks the handle unusable, so the threshold buys nothing.
- A credential the broker rejected is not a cut-off.
SaslAuthenticationFailedkeeps faulting the pipeline loudly, because a new handle built from the same configuration cannot fix it.
The consumer keeps its group id, so a rebuilt handle resumes from the group's committed offsets — no offsets are migrated and at-least-once still holds. Offsets that were stored but not yet committed are lost with the old handle, so those records are redelivered — a window that a flapping outage widens, since a commit skipped while cut off still counts against MessageSubscriptionOptions.BatchSize. While a consumer knows it is cut off it does not commit at all: Commit is synchronous with no timeout and blocks for the length of the outage. Stored offsets are cumulative, so the next commit after recovery carries them — which also means a commit is never what ends a cut-off. On a channel too quiet to deliver anything, the rebuild's rejoin is.
On the publisher side, a send that is in flight when the brokers go away fails, and the caller retries it. Nothing is dropped silently: a replaced producer is flushed and kept alive until its own sends have drained, because destroying one first would leave those callers awaiting a delivery report that never comes. How long that takes is librdkafka's message.timeout.ms, five minutes by default; lower it through Properties if a caller should find out sooner.
Replay after a rebuild
A subscription that is rebuilt before its first commit has no committed offset, so auto.offset.reset=earliest replays the channel from the retention floor. This is the same behaviour a pod restart has always had, and handlers must be idempotent regardless — but with many clients it makes the moment the broker returns a load event.
Delivery semantics
Supported: multi-partition topics, consumer groups (GroupId per handler — competing consumers within a group, isolated feeds across groups), and at-least-once delivery. Acks flow through StoreOffset and commits are explicit; a handler that fails never acks, so the next group member is redelivered the message. Make handlers idempotent.
Not supported: exactly-once. No transactions, no EOS producer/consumer wiring, no atomic consume-process-produce. If a crash between side effect and commit is unacceptable for a channel, that channel's handler must tolerate duplicates — or the design belongs on the event store, not the broker.
Threading
Each subscription's consumer owns one dedicated background thread, and every call into librdkafka — consume, commit, offset store, the dashboard's telemetry reads — runs on it. Confluent's client is synchronous throughout and offers no async form, so this is what keeps the blocking off the thread pool: a host with 30 subscriptions costs 30 dedicated threads that sit idle most of the time, rather than 30 pool threads that Kestrel and your handlers can no longer use. Sizing follows from that — subscription count drives thread count, not CPU. Each thread is named kafka: <topic> (<group>), which dotnet-stack, dotnet-dump and the Datadog profiler show in full; top -H and ps show only the first 15 bytes, because Linux truncates the kernel's copy of the name.
Because one thread owns the handle, calls from other threads queue behind whatever it is currently doing, and a consume on an idle topic waits up to 100 ms. Under steady traffic consume returns as soon as a message is there, so acks and commits barely wait. When traffic stops, the waits show up:
- Commits — the last ack and commit of a burst, and any commit issued while the topic is quiet, can each wait up to 100 ms behind a consume before they run. Offsets still reach the broker; they just trail the handler by up to that long.
- Dashboard telemetry — a Messages page load waits up to 100 ms per channel (channels are read concurrently, so the page costs the slowest one). The read itself is cheap: watermarks come from librdkafka's cache and committed offsets from what the consumer committed, with one broker round-trip per partition the first time it is read after assignment. An aborted request stops waiting, and its read is skipped if it has not started.
That wait inside consume is the subscription's only idle pause. A caught-up subscription gets a new message as soon as consume receives it, and polls again straight away when consume returns empty.
Cancellation never interrupts work already on the consumer's thread — librdkafka takes no token. A cancelled Poll returns no message, as the in-memory transport does; if the consume it started goes on to receive one, that message is kept for the next Poll rather than dropped. Acknowledge and commit ignore the token, so shutdown cannot land between them.
On shutdown the consumers are stopped together once the host's hosted services have finished, each closing (and leaving its group) within one shared ten-second budget. The consumers and their factory implement IAsyncDisposable, so the host's container waits for those closes without holding a thread. A close that fails or overruns is logged; it never throws out of the container's disposal. A stopped consumer still answers Poll, with no message, after the same short wait an idle consume takes, so a pipeline that is still running against it doesn't spin.
How it fits together
The full flow from domain event to cross-context reaction:
- Aggregate emits
BookBorrowed(domain event, stored in the event store) MessageProducerpicks it up via an event store subscription, maps it toBookBorrowedIntegrationBookBorrowedIntegrationis published to thebook-integration-eventschannelInventoryBookHandlerconsumes it and issues a command to the Inventory aggregate
Domain events stay internal. Integration events cross boundaries. The transport provides durability, ordering, and replay.
Why a broker and not a direct subscription? If you use in-process subscriptions between bounded contexts, you create invisible coupling — one context reading directly from another's event store. When you eventually extract a context into its own service, you discover that half its behaviour depends on another context's internal streams. That's not a boundary, it's a shared database with extra steps. A broker enforces the separation from day one. See Modular Monolith — Two levels of communication for the full architectural argument.
Swapping transports
The channel, mapper, and handler stay the same regardless of transport. Only the registration line changes:
// Apache Kafka
builder.Services.AddNagareKafkaTransport(builder.Configuration);
// In-memory for integration tests (no broker needed)
builder.Services.AddInMemoryTransport();
// Future: Azure Service Bus, RabbitMQ
// builder.Services.AddNagareServiceBusTransport(builder.Configuration);Testing with InMemoryTransport
For integration tests, replace Kafka with the built-in in-memory transport. Messages flow through System.Threading.Channels — no broker, no containers, no latency.
builder.Services.AddInMemoryTransport();The InMemoryTransport instance is available from DI for test assertions:
var transport = serviceProvider.GetRequiredService<InMemoryTransport>();
// Assert messages were published
var published = transport.GetPublished("book-integration-events");
Assert.Single(published);
// Wait for async producers to catch up
var messages = await transport.DrainAsync("book-integration-events", expected: 3);
// Reset between tests
transport.Clear();Each consumer gets its own copy of every message (fan-out), matching Kafka's consumer-group semantics.
Domain events are private
This is the most important rule in the messaging layer, and the one most likely to be violated when moving fast.
Domain events (BookBorrowed, BookReturned, BookLost) belong to the aggregate. They are the source of truth in the event store, they feed projections within the context, and they can evolve freely through upcasters. No other context can depend on them.
Integration events (BookBorrowedIntegration) are a public contract. They live in a shared contract assembly, they travel over the broker, and changing their schema is a breaking change for every consumer.
These are separate types for a reason. A domain event carries whatever the aggregate needs for internal state transitions. An integration event carries whatever the consuming context needs — often less, sometimes shaped differently. The mapper is the translation layer between the two.
Internal Public
──────── ──────
BookAdded(Title, Author, Isbn) → (not published — internal to the library context)
BookBorrowed(BorrowerId, At) → BookBorrowedIntegration(BookId, BorrowerId)
BookReturned(ReturnedAt) → BookReturnedIntegration(BookId)
BookLost(ReportedAt) → (not published — handled internally)If you find yourself passing a domain event type to AddMessageProducer as both the source and the target, you're exposing internals. Create a separate integration type.
Design guidelines
Keep integration events coarse-grained. They're contracts between bounded contexts.
BookBorrowedis useful.BookBorrowerIdFieldUpdatedis not.Don't share types across contexts. The publishing context defines the integration event. The consuming context can define its own type that deserializes from the same JSON.
Use the mapper to filter. Not every domain event needs to leave the context. Return
nullfromMapfor events that are purely internal.One channel per bounded context is a good starting point. You can split later if different consumers need different ordering guarantees.
Durability
The publish itself is brokered by the transactional outbox: the integration event lands in the outbox table in the same transaction as the underlying domain events, and a background dispatcher carries it to the broker with at-least-once semantics. That's why publishing doesn't need a try/catch around the broker — the dispatcher retries on a backoff ramp (up to OutboxRunnerOptions.MaxAttempts, default 10) and parks the row in the dead-letter state for triage if the budget is exhausted.