LowCodeHub.RabbitMQ 0.0.1

There is a newer version of this package available.
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
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="LowCodeHub.RabbitMQ" Version="0.0.1" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="LowCodeHub.RabbitMQ" Version="0.0.1" />
                    
Directory.Packages.props
<PackageReference Include="LowCodeHub.RabbitMQ" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add LowCodeHub.RabbitMQ --version 0.0.1
                    
#r "nuget: LowCodeHub.RabbitMQ, 0.0.1"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package LowCodeHub.RabbitMQ@0.0.1
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=LowCodeHub.RabbitMQ&version=0.0.1
                    
Install as a Cake Addin
#tool nuget:?package=LowCodeHub.RabbitMQ&version=0.0.1
                    
Install as a Cake Tool

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

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

x.Queue("orders.created", q => q.DeadLetterDefaultExchange());

Creates:

  • Main queue: orders.created with x-dead-letter-exchange: "" and x-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.dead bound to dlx.orders with routing key failed

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:

  1. OnBeforeConsumeAsync — optional pre-processing gate. Return false to skip the message and auto-ack it.
  2. ConsumeAsync — your business logic (required override).
  3. OnAfterConsumeAsync — optional post-processing hook after successful consumption.
  4. If ConsumeAsync (or any stage) throws, OnErrorAsync is 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:

  1. Reads the x-retry-count header from the message (starts at 0).
  2. If retryCount < MaxRetryAttempts, it republishes the message to the same exchange/routing-key with an incremented x-retry-count header, then acks the original delivery.
  3. 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:

  1. ConsumeAsync throws an HttpRequestException → OnErrorAsync returns Requeue.
  2. The message is republished with x-retry-count = 1 and re-delivered.
  3. Steps 1–2 repeat up to 5 times.
  4. After 5 retries the message is rejected and routed to payments.process.dlq for inspection.
  5. If ConsumeAsync throws a non-transient exception (e.g. ArgumentException), OnErrorAsync returns Reject and 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: true disables 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 guest usage
  • Enable SSL/TLS for production brokers
  • Configure cluster endpoints for high availability
  • Subscribe to LowCodeHub.RabbitMQ activity 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 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

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.0.11 129 7/9/2026
0.0.10 127 6/21/2026
0.0.4 111 5/18/2026
0.0.3 118 5/12/2026
0.0.2 124 4/23/2026
0.0.1 124 3/26/2026