Skip to main content

Curiosus.RabbitMQ

NuGet Downloads Coverage

Request-reply (RPC) client for RabbitMQ: sends a JSON request to a queue and awaits the reply on a dedicated response queue, with automatic and manual connection recovery and resending. Use it to call services that process requests from RabbitMQ. Built on the async API of RabbitMQ.Client 7.

Installation​

dotnet add package Curiosus.RabbitMQ

Usage​

RabbitMQ:
HostName: localhost
Port: 5672
UserName: guest
Password: guest
ExchangeName: ""
ClientName: billing-api
services.AddRabbitMQRPC(configuration.RabbitMQ); // validates options, registers RabbitMqRpcClientFactory

// default correlation ids come from UniqueIdGenerator (Curiosus.Tools): initialize it once per process
UniqueIdGenerator.Initialize(generatorId: 1);

// rpcClientFactory is an injected RabbitMqRpcClientFactory; CreateClientAsync connects and declares both queues
await using var client = await rpcClientFactory.CreateClientAsync(
"balance_requests",
cancellationToken: cancellationToken);

var response = await client.SendWithAutoAcknowledgeAsync<BalanceResponse, BalanceRequest>(
new BalanceRequest(accountId),
cancellationToken: cancellationToken);

The reply queue is named {requestQueue}_responses_{ClientName}; pass clientNameSuffix to CreateClientAsync when one process needs several clients for the same queue. Dispose the client (await using) to close its connection.

SendWithManualAcknowledgeAsync returns ManualAckRabbitResult<T>: process Data and await ConfirmAcknowledgeAsync() to ack the reply only after it was handled:

var result = await client.SendWithManualAcknowledgeAsync<BalanceResponse, BalanceRequest>(
new BalanceRequest(accountId),
cancellationToken: cancellationToken);

await SaveBalanceAsync(result.Data, cancellationToken);
await result.ConfirmAcknowledgeAsync(cancellationToken);

GetConsumersCountAsync returns the count of consumers of the request queue.

Messages are serialized with System.Text.Json. RabbitMqRpcClient.DefaultJsonSerializerOptions stay wire-compatible with the Newtonsoft.Json format of 1.x; pass jsonSerializerOptions to CreateClientAsync to override them.

See also​