LowCodeHub.RabbitMQ
0.0.1
See the version list below for details.
dotnet add package LowCodeHub.RabbitMQ --version 0.0.1
NuGet\Install-Package LowCodeHub.RabbitMQ -Version 0.0.1
<PackageReference Include="LowCodeHub.RabbitMQ" Version="0.0.1" />
<PackageVersion Include="LowCodeHub.RabbitMQ" Version="0.0.1" />
<PackageReference Include="LowCodeHub.RabbitMQ" />
paket add LowCodeHub.RabbitMQ --version 0.0.1
#r "nuget: LowCodeHub.RabbitMQ, 0.0.1"
#:package LowCodeHub.RabbitMQ@0.0.1
#addin nuget:?package=LowCodeHub.RabbitMQ&version=0.0.1
#tool nuget:?package=LowCodeHub.RabbitMQ&version=0.0.1
LowCodeHub.RabbitMQ
A topology-first RabbitMQ library for .NET 10. Declare queues, exchanges, bindings, and consumers in a single fluent DSL — the library compiles it to broker topology, wires up DI, and manages the consumer lifecycle automatically.
Built on RabbitMQ.Client 7.x with full async channel support.
Table of Contents
- Installation
- Quick Start
- Connection Configuration
- Topology DSL
- Queue Options
- Dead Letter Queues
- Consumers
- Publishing
- RPC (Request/Reply)
- Health Checks
- Observability
- Clustering
- SSL/TLS
- Conflict Handling
- Custom Serializer
- Architecture Overview
- Production Checklist
Installation
dotnet add package LowCodeHub.RabbitMQ
Quick Start
1. Add connection config to appsettings.json:
{
"RabbitMQ": {
"HostName": "localhost",
"Port": 5672,
"UserName": "myapp",
"Password": "secret",
"VirtualHost": "/"
}
}
2. Register services and declare topology:
builder.Services.AddRabbitMq(builder.Configuration, topology =>
{
topology.Settings(s => s
.Prefetch(20)
.QueueDefaults(q => q.Quorum().DeadLetterDefaultExchange()));
topology.Default(x =>
{
x.Queue("orders.created");
x.Queue("orders.failed");
x.Queue("payments.processed");
});
});
3. Create a consumer:
[Queue("orders.created")]
public sealed class OrderCreatedConsumer : MessageConsumer<OrderCreatedEvent>
{
public override async Task ConsumeAsync(
ConsumeContext<OrderCreatedEvent> context,
CancellationToken cancellationToken = default)
{
var order = context.Message;
// Process the order...
}
}
4. Publish a message:
public class OrderService(IMessagePublisher publisher)
{
public Task PlaceOrderAsync(OrderCreatedEvent evt, CancellationToken ct)
=> publisher.PublishWithDefaultExchangeAsync("orders.created", evt, cancellationToken: ct);
}
That's it. The library handles connection management, topology declaration, consumer wiring, retries, DLQ routing, health checks, and OpenTelemetry metrics.
Connection Configuration
Configure via appsettings.json under the RabbitMQ section:
{
"RabbitMQ": {
"HostName": "rabbitmq.internal",
"Port": 5672,
"VirtualHost": "/",
"UserName": "myapp",
"Password": "secret",
"ClientProvidedName": "MyService",
"AutomaticRecoveryEnabled": true,
"TopologyRecoveryEnabled": true,
"NetworkRecoveryIntervalSeconds": 5,
"RequestedHeartbeat": 60,
"ConnectionTimeoutSeconds": 30,
"DefaultPrefetchCount": 10
}
}
| Property | Default | Description |
|---|---|---|
HostName |
localhost |
Broker hostname or IP |
Port |
5672 |
AMQP port |
VirtualHost |
/ |
RabbitMQ virtual host |
UserName |
guest |
Authentication username |
Password |
guest |
Authentication password |
ClientProvidedName |
auto | Connection name shown in RabbitMQ management UI |
AutomaticRecoveryEnabled |
true |
Auto-reconnect on connection loss |
TopologyRecoveryEnabled |
true |
Re-declare topology after reconnect (required for cluster failover) |
NetworkRecoveryIntervalSeconds |
5 |
Delay between recovery attempts |
RequestedHeartbeat |
60 |
Heartbeat interval in seconds |
ConnectionTimeoutSeconds |
30 |
Connection timeout |
DefaultPrefetchCount |
10 |
Default QoS prefetch per consumer |
Or configure programmatically:
builder.Services.AddRabbitMq(
options =>
{
options.HostName = "rabbitmq.internal";
options.UserName = "myapp";
options.Password = "secret";
},
topology => { /* ... */ });
Topology DSL
Settings
Every topology must start with Settings(...):
topology.Settings(s => s
.Prefetch(50) // Override default prefetch count
.DeleteAndRecreateOnConflict() // Recreate conflicting entities on startup
.QueueDefaults(q => q.Quorum())); // Apply defaults to all queues
Queue Defaults
Set once, applied to every queue. Individual queues can override any setting:
topology.Settings(s => s
.QueueDefaults(q => q
.Quorum() // Default queue type
.DeadLetterDefaultExchange() // Auto-create DLQ for every queue
.Ttl(86_400_000) // 24h message TTL
.MaxLength(100_000) // Max 100k messages
.LazyMode())); // Disk-first storage
Override per queue:
topology.Default(x =>
{
x.Queue("orders.created"); // Inherits all defaults
x.Queue("temp.events", q => q.Stream()); // Override: stream instead of quorum
x.Queue("fast.queue", q => q.NoDeadLetter()); // Override: no DLQ
x.Queue("priority.queue", q => q.MaxPriority(10)); // Add priority support
});
Default Exchange Queues
For point-to-point messaging using the AMQP default exchange (""):
topology.Default(x =>
{
x.Queue("orders.created");
x.Queue("payments.processed");
x.Queue("notifications.email");
});
Direct Exchange
Route messages by exact routing key match:
topology.Direct("payments.direct", direct =>
{
direct.Queue("payments.authorized", "payment.authorized");
direct.Queue("payments.failed", "payment.failed");
});
Topic Exchange
Route messages by wildcard patterns (* = one word, # = zero or more):
topology.Topic("orders.events", topic =>
{
topic.Queue("orders.created", "order.created.*");
topic.Queue("orders.all", "order.#");
});
Fanout Exchange
Broadcast to all bound queues (routing key ignored):
topology.Fanout("notifications.fanout", fanout =>
{
fanout.Queue("notifications.email");
fanout.Queue("notifications.sms");
fanout.Queue("notifications.push");
});
Headers Exchange
Route by header matching:
topology.Headers("risk.headers", headers =>
{
headers.Queue("risk.high",
b => b.MatchAll().Header("region", "eu").Header("priority", "high"));
headers.Queue("risk.any",
b => b.MatchAny().Header("region", "eu").Header("priority", "high"));
});
Queue Options
Available on both QueueDefaults(...) and per-queue Queue("name", q => ...):
| Method | Description |
|---|---|
.Quorum() |
Quorum queue (replicated, durable — recommended) |
.Stream() |
Stream queue (append-only log) |
.Type(RabbitMqQueueType.Classic) |
Explicit queue type |
.Durable(bool) |
Survive broker restart (default: true) |
.Exclusive(bool) |
Exclusive to connection |
.AutoDelete(bool) |
Delete when last consumer disconnects |
.Ttl(int ms) |
Per-message TTL |
.MaxLength(int) |
Max message count |
.MaxLengthBytes(long) |
Max total size in bytes |
.MaxPriority(int) |
Enable priority queue (0–255) |
.LazyMode() |
Store messages to disk early |
.DeadLetterDefaultExchange() |
Route dead letters via "" exchange to {name}.dlq |
.DeadLetterNamedExchange(...) |
Route dead letters via a named exchange |
.NoDeadLetter() |
Disable dead lettering |
Dead Letter Queues
Default Exchange DLQ (Recommended)
x.Queue("orders.created", q => q.DeadLetterDefaultExchange());
Creates:
- Main queue:
orders.createdwithx-dead-letter-exchange: ""andx-dead-letter-routing-key: orders.created.dlq - DLQ:
orders.created.dlq(durable, classic)
Named Exchange DLQ
x.Queue("orders.created", q => q
.DeadLetterNamedExchange("dlx.orders", routingKey: "failed", dlqName: "orders.dead"));
Creates:
- Main queue with
x-dead-letter-exchange: dlx.orders - Exchange:
dlx.orders(direct, durable) - DLQ:
orders.deadbound todlx.orderswith routing keyfailed
Consumers
Message Consumer
[Queue("orders.created", PrefetchCount = 50, MaxRetryAttempts = 5, RequeueOnFailure = true)]
public sealed class OrderCreatedConsumer : MessageConsumer<OrderCreatedEvent>
{
private readonly IOrderService _orderService;
public OrderCreatedConsumer(IOrderService orderService)
{
_orderService = orderService;
}
public override async Task ConsumeAsync(
ConsumeContext<OrderCreatedEvent> context,
CancellationToken cancellationToken = default)
{
await _orderService.ProcessAsync(context.Message, cancellationToken);
}
}
[Queue] attribute options:
| Property | Default | Description |
|---|---|---|
Name |
(required) | Queue name to consume from |
PrefetchCount |
0 (use global) |
Override QoS prefetch for this consumer |
AutoAck |
false |
Auto-acknowledge on delivery (use false for reliability) |
MaxRetryAttempts |
3 |
Max retry count before rejecting to DLQ |
RequeueOnFailure |
false |
Republish failed messages for retry |
Consumer Lifecycle Hooks
[Queue("orders.created")]
public sealed class OrderCreatedConsumer : MessageConsumer<OrderCreatedEvent>
{
public override async Task<bool> OnBeforeConsumeAsync(
ConsumeContext<OrderCreatedEvent> context,
CancellationToken cancellationToken = default)
{
// Return false to skip processing and acknowledge
if (context.Message.Amount <= 0) return false;
return true;
}
public override async Task ConsumeAsync(
ConsumeContext<OrderCreatedEvent> context,
CancellationToken cancellationToken = default)
{
// Main processing logic
}
public override async Task OnAfterConsumeAsync(
ConsumeContext<OrderCreatedEvent> context,
CancellationToken cancellationToken = default)
{
// Post-processing (runs after successful consume)
}
}
Error Handling
Consumer Pipeline Lifecycle
Every message flows through these stages in order:
OnBeforeConsumeAsync— optional pre-processing gate. Returnfalseto skip the message and auto-ack it.ConsumeAsync— your business logic (required override).OnAfterConsumeAsync— optional post-processing hook after successful consumption.- If
ConsumeAsync(or any stage) throws,OnErrorAsyncis invoked instead.
ErrorHandlingResult Options
When ConsumeAsync throws, the framework catches the exception and calls OnErrorAsync. You return one of three values:
| Result | Behavior |
|---|---|
Acknowledge |
Ack the message — removed from the queue permanently. |
Reject |
Reject without requeue. If a DLQ is configured, the message goes there. (default) |
Requeue |
Republish the message for retry (up to MaxRetryAttempts). |
The default OnErrorAsync implementation returns ErrorHandlingResult.Reject.
[Queue] Attribute Settings
| Property | Default | Description |
|---|---|---|
MaxRetryAttempts |
3 |
Maximum times the message will be retried before final rejection. |
RequeueOnFailure |
false |
Declarative flag used alongside OnErrorAsync. |
AutoAck |
false |
If true, messages are auto-acked on delivery and all ack/reject logic is skipped. |
PrefetchCount |
0 |
Per-consumer prefetch override. 0 uses the global setting. |
How Retry Works Internally
The library does not use BasicNack(requeue: true) (which causes infinite redelivery loops). Instead it uses a republish-with-header strategy:
- Reads the
x-retry-countheader from the message (starts at0). - If
retryCount < MaxRetryAttempts, it republishes the message to the same exchange/routing-key with an incrementedx-retry-countheader, then acks the original delivery. - If
retryCount >= MaxRetryAttempts, it rejects the message (BasicReject(requeue: false)), sending it to the DLQ if one is configured.
This ensures retry is finite and safe — no infinite loops.
Accessing Retry Count
The current retry attempt is exposed on ConsumeContext<T>.RetryCount:
public override Task ConsumeAsync(ConsumeContext<OrderCreated> context, CancellationToken ct)
{
if (context.RetryCount > 0)
Console.WriteLine($"Retry attempt: {context.RetryCount}");
// business logic
return Task.CompletedTask;
}
Full Example: Retry with Conditional Error Handling + DLQ
1. Configure the queue with a DLQ:
mq.Default(x => x.Queue("payments.process", q => q
.QueueType("quorum")
.DeadLetterDefaultExchange())); // creates payments.process.dlq
2. Implement the consumer:
using LowCodeHub.RabbitMQ.Abstractions;
using LowCodeHub.RabbitMQ.Attributes;
[Queue("payments.process", MaxRetryAttempts = 5)]
public sealed class PaymentsConsumer : MessageConsumer<PaymentCommand>
{
public override async Task ConsumeAsync(
ConsumeContext<PaymentCommand> context,
CancellationToken cancellationToken = default)
{
await ProcessPaymentAsync(context.Message, cancellationToken);
}
public override Task<ErrorHandlingResult> OnErrorAsync(
ConsumeContext<PaymentCommand> context,
Exception exception,
CancellationToken cancellationToken = default)
{
// Transient errors → retry (up to MaxRetryAttempts)
if (exception is HttpRequestException or TimeoutException)
return Task.FromResult(ErrorHandlingResult.Requeue);
// Permanent errors → reject immediately (goes to DLQ)
return Task.FromResult(ErrorHandlingResult.Reject);
}
}
What happens at runtime:
ConsumeAsyncthrows anHttpRequestException→OnErrorAsyncreturnsRequeue.- The message is republished with
x-retry-count = 1and re-delivered. - Steps 1–2 repeat up to 5 times.
- After 5 retries the message is rejected and routed to
payments.process.dlqfor inspection. - If
ConsumeAsyncthrows a non-transient exception (e.g.ArgumentException),OnErrorAsyncreturnsRejectand the message goes to the DLQ immediately without any retry.
Retry Behavior
When ErrorHandlingResult.Requeue is returned, the message is republished to the same queue with an x-retry-count header. Once MaxRetryAttempts is reached, the message is rejected (sent to DLQ).
ConsumeContext Properties
| Property | Type | Description |
|---|---|---|
Message |
TMessage |
Deserialized message body |
Body |
ReadOnlyMemory<byte> |
Raw message bytes |
Properties |
IReadOnlyBasicProperties |
AMQP message properties |
DeliveryTag |
ulong |
Delivery tag for ack/nack |
Exchange |
string |
Source exchange |
RoutingKey |
string |
Routing key |
Redelivered |
bool |
Whether this is a redelivery |
ConsumerTag |
string |
Consumer tag |
MessageId |
string? |
Message ID |
CorrelationId |
string? |
Correlation ID |
RetryCount |
int |
Current retry attempt (0-based) |
Timestamp |
DateTimeOffset? |
Publish timestamp |
Publishing
Publish to Default Exchange
Point-to-point delivery to a specific queue:
await publisher.PublishWithDefaultExchangeAsync("orders.created", message, cancellationToken: ct);
Publish to Named Exchange
Route through an exchange with a routing key:
await publisher.PublishAsync("orders.events", "order.created.eu", message, cancellationToken: ct);
Publisher Confirms
Wait for broker acknowledgment (reliable publishing):
bool confirmed = await publisher.PublishWithConfirmAsync(
"orders.events", "order.created.eu", message, cancellationToken: ct);
if (!confirmed)
{
// Handle publish failure
}
Batch Publishing
Publish multiple messages on a single channel:
var messages = new[] { order1, order2, order3 };
await publisher.PublishBatchAsync("orders.events", "order.created", messages, cancellationToken: ct);
Publish Options
await publisher.PublishAsync("orders.events", "order.created", message, new PublishOptions
{
MessageId = "custom-id",
CorrelationId = correlationId,
Persistent = true, // Survive broker restart (default: true)
Priority = 5, // 0-9
Expiration = "60000", // TTL in ms
Mandatory = true, // Return unroutable messages
Headers = new() { ["tenant"] = "acme" }
}, ct);
Factory methods:
PublishOptions.WithCorrelationId("abc-123");
PublishOptions.WithExpiration(60_000);
PublishOptions.ForRpc("reply.queue", correlationId: "abc");
RPC (Request/Reply)
RPC Server
[Queue("pricing.requests")]
public sealed class PricingRpcConsumer : RpcMessageConsumer<PriceRequest, PriceResponse>
{
public override async Task<PriceResponse> HandleAsync(
RpcContext<PriceRequest> context,
CancellationToken cancellationToken = default)
{
var price = await CalculatePrice(context.Message);
return new PriceResponse { Total = price };
}
}
RPC Client
Register the RPC client:
builder.Services.AddRabbitMqRpc();
Make a call:
public class PricingService(IMessageRpcClient rpcClient)
{
public async Task<PriceResponse> GetPriceAsync(PriceRequest request, CancellationToken ct)
{
return await rpcClient.CallAsync<PriceRequest, PriceResponse>(
exchange: "",
routingKey: "pricing.requests",
request: request,
timeoutMs: 10_000,
cancellationToken: ct);
}
}
Health Checks
Automatically registered with tag rabbitmq:
app.MapHealthChecks("/health", new HealthCheckOptions
{
Predicate = check => check.Tags.Contains("rabbitmq")
});
Response includes connected endpoint and cluster name.
Observability
Built-in OpenTelemetry support via System.Diagnostics.ActivitySource and System.Diagnostics.Metrics.
Activity source: LowCodeHub.RabbitMQ
| Activity | Kind | Description |
|---|---|---|
rabbitmq.consume |
Consumer | Per-message processing span |
Meters: LowCodeHub.RabbitMQ
| Metric | Type | Description |
|---|---|---|
rabbitmq.consumer.messages.processed |
Counter | Messages successfully processed |
rabbitmq.consumer.messages.failed |
Counter | Messages that failed processing |
rabbitmq.consumer.messages.retried |
Counter | Messages republished for retry |
rabbitmq.consumer.processing.duration |
Histogram | Processing time (ms) |
rabbitmq.publisher.messages.published |
Counter | Messages published |
rabbitmq.publisher.messages.failed |
Counter | Publish failures |
rabbitmq.publisher.duration |
Histogram | Publish time (ms) |
rabbitmq.healthcheck.total |
Counter | Health check invocations |
rabbitmq.healthcheck.failed |
Counter | Failed health checks |
rabbitmq.healthcheck.duration |
Histogram | Health check time (ms) |
Subscribe in your OpenTelemetry setup:
builder.Services.AddOpenTelemetry()
.WithTracing(t => t.AddSource("LowCodeHub.RabbitMQ"))
.WithMetrics(m => m.AddMeter("LowCodeHub.RabbitMQ"));
Clustering
Configure multiple endpoints for high availability:
{
"RabbitMQ": {
"HostName": "rabbit-1",
"Port": 5672,
"UserName": "myapp",
"Password": "secret",
"Endpoints": [
{ "HostName": "rabbit-1", "Port": 5672 },
{ "HostName": "rabbit-2", "Port": 5672 },
{ "HostName": "rabbit-3", "Port": 5672 }
]
}
}
The client connects to the first available node. SSL settings are propagated to all endpoints. On connection loss, automatic recovery reconnects to any available node with topology re-declaration.
SSL/TLS
{
"RabbitMQ": {
"HostName": "rabbitmq.prod",
"Port": 5671,
"Ssl": {
"Enabled": true,
"ServerName": "rabbitmq.prod",
"CertificatePath": "/certs/client.pfx",
"CertificatePassword": "cert-password",
"AllowInsecureCertificateValidation": false
}
}
}
⚠️
AllowInsecureCertificateValidation: truedisables certificate verification. A startup warning is logged. Do not use in production.
Conflict Handling
When your topology changes (e.g., switching a queue from classic to quorum), the broker rejects re-declaration with PRECONDITION_FAILED (406).
Default behavior — fail fast on startup:
topology.Settings(s => s.DeleteAndRecreateOnConflict(false));
Auto-recreate — delete the conflicting entity and recreate it:
topology.Settings(s => s.DeleteAndRecreateOnConflict());
⚠️ Recreating a queue deletes all messages in it. Use for development/staging or planned migrations only.
Custom Serializer
Default: System.Text.Json with JsonSerializerDefaults.Web.
Replace globally:
builder.Services.UseMessageSerializer<MessagePackSerializer>();
Implement IMessageSerializer:
public sealed class MessagePackSerializer : IMessageSerializer
{
public string ContentType => "application/x-msgpack";
public byte[] Serialize<T>(T message) { /* ... */ }
public T? Deserialize<T>(ReadOnlySpan<byte> data) { /* ... */ }
public object? Deserialize(ReadOnlySpan<byte> data, Type type) { /* ... */ }
}
Architecture Overview
AddRabbitMq(config, topology => { ... })
│
├─ RabbitMqMessagingConfiguration ← Fluent DSL input
│ │
│ ▼
├─ RabbitMqTopologyCompiler ← Validates + compiles to declarations
│ │
│ ▼
├─ RabbitMqCompiledTopology ← Exchanges, queues, bindings, name mappings
│
├─ DI Registration
│ ├─ IRabbitMqConnectionManager ← Connection + channel lifecycle (singleton)
│ ├─ IMessagePublisher ← Publish API (scoped)
│ ├─ IMessageRpcClient ← RPC client (opt-in, singleton)
│ ├─ IMessageSerializer ← JSON by default (singleton)
│ └─ RabbitMqHealthCheck ← ASP.NET health check
│
└─ RabbitMqConsumerHostedService ← IHostedService
├─ TopologyInitializer ← Declares topology on broker at startup
├─ Consumer auto-discovery ← Scans assemblies for [Queue] consumers
├─ Per-queue dedicated channel ← Isolated QoS per consumer
└─ Channel recovery ← Auto-restart consumers on channel loss
Production Checklist
- Use quorum queues for business-critical workloads
- Configure dead letter queues for failed message inspection
- Set
AutoAck = false(default) for reliable message processing - Keep consumers idempotent — at-least-once delivery under retries
- Use publisher confirms for critical publishes
- Set proper credentials — startup warns on
guestusage - Enable SSL/TLS for production brokers
- Configure cluster endpoints for high availability
- Subscribe to
LowCodeHub.RabbitMQactivity source and meter for observability - Set
DeleteAndRecreateOnConflict(false)in production (default) - Review startup topology report logs for warnings
License
Internal package — see repository license.
| 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.Configuration.Binder (>= 10.0.0)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Diagnostics.HealthChecks (>= 10.0.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.0)
- Microsoft.Extensions.Options (>= 10.0.0)
- RabbitMQ.Client (>= 7.1.2)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.