Sanyappc.Extensions.RabbitMq
0.1.1
dotnet add package Sanyappc.Extensions.RabbitMq --version 0.1.1
NuGet\Install-Package Sanyappc.Extensions.RabbitMq -Version 0.1.1
<PackageReference Include="Sanyappc.Extensions.RabbitMq" Version="0.1.1" />
<PackageVersion Include="Sanyappc.Extensions.RabbitMq" Version="0.1.1" />
<PackageReference Include="Sanyappc.Extensions.RabbitMq" />
paket add Sanyappc.Extensions.RabbitMq --version 0.1.1
#r "nuget: Sanyappc.Extensions.RabbitMq, 0.1.1"
#:package Sanyappc.Extensions.RabbitMq@0.1.1
#addin nuget:?package=Sanyappc.Extensions.RabbitMq&version=0.1.1
#tool nuget:?package=Sanyappc.Extensions.RabbitMq&version=0.1.1
Sanyappc.Extensions.RabbitMq
A .NET library for publishing and consuming RabbitMQ messages. Supports typed JSON messaging, manual acknowledgement, request/reply via Direct Reply-to, configurable reply timeout, automatic connection recovery, multiple broker connections, well-typed exceptions, health checks, and built-in OpenTelemetry tracing, metrics and structured logs following messaging semantic conventions. Ships a Roslyn analyzer that flags polymorphic message types passed or read as a derived type at compile time.
Installation
dotnet add package Sanyappc.Extensions.RabbitMq
Configuration
The library binds options from the RabbitMq configuration section.
appsettings.json
{
"RabbitMq": {
"Hostname": "localhost",
"Port": -1,
"Username": "guest",
"Password": "guest",
"ReplyTimeout": "00:00:05",
"RecoveryInterval": "00:00:05",
"RecoveryTimeout": "00:01:00"
}
}
Port: -1uses the RabbitMQ default port (5672). Durations areTimeSpanvalues, writtenhh:mm:ssin configuration.ReplyTimeoutis how longRequestAsyncwaits for a reply. Default is00:00:05;Timeout.InfiniteTimeSpan(-00:00:00.001in configuration) waits indefinitely.RecoveryIntervalis the wait between reconnect attempts after a lost connection. Default is00:00:05.RecoveryTimeoutis how long a consumer waits for its connection to come back before it gives up. Default is00:01:00; it must be greater thanRecoveryInterval. See Connection recovery.
Environment variables
Use __ as the section separator:
RabbitMq__Hostname=localhost
RabbitMq__Port=5672
RabbitMq__Username=guest
RabbitMq__Password=guest
RabbitMq__ReplyTimeout=00:00:30
RabbitMq__RecoveryInterval=00:00:05
RabbitMq__RecoveryTimeout=00:01:00
Programmatic (code)
services.AddRabbitMqService(options =>
{
options.Hostname = "localhost";
options.Username = "guest";
options.Password = "guest";
});
The delegate takes precedence over configuration. It is useful for overriding specific values regardless of the config file.
Registration
builder.Services.AddRabbitMqService();
Multiple brokers
To connect to more than one RabbitMQ broker, use the named overload. Each name registers an independent set of keyed services.
builder.Services.AddRabbitMqService("broker1", o => o.Hostname = "rabbit1");
builder.Services.AddRabbitMqService("broker2", o => o.Hostname = "rabbit2");
Config binding uses RabbitMq:{name} as the section:
{
"RabbitMq": {
"broker1": { "Hostname": "rabbit1", "Username": "guest", "Password": "guest" },
"broker2": { "Hostname": "rabbit2", "Username": "guest", "Password": "guest" }
}
}
Inject a named publisher using [FromKeyedServices]:
public class MyService([FromKeyedServices("broker2")] IRabbitMqPublishService publisher)
{
}
Register a named consumer:
builder.Services.AddRabbitMqConsumer<PaymentProcessor>("broker2", "payments");
Publishing
Inject IRabbitMqPublishService and call PublishAsync. Messages are serialized as JSON.
public class OrderService(IRabbitMqPublishService publisher)
{
public Task SendOrderAsync(Order order, CancellationToken ct) =>
publisher.PublishAsync("orders", order, cancellationToken: ct);
}
Raw bytes are also supported:
await publisher.PublishAsync("orders", bytes, ct);
Consuming
Implement IRabbitMqMessageProcessingService for your message handler. The queue is consumed with prefetch=1 and manual acknowledgement — you must call AckAsync or RejectAsync on every message.
public class OrderProcessor : IRabbitMqMessageProcessingService
{
public async Task ProcessMessageAsync(RabbitMqMessage message, CancellationToken ct)
{
Order order = message.GetBody<Order>();
// process...
await message.AckAsync(ct);
// or: await message.RejectAsync(requeue: false, ct);
}
}
Register the consumer in DI. This starts a hosted service that runs for the lifetime of the application:
builder.Services.AddRabbitMqConsumer<OrderProcessor>("orders");
AddRabbitMqConsumer and AddRabbitMqRpcConsumer both call AddRabbitMqService internally, so the explicit call is optional when using only consumers.
Multiple queues:
builder.Services.AddRabbitMqConsumer<OrderProcessor>("orders");
builder.Services.AddRabbitMqConsumer<PaymentProcessor>("payments");
Request / Reply
For synchronous RPC over RabbitMQ using Direct Reply-to:
Caller — with response body:
InvoiceResponse invoice = await publisher.RequestAsync<OrderRequest, InvoiceResponse>(
"invoices", new OrderRequest { OrderId = 42 }, cancellationToken: ct);
Caller — no response body (acknowledgement only):
await publisher.RequestAsync<DispatchCommand>("dispatch", new DispatchCommand { OrderId = 42 }, cancellationToken: ct);
Handler:
Implement IRabbitMqRpcMessageProcessingService. ReplyAsync sends the reply and acknowledges the message — no separate AckAsync call is needed.
public class InvoiceProcessor : IRabbitMqRpcMessageProcessingService
{
public async Task ProcessMessageAsync(RabbitMqRpcMessage message, CancellationToken ct)
{
OrderRequest request = message.GetBody<OrderRequest>();
InvoiceResponse response = new() { /* ... */ };
await message.ReplyAsync(response, cancellationToken: ct);
}
}
To reply with no body (acknowledgement only), call the parameterless overload:
await message.ReplyAsync(ct);
Register with AddRabbitMqRpcConsumer:
builder.Services.AddRabbitMqRpcConsumer<InvoiceProcessor>("invoices");
Returning an error to the caller:
Call ReplyErrorAsync instead of ReplyAsync. It sends the error message back and acknowledges the delivery. The caller receives a RabbitMqRequestRejectedException.
public async Task ProcessMessageAsync(RabbitMqRpcMessage message, CancellationToken ct)
{
OrderRequest request = message.GetBody<OrderRequest>();
if (!IsValid(request))
{
await message.ReplyErrorAsync("Invalid request", ct);
return;
}
await message.ReplyAsync(new InvoiceResponse { /* ... */ }, cancellationToken: ct);
}
Fire-and-forget messages on an RPC queue:
If a message arrives without a ReplyTo header (sent fire-and-forget to the same queue), ReplyAsync skips the publish and only acknowledges. The handler code does not need to change.
Polymorphic messages
System.Text.Json writes a type discriminator only when the value is serialized through the polymorphic base type, and checks one only when deserialized through it. A [JsonDerivedType] attribute is all it takes to make a base polymorphic; [JsonPolymorphic] only customises the discriminator.
[JsonPolymorphic(TypeDiscriminatorPropertyName = "kind")]
[JsonDerivedType(typeof(Block), "block")]
[JsonDerivedType(typeof(Delete), "delete")]
public abstract record Deactivate { /* ... */ }
// Discriminator dropped — derived static type, also inside a List<> or an array
Deactivate.Delete msg = new() { /* ... */ };
await publisher.PublishAsync("q", msg, ct);
// Discriminator written — variable typed as the polymorphic base
Deactivate msg = new Deactivate.Delete { /* ... */ };
await publisher.PublishAsync("q", msg, ct);
// Discriminator ignored — a "block" payload is accepted as a Delete with default members
Deactivate.Delete read = message.GetBody<Deactivate.Delete>();
// Discriminator checked — the result is the type the payload names
Deactivate read = message.GetBody<Deactivate>();
The package ships a Roslyn analyzer that flags both unsafe forms at compile time: SANYRMQ001 for a derived type passed to PublishAsync, RequestAsync, ReplyAsync or SerializeBody, and SANYRMQ002 for a derived type read with GetBody, DeserializeBody or as the reply type of RequestAsync. Suppress them where you intentionally want concrete-type serialization.
Health checks
Register a health check that verifies broker connectivity by opening a channel. Requires AddRabbitMqService() to have been called first (or AddRabbitMqConsumer / AddRabbitMqRpcConsumer, which call it internally):
builder.Services.AddRabbitMqService();
builder.Services.AddHealthChecks()
.AddRabbitMq();
For a named connection, pass the connection name. The health check name defaults to rabbitmq:{connectionName}:
builder.Services.AddHealthChecks()
.AddRabbitMq("broker1")
.AddRabbitMq("broker2");
Override the health check name with the second parameter:
builder.Services.AddHealthChecks()
.AddRabbitMq("broker1", name: "primary-broker");
Connection recovery
A lost broker connection does not stop a consumer. The client reconnects every RecoveryInterval, redeclares the queues, restores the channels and their consumers, and ConsumeAsync / ConsumeRpcAsync carry on in the same process. A message that was delivered but not yet acknowledged when the connection dropped is redelivered, so processing must tolerate a repeat. Publishing during the outage throws RabbitMqUnavailableException; once the connection is back, the next publish uses it.
A consumer gives up and throws RabbitMqUnavailableException when:
- its connection is not back within
RecoveryTimeout. The hosted consumers registered byAddRabbitMqConsumer/AddRabbitMqRpcConsumerthen stop the host with exit code1, so an orchestrator restarts the process; - the broker closes its channel while the connection stays open, for example after an acknowledgement with an unknown delivery tag. The client never reopens such a channel, so the consumer fails at once instead of waiting.
The channel factory logs every loss and recovery: the loss is a warning (event 3), each failed attempt and the recovery are information (events 4 and 5), and every consumer notes its channel going down and coming back (events 42 and 47). See Log events.
Each recovery is also measured by sanyappc.rabbitmq.connection.recovery.duration (see Metrics), so recoveries stay visible on a dashboard although nothing restarts.
Error handling
All library errors derive from RabbitMqException, so you can catch the base type or a specific subtype:
| Exception | When thrown |
|---|---|
RabbitMqUnavailableException |
Broker is unreachable, a consumer's connection is not recovered within RecoveryTimeout, or the broker closes a consumer's channel while the connection stays open |
RabbitMqTimeoutException |
RequestAsync did not receive a reply within ReplyTimeout |
RabbitMqRequestRejectedException |
RequestAsync received an error reply from the handler via ReplyErrorAsync |
try
{
await publisher.PublishAsync("orders", order, cancellationToken: ct);
}
catch (RabbitMqUnavailableException ex)
{
// broker down — retry, circuit-break, or return 503
}
try
{
InvoiceResponse invoice = await publisher.RequestAsync<OrderRequest, InvoiceResponse>(
"invoices", request, cancellationToken: ct);
}
catch (RabbitMqRequestRejectedException ex)
{
// handler called ReplyErrorAsync — ex.Message contains the error
}
catch (RabbitMqTimeoutException)
{
// no reply within ReplyTimeout — return 504
}
catch (RabbitMqUnavailableException)
{
// broker down — return 503
}
OpenTelemetry
Tracing
The library creates spans for publish, request, and receive operations using the activity source name exposed by RabbitMqTelemetry.ActivitySourceName.
builder.Services.AddOpenTelemetry()
.WithTracing(tracing => tracing
.AddSource(RabbitMqTelemetry.ActivitySourceName)
.AddOtlpExporter());
Spans follow OpenTelemetry messaging semantic conventions:
| Operation | Span name | Kind | messaging.operation.type |
|---|---|---|---|
PublishAsync |
send {queue} |
Producer | send |
RequestAsync |
send {queue} |
Client | send |
ConsumeAsync |
process {queue} |
Consumer | process |
ConsumeRpcAsync |
process {queue} |
Consumer | process |
Each span carries the messaging attributes (messaging.system, messaging.destination.name, messaging.operation.name, messaging.operation.type, messaging.rabbitmq.destination.routing_key, server.address, server.port, messaging.message.body.size) plus messaging.rabbitmq.message.delivery_tag, messaging.message.id, and messaging.message.conversation_id on the consumer side when available. All sampling-relevant attributes are set at activity creation time. On failure, the span sets error.type, records the exception via Activity.AddException, and sets the status to Error.
Metrics
The library records messaging metrics using the meter name exposed by RabbitMqTelemetry.MeterName.
builder.Services.AddOpenTelemetry()
.WithMetrics(metrics => metrics
.AddMeter(RabbitMqTelemetry.MeterName)
.AddOtlpExporter());
| Instrument | Type | Unit | When recorded |
|---|---|---|---|
messaging.client.sent.messages |
Counter | {message} |
Once per attempted PublishAsync or RequestAsync (success or failure) |
messaging.client.consumed.messages |
Counter | {message} |
On each delivered message |
messaging.client.operation.duration |
Histogram | s |
Per PublishAsync / RequestAsync call (success or failure) |
messaging.process.duration |
Histogram | s |
Per invocation of ProcessMessageAsync |
sanyappc.rabbitmq.connection.recovery.duration |
Histogram | s |
Once per recovered connection, from the loss until the connection, its channels and their consumers are back |
Every messaging instrument includes messaging.system, messaging.destination.name, messaging.operation.name, messaging.operation.type, messaging.rabbitmq.destination.routing_key, server.address, and server.port. Failed operations additionally set error.type to one of timeout, request_rejected, broker_unavailable, or the fully qualified exception type for unexpected errors. The recovery histogram carries only messaging.system, server.address, and server.port, since one connection serves every queue.
Log events
Every message the library logs carries a stable event id, so a sink can filter or alert on an event without matching its text. Ids 1–19 belong to the connection, 20–39 to publishing, 40–59 to consuming. The placeholders are the log record's attribute names: the same semantic convention names the spans and metrics carry (messaging.destination.name, server.address, server.port), and sanyappc.rabbitmq.* for what has no convention.
| Id | Level | Message |
|---|---|---|
| 1 | Information | RabbitMQ connection established to {server.address}:{server.port} |
| 2 | Information | RabbitMQ connection disposed |
| 3 | Warning | RabbitMQ connection to {server.address}:{server.port} lost: {sanyappc.rabbitmq.shutdown.reason}. Reconnecting every {sanyappc.rabbitmq.connection.recovery.interval} |
| 4 | Information | RabbitMQ connection to {server.address}:{server.port} could not be recovered yet: {sanyappc.rabbitmq.connection.recovery.error}, once per failed attempt |
| 5 | Information | RabbitMQ connection to {server.address}:{server.port} recovered after {sanyappc.rabbitmq.connection.recovery.duration} |
| 20 | Debug | Publishing message to queue {messaging.destination.name} |
| 21 | Error | RabbitMQ broker unavailable while publishing to queue {messaging.destination.name}, with the exception |
| 22 | Debug | Sending request to queue {messaging.destination.name}, awaiting reply |
| 23 | Warning | RabbitMQ request to queue {messaging.destination.name} timed out after {sanyappc.rabbitmq.reply.timeout} |
| 24 | Error | RabbitMQ broker unavailable during request to queue {messaging.destination.name}, with the exception |
| 40 | Debug | Received message from queue {messaging.destination.name} |
| 41 | Error | Error processing message from queue {messaging.destination.name}, with the exception |
| 42 | Information | Channel shut down unexpectedly: {sanyappc.rabbitmq.shutdown.reason} |
| 43 | Error | RabbitMQ broker unavailable while consuming from queue {messaging.destination.name}, with the exception |
| 44 | Debug | Received RPC message from queue {messaging.destination.name} |
| 45 | Error | Error processing RPC message from queue {messaging.destination.name}, with the exception |
| 46 | Error | RabbitMQ broker unavailable while consuming RPC from queue {messaging.destination.name}, with the exception |
| 47 | Debug | Resumed consuming from queue {messaging.destination.name} after the connection recovered |
Log correlation
Each consumed message is processed inside a logger scope, so every line logged meanwhile, by the library or by your processor, carries the message's own attributes in structured sinks (Seq, Loki, Application Insights, etc.). On the simple console the scope renders as queue orders, message 7f3a, delivery tag 42.
| Key | Value |
|---|---|
messaging.destination.name |
Queue name |
messaging.message.id |
AMQP message ID (null if not set by publisher) |
messaging.rabbitmq.message.delivery_tag |
Per-channel delivery sequence number |
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net10.0 is compatible. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net10.0
- Microsoft.Extensions.Diagnostics.HealthChecks (>= 10.0.12)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.12)
- Microsoft.Extensions.Logging (>= 10.0.12)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.12)
- Microsoft.Extensions.Options.DataAnnotations (>= 10.0.12)
- Microsoft.Extensions.Telemetry.Abstractions (>= 10.10.0)
- RabbitMQ.Client (>= 7.2.2)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 0.1.1 | 92 | 9/22/2026 |
| 0.1.0 | 93 | 9/22/2026 |
| 0.0.11 | 94 | 9/14/2026 |
| 0.0.10 | 134 | 5/9/2026 |
| 0.0.9 | 133 | 3/25/2026 |
| 0.0.8 | 119 | 3/24/2026 |
| 0.0.7 | 131 | 3/19/2026 |
| 0.0.6 | 121 | 3/18/2026 |
| 0.0.5 | 121 | 3/18/2026 |
| 0.0.4 | 128 | 3/17/2026 |
| 0.0.3 | 158 | 3/16/2026 |
| 0.0.2 | 177 | 3/14/2026 |
| 0.0.1 | 155 | 3/14/2026 |
| 0.0.1-alpha.9 | 91 | 3/12/2026 |
| 0.0.1-alpha.8 | 248 | 12/13/2024 |
| 0.0.1-alpha.7 | 119 | 12/13/2024 |
| 0.0.1-alpha.6 | 118 | 12/12/2024 |
| 0.0.1-alpha.5 | 113 | 12/10/2024 |
| 0.0.1-alpha.4 | 115 | 12/6/2024 |
| 0.0.1-alpha.3 | 106 | 12/2/2024 |