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

Information
Class: AsyncResponse.Transports.AzureServiceBus.QueuedAzureServiceBusMessageDispatcher
Assembly: AsyncResponse.Transports.AzureServiceBus
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.AzureServiceBus/AzureServiceBusMessageDispatcher.cs
Line coverage
100%
Covered lines: 107
Uncovered lines: 0
Coverable lines: 107
Total lines: 393
Line coverage: 100%
Branch coverage
100%
Covered branches: 6
Total branches: 6
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
get_PendingCount()100%11100%
get_RunningCount()100%11100%
get_CanAcceptMore()100%11100%
get_FreeCapacity()100%11100%
WaitForCapacityAsync()100%2275%
HandleAsync()100%22100%
DisposeAsync()100%22100%
RunWorkerAsync()100%11100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.AzureServiceBus/AzureServiceBusMessageDispatcher.cs

#LineLine coverage
 1using Microsoft.Extensions.Logging;
 2using System.Diagnostics;
 3using System.Threading.Channels;
 4
 5namespace AsyncResponse.Transports.AzureServiceBus;
 6
 7internal enum AzureServiceBusSubscriberRole
 8{
 9    Worker,
 10    ResponseIngress
 11}
 12
 13internal abstract class AzureServiceBusMessageDispatcher : IAsyncDisposable
 14{
 15    private readonly Func<AzureServiceBusTransportDelivery, CancellationToken, Task> _handler;
 16    private readonly AzureServiceBusAsyncResponseOptions _transportOptions;
 17    private readonly AzureServiceBusSubscriberOptions _subscriberOptions;
 18    private readonly string _queue;
 19    private readonly AzureServiceBusSubscriberRole _role;
 20
 21    protected AzureServiceBusMessageDispatcher(
 22        Func<AzureServiceBusTransportDelivery, CancellationToken, Task> handler,
 23        AzureServiceBusAsyncResponseOptions transportOptions,
 24        AzureServiceBusSubscriberOptions subscriberOptions,
 25        ILogger logger,
 26        string queue,
 27        AzureServiceBusSubscriberRole role)
 28    {
 29        _handler = handler;
 30        _transportOptions = transportOptions;
 31        _subscriberOptions = subscriberOptions;
 32        Logger = logger;
 33        _queue = queue;
 34        _role = role;
 35    }
 36
 37    protected ILogger Logger { get; }
 38    protected int MaxDeliveryAttempts => _subscriberOptions.MaxDeliveryAttempts;
 39
 40    /// <summary>Creates the dispatcher configured by the subscriber options.</summary>
 41    public static AzureServiceBusMessageDispatcher Create(
 42        Func<AzureServiceBusTransportDelivery, CancellationToken, Task> handler,
 43        AzureServiceBusAsyncResponseOptions transportOptions,
 44        AzureServiceBusSubscriberOptions subscriberOptions,
 45        ILogger logger,
 46        string queue,
 47        AzureServiceBusSubscriberRole role)
 48    {
 49        AzureServiceBusOptionsValidator.ValidateSubscriber(transportOptions, subscriberOptions, role);
 50
 51        return subscriberOptions.AckMode == AzureServiceBusAckMode.AckAfterHandlerCompletes
 52            ? new AwaitingAzureServiceBusMessageDispatcher(
 53                handler,
 54                transportOptions,
 55                subscriberOptions,
 56                logger,
 57                queue,
 58                role)
 59            : new QueuedAzureServiceBusMessageDispatcher(
 60                handler,
 61                transportOptions,
 62                subscriberOptions,
 63                logger,
 64                queue,
 65                role);
 66    }
 67
 68    /// <summary>Validates the supplied subscriber options.</summary>
 69    public static void ValidateOptions(
 70        AzureServiceBusAsyncResponseOptions transportOptions,
 71        AzureServiceBusSubscriberOptions subscriberOptions,
 72        AzureServiceBusSubscriberRole role)
 73        => AzureServiceBusOptionsValidator.ValidateSubscriber(transportOptions, subscriberOptions, role);
 74
 75    /// <summary>Handles the delivered message.</summary>
 76    public abstract Task HandleAsync(
 77        AzureServiceBusTransportDelivery delivery,
 78        CancellationToken subscriberCancellationToken);
 79
 80    /// <summary>
 81    /// Whether the dispatcher can accept more deliveries right now. Awaiting dispatchers always can
 82    /// (handlers run inline); the queued dispatcher returns <c>false</c> while its bounded queue is
 83    /// saturated so the receive loop stops pulling messages instead of receiving and abandoning them —
 84    /// every abandon burns <c>DeliveryCount</c> toward the entity's MaxDeliveryCount.
 85    /// </summary>
 86    public virtual bool CanAcceptMore => true;
 87
 88    /// <summary>
 89    /// Number of deliveries the dispatcher can accept right now. The receive loop requests at most
 90    /// this many messages per receive in early-ACK mode so a burst never overflows the background queue.
 91    /// </summary>
 92    public virtual int FreeCapacity => int.MaxValue;
 93
 94    /// <summary>
 95    /// Waits until the dispatcher can accept at least one more delivery. Completes immediately for
 96    /// awaiting dispatchers; the queued dispatcher waits for a background worker to free a slot.
 97    /// </summary>
 98    public virtual ValueTask WaitForCapacityAsync(CancellationToken cancellationToken) => ValueTask.CompletedTask;
 99
 100    /// <summary>Releases resources held by this instance.</summary>
 101    public virtual ValueTask DisposeAsync() => ValueTask.CompletedTask;
 102
 103    protected async Task ExecuteHandlerAsync(
 104        AzureServiceBusTransportDelivery delivery,
 105        CancellationToken cancellationToken,
 106        bool logFailures = true)
 107    {
 108        using var activity = AsyncResponseDiagnostics.StartActivity(
 109            "asyncresponse.azure_service_bus.receive",
 110            ActivityKind.Consumer);
 111        activity?.SetTag("asyncresponse.transport", "azure_service_bus");
 112        activity?.SetTag("asyncresponse.azure_service_bus.role", _role.ToString());
 113        activity?.SetTag("asyncresponse.azure_service_bus.ack_mode", _subscriberOptions.AckMode.ToString());
 114        activity?.SetTag("messaging.system", "azure_service_bus");
 115        activity?.SetTag("messaging.destination.name", _queue);
 116        activity?.SetTag("messaging.message.id", delivery.MessageId);
 117        activity?.SetTag("messaging.azure_service_bus.sequence_number", delivery.SequenceNumber);
 118
 119        if (!string.IsNullOrWhiteSpace(delivery.CorrelationId))
 120            AsyncResponseDiagnostics.SetCorrelationId(activity, delivery.CorrelationId);
 121
 122        try
 123        {
 124            await _handler(delivery, cancellationToken).ConfigureAwait(false);
 125        }
 126        catch (Exception ex)
 127        {
 128            if (logFailures)
 129                Logger.LogError(ex, "Azure Service Bus message handling failed for message {MessageId}.", delivery.Messa
 130            AsyncResponseDiagnostics.SetError(activity, ex);
 131            throw;
 132        }
 133    }
 134
 135    protected async ValueTask NotifyBackgroundFailureAsync(
 136        AzureServiceBusTransportDelivery delivery,
 137        Exception exception,
 138        string queue,
 139        AzureServiceBusSubscriberRole role)
 140    {
 141        var callback = _subscriberOptions.OnBackgroundFailure;
 142        if (callback is null)
 143            return;
 144
 145        try
 146        {
 147            var context = new AzureServiceBusBackgroundFailureContext(
 148                queue,
 149                role.ToString(),
 150                delivery.SequenceNumber,
 151                delivery.MessageId,
 152                delivery.CorrelationId ?? TryReadApplicationCorrelationId(delivery),
 153                exception);
 154            await callback(context).ConfigureAwait(false);
 155        }
 156        catch (Exception callbackException)
 157        {
 158            Logger.LogError(
 159                callbackException,
 160                "Azure Service Bus background failure callback failed for already-completed message {MessageId} on {Queu
 161                delivery.MessageId,
 162                queue);
 163        }
 164    }
 165
 166    private string? TryReadApplicationCorrelationId(AzureServiceBusTransportDelivery delivery)
 167    {
 168        if (!string.IsNullOrWhiteSpace(_transportOptions.CorrelationIdProperty)
 169            && delivery.ApplicationProperties.TryGetValue(_transportOptions.CorrelationIdProperty, out var value))
 170        {
 171            return AzureServiceBusCorrelationIdExtractor.TryConvertProperty(value);
 172        }
 173
 174        return null;
 175    }
 176}
 177
 178internal sealed class AwaitingAzureServiceBusMessageDispatcher(
 179    Func<AzureServiceBusTransportDelivery, CancellationToken, Task> handler,
 180    AzureServiceBusAsyncResponseOptions transportOptions,
 181    AzureServiceBusSubscriberOptions subscriberOptions,
 182    ILogger logger,
 183    string queue,
 184    AzureServiceBusSubscriberRole role)
 185    : AzureServiceBusMessageDispatcher(handler, transportOptions, subscriberOptions, logger, queue, role)
 186{
 187    /// <summary>Handles the delivered message.</summary>
 188    public override async Task HandleAsync(
 189        AzureServiceBusTransportDelivery delivery,
 190        CancellationToken subscriberCancellationToken)
 191    {
 192        try
 193        {
 194            await ExecuteHandlerAsync(delivery, subscriberCancellationToken).ConfigureAwait(false);
 195            await delivery.CompleteAsync().ConfigureAwait(false);
 196        }
 197        catch (Exception ex)
 198        {
 199            if (MaxDeliveryAttempts > 0 && delivery.DeliveryCount >= MaxDeliveryAttempts)
 200            {
 201                await delivery.DeadLetterAsync(
 202                    "AsyncResponseHandlerFailed",
 203                    ex.Message).ConfigureAwait(false);
 204                return;
 205            }
 206
 207            await delivery.AbandonAsync().ConfigureAwait(false);
 208        }
 209    }
 210}
 211
 212internal sealed class QueuedAzureServiceBusMessageDispatcher : AzureServiceBusMessageDispatcher
 213{
 214    private readonly Channel<AzureServiceBusTransportDelivery> _queue;
 215    private readonly Task[] _workers;
 3216    private readonly CancellationTokenSource _drainCancellation = new();
 217    private readonly TimeSpan _drainTimeout;
 218    private readonly int _capacity;
 219    private readonly string _queueName;
 220    private readonly AzureServiceBusSubscriberRole _role;
 221    private int _pendingCount;
 222    private int _runningCount;
 223    private int _disposeStarted;
 224
 225    /// <summary>Creates an ACK-after-enqueue dispatcher with a bounded background queue.</summary>
 226    public QueuedAzureServiceBusMessageDispatcher(
 227        Func<AzureServiceBusTransportDelivery, CancellationToken, Task> handler,
 228        AzureServiceBusAsyncResponseOptions transportOptions,
 229        AzureServiceBusSubscriberOptions subscriberOptions,
 230        ILogger logger,
 231        string queue,
 232        AzureServiceBusSubscriberRole role)
 3233        : base(handler, transportOptions, subscriberOptions, logger, queue, role)
 234    {
 3235        _drainTimeout = subscriberOptions.BackgroundDrainTimeout;
 3236        _capacity = subscriberOptions.BackgroundQueueCapacity;
 3237        _queueName = queue;
 3238        _role = role;
 3239        _queue = Channel.CreateBounded<AzureServiceBusTransportDelivery>(new BoundedChannelOptions(subscriberOptions.Bac
 3240        {
 3241            AllowSynchronousContinuations = false,
 3242            FullMode = BoundedChannelFullMode.Wait,
 3243            SingleReader = subscriberOptions.BackgroundWorkerCount == 1,
 3244            SingleWriter = false
 3245        });
 246
 3247        _workers = Enumerable.Range(0, subscriberOptions.BackgroundWorkerCount)
 3248            .Select(workerIndex => Task.Run(() => RunWorkerAsync(workerIndex)))
 3249            .ToArray();
 250
 3251        Logger.LogInformation(
 3252            "Created Azure Service Bus ACK-after-enqueue dispatcher for {Queue} with {WorkerCount} worker(s), queue capa
 3253            _queueName,
 3254            subscriberOptions.BackgroundWorkerCount,
 3255            subscriberOptions.BackgroundQueueCapacity,
 3256            _drainTimeout);
 3257    }
 258
 3259    internal int PendingCount => Volatile.Read(ref _pendingCount);
 3260    internal int RunningCount => Volatile.Read(ref _runningCount);
 261
 3262    public override bool CanAcceptMore => Volatile.Read(ref _pendingCount) < _capacity;
 263
 3264    public override int FreeCapacity => Math.Max(0, _capacity - Volatile.Read(ref _pendingCount));
 265
 266    public override async ValueTask WaitForCapacityAsync(CancellationToken cancellationToken)
 267    {
 268        // WaitToWriteAsync completes when the bounded channel has room (or the channel is completed
 269        // during dispose, in which case there is nothing left to gate).
 3270        while (!CanAcceptMore)
 271        {
 2272            if (!await _queue.Writer.WaitToWriteAsync(cancellationToken).ConfigureAwait(false))
 1273                return;
 274        }
 3275    }
 276
 277    /// <summary>Handles the delivered message.</summary>
 278    public override async Task HandleAsync(
 279        AzureServiceBusTransportDelivery delivery,
 280        CancellationToken subscriberCancellationToken)
 281    {
 3282        Interlocked.Increment(ref _pendingCount);
 3283        if (!_queue.Writer.TryWrite(delivery))
 284        {
 285            // The receive loop gates on free capacity, so this only covers the residual race between
 286            // its capacity check and this write. The abandon burns one DeliveryCount, but the loop
 287            // never receives while saturated, so a healthy message cannot repeat this path toward
 288            // the entity's MaxDeliveryCount.
 2289            Interlocked.Decrement(ref _pendingCount);
 2290            Logger.LogWarning(
 2291                "Azure Service Bus background queue rejected message {MessageId} for {Queue}; abandoning for redelivery.
 2292                delivery.MessageId,
 2293                _queueName,
 2294                PendingCount,
 2295                RunningCount);
 2296            await delivery.AbandonAsync().ConfigureAwait(false);
 3297            return;
 298        }
 299
 300        // The delivery now belongs to a background worker, which decrements _pendingCount when it dequeues.
 301        // Do not touch the counter or abandon here, even if the Complete below fails — the message is already
 302        // executing in-process and abandoning it would trigger a duplicate execution via redelivery.
 303        try
 304        {
 3305            await delivery.CompleteAsync().ConfigureAwait(false);
 3306        }
 2307        catch (Exception ex)
 308        {
 2309            Logger.LogError(
 2310                ex,
 2311                "Failed to complete Azure Service Bus message {MessageId} for {Queue} after enqueue; it is being process
 2312                delivery.MessageId,
 2313                _queueName);
 3314        }
 3315    }
 316
 317    /// <summary>Releases resources held by this instance.</summary>
 318    public override async ValueTask DisposeAsync()
 319    {
 3320        if (Interlocked.Exchange(ref _disposeStarted, 1) != 0)
 3321            return;
 322
 3323        Logger.LogInformation(
 3324            "Draining Azure Service Bus ACK-after-enqueue dispatcher for {Queue}. Pending={PendingCount}, Running={Runni
 3325            _queueName,
 3326            PendingCount,
 3327            RunningCount);
 3328        _queue.Writer.TryComplete();
 329
 330        try
 331        {
 3332            await Task.WhenAll(_workers).WaitAsync(_drainTimeout).ConfigureAwait(false);
 3333            _drainCancellation.Dispose();
 3334        }
 2335        catch (TimeoutException ex)
 336        {
 2337            _drainCancellation.Cancel();
 2338            Logger.LogWarning(
 2339                ex,
 2340                "Timed out while draining Azure Service Bus ACK-after-enqueue dispatcher for {Queue}. Pending={PendingCo
 2341                _queueName,
 2342                PendingCount,
 2343                RunningCount);
 344
 2345            _ = Task.WhenAll(_workers).ContinueWith(
 3346                _ => _drainCancellation.Dispose(),
 2347                CancellationToken.None,
 2348                TaskContinuationOptions.ExecuteSynchronously,
 2349                TaskScheduler.Default);
 3350        }
 3351    }
 352
 353    private async Task RunWorkerAsync(int workerIndex)
 354    {
 3355        await foreach (var delivery in _queue.Reader.ReadAllAsync().ConfigureAwait(false))
 356        {
 3357            Interlocked.Decrement(ref _pendingCount);
 3358            Interlocked.Increment(ref _runningCount);
 359
 360            try
 361            {
 3362                Logger.LogDebug(
 3363                    "Azure Service Bus background worker {WorkerIndex} handling message {MessageId} for {Queue}. Pending
 3364                    workerIndex,
 3365                    delivery.MessageId,
 3366                    _queueName,
 3367                    PendingCount,
 3368                    RunningCount);
 3369                await ExecuteHandlerAsync(
 3370                    delivery,
 3371                    _drainCancellation.Token,
 3372                    logFailures: false).ConfigureAwait(false);
 3373            }
 3374            catch (Exception ex)
 375            {
 2376                Logger.LogError(
 2377                    ex,
 2378                    "Azure Service Bus background handler failed for already-completed message {MessageId} on {Queue}.",
 2379                    delivery.MessageId,
 2380                    _queueName);
 2381                await NotifyBackgroundFailureAsync(
 2382                    delivery,
 2383                    ex,
 2384                    _queueName,
 2385                    _role).ConfigureAwait(false);
 386            }
 387            finally
 388            {
 3389                Interlocked.Decrement(ref _runningCount);
 390            }
 3391        }
 3392    }
 393}