Messaging
Each test gets a broker client. Publish a message, then await the matching message with a predicate and a timeout:
[ProtoTest]
public async Task Paying_an_invoice_publishes_an_event()
{
var messages = Proto.Context.Messaging();
await messages.PublishAsync("invoice.paid", """{"id":42}""", contentType: "application/json");
var message = await messages.AwaitAsync(
"invoice.paid",
candidate => candidate.Payload!.Contains("\"id\":42"));
message.Should.MatchShape(new { id = 42 });
}
Run it with dotnet test. A green run prints the passed test. Failing to arrive is a TimeoutException and a test failure, not a sleep.
What it adds
Without an adapter the client runs against an in-memory broker, so the API works anywhere; ProtoTest.Messaging.RabbitMq replaces it with RabbitMQ, ProtoTest.Messaging.RabbitMq.Testcontainers owns a broker for the whole run, and ProtoTest.Messaging.MassTransit bridges the surface to an in-process application's MassTransit test harness - or, through its MassTransitEnvelope helper, speaks the MassTransit wire envelope over any adapter, published applications included.
| Package | Does |
|---|---|
ProtoTest.Messaging | the surface: publish, await, Tap/Declare, in-memory broker |
ProtoTest.Messaging.RabbitMq | a real RabbitMQ broker as the adapter |
ProtoTest.Messaging.RabbitMq.Testcontainers | a broker container owned by the run |
ProtoTest.Messaging.MassTransit | the application's in-process MassTransit harness as the adapter |
Which brokers
| Broker | Package | What it is |
|---|---|---|
| In-memory | ProtoTest.Messaging | The default test double. It keeps history for the run, so the API works with no broker. It registers no Broker capability, so capability-gated tests skip instead of passing against the double |
| RabbitMQ | ProtoTest.Messaging.RabbitMq (+ .Testcontainers to own one) | A real broker as the adapter, with per-test tap queues, suite-owned topology via Declare, and the container recipe under Owning a broker |
| MassTransit | ProtoTest.Messaging.MassTransit | The application's in-process MassTransit harness as the adapter, or the MassTransitEnvelope wire envelope over any adapter. See the MassTransit page |
Kafka, Azure Service Bus and AWS (MSK, SQS/SNS, EventBridge) are not shipped. There is no plan to announce here either: the extension path is the adapter seam. Implement IProtoMessageBroker (publish, one consumer per test, declare) with a consumer derived from ProtoMessageConsumerBase, register it with UseBroker, and the test-side API, the trace and the Broker capability come with it. The full contract is on the Adapter contract.
Install
dotnet add package ProtoTest.Messaging
dotnet add package ProtoTest.Messaging.RabbitMq
dotnet add package ProtoTest.Messaging.RabbitMq.Testcontainers
dotnet add package ProtoTest.Messaging.MassTransit
The packages target .NET 8, 9 and 10 (the project template defaults to net10.0; pass --framework net8.0 or --framework net9.0 for an older runtime). The first is the capability. The others are optional adapters: add RabbitMQ to talk to a real broker (with the container package when the run should start one), or ProtoTest.Messaging.MassTransit to use an in-process application's MassTransit test harness. See MassTransit. ProtoTest.Messaging.RabbitMq and ProtoTest.Messaging.MassTransit bring ProtoTest.Messaging with them.
Compose
builder.AddMessaging(); // in-memory broker
builder.AddMessaging(messaging => messaging.UseRabbitMq()); // RabbitMQ
builder.AddMessaging(messaging => messaging.UseMassTransit<Program>()); // the app's MassTransit harness
IProtoHostBuilder AddMessaging(this IProtoHostBuilder builder, Action<ProtoMessagingBuilder>? configure = null);
UseBroker is the adapter seam; UseRabbitMq and UseMassTransit are the built-in implementations of it.
Tap declares the destinations this suite awaits, in code, so an adapter can bind each test's tap during setup. See Destinations. Declare declares the destinations this suite owns, in code, so an adapter creates them during setup before any tap binds. See Suite-owned topology.
When an adapter is configured, AddMessaging also registers the Messaging capability with kind broker. It always registers the options, the per-test client initializer, and the run-scoped broker resource.
A repeated AddMessaging call runs its configure callback again, so a later call can add an adapter to an adapter-less first call or extend attachment options. Infrastructure stays idempotent: one options object, one broker holder, one initializer, one capability and one run resource. The first adapter configured wins.
A configured adapter is what makes the Broker capability true. The in-memory default registers none, so [RequiresCapability(ProtoCapabilityKinds.Broker)] skips where no real broker is configured instead of passing against the double (see skip conditions).
The tasks
The Paying_an_invoice_publishes_an_event test above is the whole pattern: publish, await the match, assert the shape.
Going further
RabbitMQ topology
| Fact | Rule |
|---|---|
| Publish target | the exchange named like the destination, under the message routing key or the destination itself |
| Content | ContentType defaults to application/json, messages are non-persistent, null headers are dropped |
| Delivery | the delivery routing key fills ProtoMessage.RoutingKey; a keyed await binds its key on the destination tap |
| Connection | one connection and one publish channel for the run, created lazily; each test consumer owns one channel per destination |
| Exchange tap | one exclusive, auto-delete queue per destination (prototest-{guid}), bound with the destination key and #: direct exchanges match the key, topic exchanges match the # catch-all, fanout and headers exchanges match every message; bound during setup for a tapped destination, just in time at the first await otherwise |
| Keyed await | binds the key on the destination tap; a pre-bound tap keeps what arrived before the keyed await |
| Queue destination | consumed as it exists, checked passively and left untouched |
Each tap queue is deleted when the test consumer is disposed, and an exclusive tap never competes with the application's own consumers. A tap binds only to an exchange that exists; the adapter never guesses. The application declares its topology before messaging prepares. The demo declares its event exchanges at application startup and registers AddMessaging last on purpose, so the messaging initializer binds after the in-process application's initializer has created them:
builder
.AddApplication(NorthstarTargets.Api, app => app.AddAspNetCoreServer<Program>())
.AddMessaging(messaging => messaging.UseRabbitMq());
A suite that owns the broker itself declares the destinations it owns instead.
Suite-owned topology
When the run owns the broker and publishes its own events, nothing else declares the destinations and Declare is the documented path:
builder.AddMessaging(messaging => messaging
// The suite owns these destinations: create them before any tap binds or the act publishes.
.Declare("invoice.paid", "invoice.shipped")
.Tap("invoice.paid")
.UseRabbitMq());
Declare creates each destination on the broker during test setup, before any tap is prepared: a fanout, durable, non-auto-delete exchange (the shape the sample application declares for its event exchanges) sent to the broker once per run. A declaration is idempotent: a destination that already exists with that shape is left as it is, repeated Declare calls are deduped, and every later test repeats a no-op. A declaration the broker refuses, when the name already exists with another type or durability, fails setup with the destination named, instead of letting tests publish into a mismatch. Declare the destinations this suite owns, not the application's: declaring one the application declares with other properties is a conflict, not a fix.
A destination whose exchange the application and the suite both leave undeclared fails only the tests that await it, not the whole class: preparing that tap cannot bind, and the first AwaitAsync on the destination throws the named error while the other tests run normally.
Queue destinations
An exchange destination is awaited through a test-owned tap queue bound to it. A destination that names a queue instead - queue:{name}, built with ProtoDestination.Queue(name) - is consumed directly, which is what a dead-letter queue needs: the DLQ is a queue, not an exchange, and the product's own dead-letter bindings are what feed it.
// The product's dead-letter bindings carry the poison into billing.session-ended.dlq;
// the test awaits the queue itself.
var dead = await Proto.Context.Messaging().AwaitAsync(
ProtoDestination.Queue("billing.session-ended.dlq"),
message => message.ReadRequired<SessionEnded>().SessionId == sessionId);
A queue destination is consumed, never declared: the adapter verifies the queue exists - a missing queue fails naming it, like a missing exchange - consumes it on the test's own channel, and leaves the queue as it is when the test ends.
Tap accepts a queue destination too, so the consume starts during setup instead of at the first await; Declare refuses the queue form, because the component that owns the queue creates it. A consumed delivery carries the queue destination on ProtoMessage.Destination and the transport's routing key on ProtoMessage.RoutingKey.
Consuming a shared queue has two consequences. A queue await reads what the queue already holds, unlike an exchange tap that only sees what arrives after it binds - and it removes what it reads, so a second test (or a live consumer) awaiting the same queue no longer sees it.
Await the queue directly when the test owns it (a dead-letter queue no one else reads, a serial suite); prefer the exchange that feeds it when the suite runs in parallel.
A broker whose model has no queues - the in-memory broker, the MassTransit harness - refuses a queue destination with an error naming the transport and the exchange to await instead.
Repeats and consumption
A queue await has no exchange to bind, so a keyed await on a queue destination filters the deliveries the queue hands over by the routing key the transport recorded.
Each await consumes the message it matches. Awaits on one consumer are serialized in call order, and a delivery that matches no awaited predicate is not consumed: it stays available to a later await on the same consumer, so concurrent awaits on one destination neither lose nor steal each other's messages and every matched message is consumed exactly once.
Consumption is tracked per destination, so an await on one destination never hides another destination's first delivery, even though each destination's tap numbers its deliveries from zero.
The in-memory broker behaves the same way: messages live for the run and are ordered, each consumer snapshots the broker position when it is created, so only messages published after its test started can match, and each matched message is consumed once.
A predicate that throws fails only the await that owns it. A timeout is a TimeoutException; awaiting on a tap whose exchange is missing (or a queue that does not exist) is an InvalidOperationException naming the destination, and a Declare the broker refuses fails setup with the destination named; an unreachable broker is an InvalidOperationException naming the sanitized address and ProtoTest:Messaging:RabbitMq:ConnectionString.
Owning a broker
When the run should start RabbitMQ itself, register the container as infrastructure and let the host fill the keys:
builder.AddInfrastructure(
"MessagingBroker",
chain => chain
.UseConfigured()
.UseContainer(RabbitMqBroker.Container()),
RabbitMqOptions.ConnectionStringSetting, // ProtoTest:Messaging:RabbitMq:ConnectionString
"Messaging:RabbitMq:ConnectionString"); // what the application reads
builder.AddMessaging(messaging => messaging
// The suite owns the broker: declare the destinations it publishes to itself.
.Declare("invoice.paid")
.Tap("invoice.paid")
.UseRabbitMq());
The target's UseContainer provider starts the container with the host and fills every key with the started connection string, so the adapter and the application under test reach the same broker.
The provider starts before any test-level skip condition is evaluated, so a machine without a container runtime fails the run at start.
RabbitMqBroker.Container() creates the resource without starting it; Start() starts now or throws with the reason; TryStart(configure) reports the reason in its result instead, for a fixture that decides before registering infrastructure. The default image is rabbitmq:3, configurable through the builder passed to Container.
Registering with AddResource only owns the release: it neither starts the container nor fills settings. With no application initializer to declare the event exchanges, the suite declares its own with Declare (see Suite-owned topology), with no raw broker client in the suite.
MassTransit bridge
An application that composes AddMassTransitTestHarness can be the broker itself: ProtoTest.Messaging.MassTransit publishes and awaits over the application's in-process ITestHarness, so the events the application publishes through its own IPublishEndpoint are the ones a test awaits.
A destination names a message contract type (its full name, short name or urn:message: URN); Declare is a no-op because MassTransit owns message topology; and the Broker capability is declared while the application is hosted in-process (the UseBrokerWhenInProcess seam), so a published application skips instead of failing.
For a published application - or any suite that talks to the broker itself - the package's MassTransitEnvelope builds and reads the MassTransit wire envelope (application/vnd.masstransit+json) through whichever adapter is configured, no harness needed. The MassTransit page has the registration, the ordering rule, the envelope interop and the limits.
Attachments
CaptureAttachments records every published payload and every payload matched by an await as a test attachment:
builder.AddMessaging(messaging => messaging.CaptureAttachments());
A publish attaches message-publish-{destination}-{sequence}-payload after the broker call succeeded; a matched await attaches message-receive-{destination}-{sequence}-payload. {sequence} is the client's per-test capture number, so repeated captures on one destination stay distinct. Payloads are redacted with the shared JSON rules (password, token, secret, … become [REDACTED]) and truncated to MaxDiagnosticBodyLength; the trace's Message section is redacted with the same rules. An explicit content type wins, otherwise a payload starting with { or [ is application/json and everything else is text/plain. Without CaptureAttachments nothing is captured. Capture never fails the operation: a failed capture is a messaging.attachment.failed event and the publish or await still succeeds.
Reference: options and API
| Key | Option | Type | Default |
|---|---|---|---|
ProtoTest:Messaging:DefaultTimeout | MessagingOptions.DefaultTimeout | TimeSpan | 10 seconds |
ProtoTest:Messaging:Destinations | MessagingOptions.Destinations | IList<string> | empty |
ProtoTest:Messaging:DeclaredDestinations | MessagingOptions.DeclaredDestinations | IList<string> | empty |
ProtoTest:Messaging:Attachments:CapturePublishedPayloads | MessagingAttachmentOptions.CapturePublishedPayloads | bool | true |
ProtoTest:Messaging:Attachments:CaptureReceivedPayloads | MessagingAttachmentOptions.CaptureReceivedPayloads | bool | true |
ProtoTest:Messaging:Attachments:RedactSensitiveData | JsonDiagnosticOptions.RedactSensitiveData | bool | true |
ProtoTest:Messaging:Attachments:MaxDiagnosticBodyLength | JsonDiagnosticOptions.MaxDiagnosticBodyLength | int | 65536 (64 KiB) |
ProtoTest:Messaging:Attachments:SensitiveJsonProperties | JsonDiagnosticOptions.SensitiveJsonProperties | List<string> | password, token, access_token, refresh_token, secret, apiKey, api_key, authorization, cookie, connectionString, clientSecret, client_secret, id_token |
ProtoTest:Messaging:RabbitMq:ConnectionString | RabbitMqOptions.ConnectionString | string | amqp://guest:guest@localhost:5672/ |
MessagingAttachmentOptions derives from JsonDiagnosticOptions and binds from ProtoTest:Messaging:Attachments; RabbitMqOptions binds from ProtoTest:Messaging:RabbitMq. Code configuration runs first and the configuration section binds over it. For the connection string the order is: an explicit configuration value wins, then the value a started container filled, then your code callback or the default. Destinations and DeclaredDestinations are lists, so configuration adds its entries after the code-declared ones and the set a test prepares or declares is deduped.
ProtoTest:Messaging:Broker is not a library option. The demo reads it itself (ProtoTest:Messaging:Broker=container) to decide whether to register a container; the messaging packages never look at that key.
public static ProtoMessageClient Messaging(this ProtoExecutionContext context, string? name = null);
Task PublishAsync(string destination, string? payload = null,
IReadOnlyDictionary<string, string?>? headers = null, string? contentType = null,
CancellationToken cancellationToken = default);
Task<ProtoMessage> AwaitAsync(string destination, Func<ProtoMessage, bool> predicate,
TimeSpan? timeout = null, CancellationToken cancellationToken = default);
AwaitAsync returns the first message on the destination that matches the predicate; when no timeout is given it uses MessagingOptions.DefaultTimeout. Messaging() throws when the host was not composed with AddMessaging. ProtoMessage is the broker-agnostic shape every adapter maps onto (Destination, Payload, Headers, ContentType, RoutingKey). The payload reads are typed: message.ReadAsJson<T>() deserializes with ProtoJsonDefaults.Reader and returns default for an empty payload; message.ReadRequired<T>() and message.ReadRequired<T>(jsonPath) throw MessagingAssertionException naming the destination when the payload is empty, JSON null, or the path is missing.
ProtoMessagingBuilder UseBroker(this ProtoMessagingBuilder messaging,
Func<IServiceProvider, IProtoMessageBroker> factory, params string[] addressKeys);
ProtoMessagingBuilder UseBrokerWhenInProcess(this ProtoMessagingBuilder messaging,
Func<IServiceProvider, IProtoMessageBroker> factory, string application, params string[] configuredKeys);
ProtoMessagingBuilder Tap(this ProtoMessagingBuilder messaging, params string[] destinations);
ProtoMessagingBuilder Declare(this ProtoMessagingBuilder messaging, params string[] destinations);
ProtoMessagingBuilder UseRabbitMq(this ProtoMessagingBuilder messaging, Action<RabbitMqOptions>? configure = null);
ProtoMessagingBuilder UseMassTransit<TProgram>(this ProtoMessagingBuilder messaging, string application = "Default");
UseBroker declares the Broker capability only while at least one of its addressKeys can provide an address. UseBrokerWhenInProcess declares it only while the named application is served in-process. A repeated AddMessaging call runs its configure callback again; infrastructure stays idempotent and the first adapter configured wins. When an adapter is configured, AddMessaging also registers the Messaging capability with kind broker.
A routing key names the address a publish takes inside the exchange, and the address a message was delivered under. AwaitAsync(exchange, routingKey, predicate, …) matches only that key; without a routing key the plain overloads publish under the destination itself. The key is bound when the keyed await starts, so a direct exchange only delivers a message published after that; pre-bind the destination with Tap to keep the act-then-await flow reliable there. MassTransit addresses message contract types rather than broker routing, so its adapter rejects a routing key.
message.Should.MatchShape(shape) matches the payload with the same shape matcher as REST, GraphQL and gRPC; MatchShape(shape, exact: true) is the exhaustive form. A mismatch throws MessagingAssertionException whose message starts with the destination, keeping the shared JsonShapeMismatchException as InnerException. An adapter implements two interfaces: the capability owns the broker resource and the test-side API, and the adapter owns the client technology. The full contract lives on Adapter contract.
Destinations
Destinations a suite awaits are declared before the run, so the RabbitMQ adapter can bind each test's own tap during setup. Declare them in code with Tap:
builder.AddMessaging(messaging => messaging
.CaptureAttachments()
// Pre-bind the test's tap before the system under test publishes: the worker can publish
// invoice.issued before a test reaches its first AwaitAsync.
.Tap("invoice.issued", "invoice.paid")
.UseRabbitMq());
Tap takes one or more destinations; repeated calls compose and values already declared are not added twice. Configuration under ProtoTest:Messaging:Destinations still binds over the code values, so an environment can add its own:
{
"ProtoTest": {
"Messaging": {
"Destinations": [ "invoice.paid" ],
"DefaultTimeout": "00:00:15"
}
}
}
From the moment a tap is declared, anything the application publishes is queued for that test, so the usual act-then-await order works. Binding at await time instead would miss everything published in between, which is exactly what happens for an undeclared destination: the consumer binds just in time and can only see later messages. Tap is a reliability declaration: pre-bind every destination the act publishes to. The in-memory broker needs no declaration because it keeps its own history.
A destination the suite itself owns, which nothing else declares, must exist before a tap binds and before the suite publishes to it. Declare it with Declare (see Suite-owned topology):
builder.AddMessaging(messaging => messaging
// The suite publishes its own events on these destinations: create them at prepare.
.Declare("invoice.paid", "invoice.shipped")
.Tap("invoice.paid")
.UseRabbitMq());
In the trace and coverage
messaging.publish · invoice.paid # system InMemory, payload redacted
└─ messaging.published observation (target InMemory)
messaging.await · invoice.paid + routing key "session.ended"
├─ messaging.timeout_ms = 10000
├─ tap bound during setup (Tap), or just in time at the first await
└─ messaging.receive observation on match, TimeoutException naming destination and key on timeout
Every publish and await is recorded:
- Operations
messaging.publishandmessaging.awaitwithmessaging.system(the broker name),messaging.destination,messaging.routing_keywhen the call named one, andmessaging.timeout_mson the await. The payload is recorded as a redactedMessagecode section. - Observations
messaging.publishedfor a successful publish andmessaging.receivefor a matched await; the target is the broker name (InMemoryorRabbitMQ), the identifier is the destination, and the metadata carriesmessaging.system.Should.MatchShapeadds amessaging.contract.shapeobservation carryingMessagingShapeMatchData(Destination, MatchedProperties). - Resources: the run-scoped
messaging:brokerresource with kindbroker, and the per-testmessaging:consumer:Defaultresource with kindconsumer. The container addsbroker:rabbitmq. - The event
messaging.attachment.failedwithattachment.namewhen a capture cannot be registered.
The observations are trace evidence, not a coverage promise: ProtoTest.Messaging ships no collector, so destinations are never aggregated into a report unless you register a collector of your own with the broker's target name (see coverage).
Skip
[RequiresCapability(ProtoCapabilityKinds.Broker)] guards tests that need a real broker, with a reason naming the configuration key:
[RequiresCapability(
ProtoCapabilityKinds.Broker,
Reason = "No broker is configured; set ProtoTest:Messaging:RabbitMq:ConnectionString.")]
AddMessaging registers the Messaging capability only when an adapter is configured. With UseRabbitMq the capability is conditional on ProtoTest:Messaging:RabbitMq:ConnectionString: a run with a configured key or a broker container that declares it keeps the capability, while a run with neither loses it and gated tests skip instead of failing at setup or first publish. A callback that sets RabbitMqOptions.ConnectionString in code provides the address without a key and keeps the capability; an adapter registered through UseBroker without address keys keeps the unconditional declaration.
Limits
-
Repeated
AddMessagingcomposes; the first adapter wins. A later call runs itsconfigureagain and can add an adapter or extend attachment options, but it never replaces the first adapter. -
The in-memory broker is a test double. It registers no
Brokercapability; exchange destinations already exist there, so itsPrepareAsyncandDeclareAsyncdo nothing, and a queue destination is refused because it has no queues. -
No history on RabbitMQ. An exchange tap holds only what arrived after it was declared; a destination declared just in time at the await sees only later messages. A queue destination reads the queue's backlog as well, but removes what it reads.
-
An unmatched delivery stays for a later await. A delivery that matched no awaited predicate is kept for a later await on the same consumer rather than consumed, so concurrent awaits on one destination cannot steal each other's messages; disposing the consumer drops whatever it never matched. Match on the destination and the start of the payload rather than re-awaiting a message another await already consumed.
-
UTF-8 strings only.
ProtoMessage.Payloadis astring?; there is no binary payload API. -
TapandDeclareare run-scoped. They declare destinations for every test in the run; there is no per-test destination declaration onProtoMessageClient. A queue destination is tapped, never declared:Declarecreates exchanges and refuses the queue form. -
Declarecovers the destination, not a binding. On RabbitMQ it creates the exchange the adapter publishes to; the per-test tap queue and its bindings stay the adapter's. An(exchange, routingKey)pair is addressed per publish (PublishAsync(exchange, routingKey, payload)) and per await (AwaitAsync(exchange, routingKey, predicate, …)), and a routing key is bound when the keyed await starts. A routing key on a queue await filters the consumed deliveries instead of binding anything. -
A queue destination is consumed, not tapped. The adapter consumes the named queue and leaves it as it is, so the await reads the queue's backlog and removes what it reads: a second test (or a live consumer) awaiting the same queue competes for its deliveries. Await the queue directly when the test owns it; prefer the exchange that feeds it where the suite runs in parallel or a consumer is live.
-
Configuration adds to
TapandDeclare, it does not replace them. BecauseDestinationsandDeclaredDestinationsare lists, an environment that exportsProtoTest__Messaging__Destinations__0adds a destination; it cannot withdraw a code-declared one. -
Capture is opt-in. Payload attachments exist only after
CaptureAttachments. -
Destinations are evidence, not coverage. No Messaging collector ships (a decision, not a gap); the
messaging.published,messaging.receiveandmessaging.contract.shapeobservations reach a report only through a collector a suite registers. -
One run connection, serialized consumers. RabbitMQ uses a single connection and one publish channel; every consumer owns a channel per destination and awaits on one consumer serialize in call order. An exchange tap's exchange must exist, because the application declares its topology or the suite declares its own with
Declare, and there is no retry or backoff.
Links
- The sample's event journey:
samples/Northstar.ProtoTest/BrokerJourney.csand the host wiring insamples/Northstar.ProtoTest/Setup.cs. - Recipe: API publishes an event.
- Related: Coverage and observations, ProtoTrace, Infrastructure, Skip conditions.