< Summary - AsyncResponse (Release / net8.0+net10.0 / unit+integration)

Information
Class: AsyncResponse.Transports.RabbitMQ.RabbitMqTopology
Assembly: AsyncResponse.Transports.RabbitMQ
File(s): /_/src/Transports/AsyncResponse.Transports.RabbitMQ/RabbitMqTopology.cs
Line coverage
98%
Covered lines: 53
Uncovered lines: 1
Coverable lines: 54
Total lines: 130
Line coverage: 98.1%
Branch coverage
94%
Covered branches: 17
Total branches: 18
Branch coverage: 94.4%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
EnsureWorkerAsync(...)100%22100%
EnsureResponseAsync(...)100%22100%
DeclareQueueTopologyAsync()50%2285.71%
EnsureDeadLetterTopologyAsync()100%88100%
CreatePersistentJsonProperties(...)100%44100%

File(s)

/_/src/Transports/AsyncResponse.Transports.RabbitMQ/RabbitMqTopology.cs

#LineLine coverage
 1using RabbitMQ.Client;
 2
 3namespace AsyncResponse.Transports.RabbitMQ;
 4
 5internal static class RabbitMqTopology
 6{
 7    private const string DirectExchange = "direct";
 8
 9    /// <summary>Ensures the required resource exists.</summary>
 10    public static Task EnsureWorkerAsync(
 11        IRabbitMqChannel channel,
 12        RabbitMqAsyncResponseOptions options,
 13        CancellationToken cancellationToken = default)
 14    {
 43915        if (!options.DeclareTopology)
 216            return Task.CompletedTask;
 17
 43718        var exchange = RabbitMqOptionsValidator.Required(options.WorkerExchange, nameof(options.WorkerExchange));
 43719        var queue = RabbitMqOptionsValidator.Required(options.WorkerQueue, nameof(options.WorkerQueue));
 43720        var routingKey = RabbitMqOptionsValidator.Required(options.WorkerRoutingKey, nameof(options.WorkerRoutingKey));
 21
 43722        return DeclareQueueTopologyAsync(channel, options, exchange, queue, routingKey, cancellationToken);
 23    }
 24
 25    /// <summary>Ensures the required resource exists.</summary>
 26    public static Task EnsureResponseAsync(
 27        IRabbitMqChannel channel,
 28        RabbitMqAsyncResponseOptions options,
 29        CancellationToken cancellationToken = default)
 30    {
 20031        if (!options.DeclareTopology)
 232            return Task.CompletedTask;
 33
 19834        var exchange = RabbitMqOptionsValidator.Required(options.ResponseExchange, nameof(options.ResponseExchange));
 19835        var queue = RabbitMqOptionsValidator.Required(options.ResponseQueue, nameof(options.ResponseQueue));
 19836        var routingKey = RabbitMqOptionsValidator.Required(options.ResponseRoutingKey, nameof(options.ResponseRoutingKey
 37
 19838        return DeclareQueueTopologyAsync(channel, options, exchange, queue, routingKey, cancellationToken);
 39    }
 40
 41    private static async Task DeclareQueueTopologyAsync(
 42        IRabbitMqChannel channel,
 43        RabbitMqAsyncResponseOptions options,
 44        string exchange,
 45        string queue,
 46        string routingKey,
 47        CancellationToken cancellationToken)
 48    {
 63549        await channel.ExchangeDeclareAsync(exchange, DirectExchange, durable: true, autoDelete: false, cancellationToken
 50
 63351        var queueArguments = await EnsureDeadLetterTopologyAsync(channel, options, routingKey, cancellationToken).Config
 52
 53        // The park queue is declared on its own, never bound: it is reached only through the default
 54        // exchange (routing key = queue name) by the park-at-cap publish. Bound to the dead-letter
 55        // exchange — as DeadLetterQueue is — it would also collect one copy per retry hop of a
 56        // TTL-retry cycle, including the hops of messages that later succeed. Declared with or
 57        // without a dead-letter exchange: a broker policy can supply that exchange, and a park
 58        // destination nobody declared turns every park into a failed publish.
 63359        if (!string.IsNullOrWhiteSpace(options.ParkQueue))
 060            await channel.QueueDeclareAsync(options.ParkQueue, durable: true, exclusive: false, autoDelete: false, argum
 61
 63362        await channel.QueueDeclareAsync(queue, durable: true, exclusive: false, autoDelete: false, queueArguments, cance
 63363        await channel.QueueBindAsync(queue, exchange, routingKey, cancellationToken).ConfigureAwait(false);
 63364    }
 65
 66    /// <summary>
 67    /// Declares the dead-letter exchange/queue/binding when configured and returns the arguments to apply to the
 68    /// source queue (<c>x-dead-letter-exchange</c> and optionally <c>x-dead-letter-routing-key</c>), or null when
 69    /// no dead-letter exchange is configured.
 70    /// </summary>
 71    private static async Task<IDictionary<string, object?>?> EnsureDeadLetterTopologyAsync(
 72        IRabbitMqChannel channel,
 73        RabbitMqAsyncResponseOptions options,
 74        string sourceRoutingKey,
 75        CancellationToken cancellationToken)
 76    {
 63377        if (string.IsNullOrWhiteSpace(options.DeadLetterExchange))
 61978            return null;
 79
 1480        var deadLetterExchange = options.DeadLetterExchange;
 1481        await channel.ExchangeDeclareAsync(deadLetterExchange, DirectExchange, durable: true, autoDelete: false, cancell
 82
 1483        if (!string.IsNullOrWhiteSpace(options.DeadLetterQueue))
 84        {
 1085            var bindingRoutingKey = string.IsNullOrWhiteSpace(options.DeadLetterRoutingKey)
 1086                ? sourceRoutingKey
 1087                : options.DeadLetterRoutingKey;
 88
 1089            await channel.QueueDeclareAsync(options.DeadLetterQueue, durable: true, exclusive: false, autoDelete: false,
 1090            await channel.QueueBindAsync(options.DeadLetterQueue, deadLetterExchange, bindingRoutingKey, cancellationTok
 1091        }
 92
 1493        var arguments = new Dictionary<string, object?>(StringComparer.Ordinal)
 1494        {
 1495            ["x-dead-letter-exchange"] = deadLetterExchange
 1496        };
 97
 1498        if (!string.IsNullOrWhiteSpace(options.DeadLetterRoutingKey))
 899            arguments["x-dead-letter-routing-key"] = options.DeadLetterRoutingKey;
 100
 14101        return arguments;
 633102    }
 103
 104    /// <summary>Creates the requested resource.</summary>
 105    public static BasicProperties CreatePersistentJsonProperties(string? correlationId, string correlationHeader)
 106    {
 463107        var properties = new BasicProperties
 463108        {
 463109            ContentType = "application/json",
 463110            DeliveryMode = DeliveryModes.Persistent,
 463111            Persistent = true,
 463112            MessageId = Guid.NewGuid().ToString("N"),
 463113            Timestamp = new AmqpTimestamp(DateTimeOffset.UtcNow.ToUnixTimeSeconds())
 463114        };
 115
 463116        if (!string.IsNullOrWhiteSpace(correlationId))
 117        {
 135118            properties.CorrelationId = correlationId;
 135119            if (!string.IsNullOrWhiteSpace(correlationHeader))
 120            {
 133121                properties.Headers = new Dictionary<string, object?>(StringComparer.Ordinal)
 133122                {
 133123                    [correlationHeader] = correlationId
 133124                };
 125            }
 126        }
 127
 463128        return properties;
 129    }
 130}