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

Information
Class: AsyncResponse.Transports.SQS.SqsMessageDispatcher
Assembly: AsyncResponse.Transports.SQS
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.SQS/SqsMessageDispatcher.cs
Line coverage
98%
Covered lines: 90
Uncovered lines: 1
Coverable lines: 91
Total lines: 408
Line coverage: 98.9%
Branch coverage
71%
Covered branches: 20
Total branches: 28
Branch coverage: 71.4%
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_RedeliveryDelay()100%11100%
Create(...)100%22100%
ValidateOptions(...)100%11100%
get_CanAcceptMore()100%11100%
get_FreeCapacity()100%11100%
WaitForCapacityAsync(...)100%11100%
DisposeAsync()100%11100%
ExecuteHandlerAsync()61.11%1818100%
TryChangeVisibilityAsync()100%11100%
NotifyBackgroundFailureAsync()50%2295%
TryReadCorrelationId(...)100%66100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.SQS/SqsMessageDispatcher.cs

#LineLine coverage
 1using Microsoft.Extensions.Logging;
 2using System.Diagnostics;
 3using System.Threading.Channels;
 4
 5namespace AsyncResponse.Transports.SQS;
 6
 7internal enum SqsSubscriberRole
 8{
 9    Worker,
 10    ResponseIngress
 11}
 12
 13internal abstract class SqsMessageDispatcher : IAsyncDisposable
 14{
 15    private readonly Func<SqsTransportDelivery, CancellationToken, Task> _handler;
 16    private readonly SqsAsyncResponseOptions _transportOptions;
 17    private readonly SqsSubscriberOptions _subscriberOptions;
 18    private readonly string _queue;
 19    private readonly SqsSubscriberRole _role;
 20
 321    protected SqsMessageDispatcher(
 322        Func<SqsTransportDelivery, CancellationToken, Task> handler,
 323        SqsAsyncResponseOptions transportOptions,
 324        SqsSubscriberOptions subscriberOptions,
 325        ILogger logger,
 326        string queue,
 327        SqsSubscriberRole role)
 28    {
 329        _handler = handler;
 330        _transportOptions = transportOptions;
 331        _subscriberOptions = subscriberOptions;
 332        Logger = logger;
 333        _queue = queue;
 334        _role = role;
 335    }
 36
 37    protected ILogger Logger { get; }
 238    protected TimeSpan? RedeliveryDelay => _subscriberOptions.RedeliveryDelay;
 39
 40    /// <summary>Creates the dispatcher configured by the subscriber options.</summary>
 41    public static SqsMessageDispatcher Create(
 42        Func<SqsTransportDelivery, CancellationToken, Task> handler,
 43        SqsAsyncResponseOptions transportOptions,
 44        SqsSubscriberOptions subscriberOptions,
 45        ILogger logger,
 46        string queue,
 47        SqsSubscriberRole role)
 48    {
 349        SqsOptionsValidator.ValidateSubscriber(transportOptions, subscriberOptions, role);
 50
 351        return subscriberOptions.AckMode == SqsAckMode.AckAfterHandlerCompletes
 352            ? new AwaitingSqsMessageDispatcher(
 353                handler,
 354                transportOptions,
 355                subscriberOptions,
 356                logger,
 357                queue,
 358                role)
 359            : new QueuedSqsMessageDispatcher(
 360                handler,
 361                transportOptions,
 362                subscriberOptions,
 363                logger,
 364                queue,
 365                role);
 66    }
 67
 68    /// <summary>Validates the supplied subscriber options.</summary>
 69    public static void ValidateOptions(
 70        SqsAsyncResponseOptions transportOptions,
 71        SqsSubscriberOptions subscriberOptions,
 72        SqsSubscriberRole role)
 373        => SqsOptionsValidator.ValidateSubscriber(transportOptions, subscriberOptions, role);
 74
 75    /// <summary>Handles the delivered message.</summary>
 76    public abstract Task HandleAsync(
 77        SqsTransportDelivery 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 releasing them —
 84    /// SQS counts every receive toward the queue's redrive policy.
 85    /// </summary>
 286    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>
 392    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>
 398    public virtual ValueTask WaitForCapacityAsync(CancellationToken cancellationToken) => ValueTask.CompletedTask;
 99
 100    /// <summary>Releases resources held by this instance.</summary>
 3101    public virtual ValueTask DisposeAsync() => ValueTask.CompletedTask;
 102
 103    protected async Task ExecuteHandlerAsync(
 104        SqsTransportDelivery delivery,
 105        CancellationToken cancellationToken,
 106        bool logFailures = true)
 107    {
 3108        using var activity = AsyncResponseDiagnostics.StartActivity(
 3109            "asyncresponse.sqs.receive",
 3110            ActivityKind.Consumer);
 3111        activity?.SetTag("asyncresponse.transport", "aws_sqs");
 3112        activity?.SetTag("asyncresponse.sqs.role", _role.ToString());
 3113        activity?.SetTag("asyncresponse.sqs.ack_mode", _subscriberOptions.AckMode.ToString());
 3114        activity?.SetTag("messaging.system", "aws_sqs");
 3115        activity?.SetTag("messaging.destination.name", _queue);
 3116        activity?.SetTag("messaging.message.id", delivery.MessageId);
 3117        activity?.SetTag("messaging.aws_sqs.receive_count", delivery.ReceiveCount);
 118
 3119        if (TryReadCorrelationId(delivery) is { } correlationId)
 3120            AsyncResponseDiagnostics.SetCorrelationId(activity, correlationId);
 121
 122        try
 123        {
 3124            await _handler(delivery, cancellationToken).ConfigureAwait(false);
 3125        }
 2126        catch (Exception ex)
 127        {
 2128            if (logFailures)
 2129                Logger.LogError(ex, "SQS message handling failed for message {MessageId}.", delivery.MessageId);
 2130            AsyncResponseDiagnostics.SetError(activity, ex);
 2131            throw;
 132        }
 3133    }
 134
 135    /// <summary>
 136    /// Best-effort <c>ChangeMessageVisibility</c>: the receipt handle may already be expired or the
 137    /// message deleted by a competing consumer, and either way SQS redelivery still owns the retry.
 138    /// </summary>
 139    protected async ValueTask TryChangeVisibilityAsync(SqsTransportDelivery delivery, TimeSpan delay)
 140    {
 141        try
 142        {
 2143            await delivery.ChangeVisibilityAsync(delay).ConfigureAwait(false);
 2144        }
 2145        catch (Exception ex)
 146        {
 2147            Logger.LogWarning(
 2148                ex,
 2149                "Failed to change visibility of SQS message {MessageId} on {Queue}; it stays invisible until the visibil
 2150                delivery.MessageId,
 2151                _queue);
 2152        }
 2153    }
 154
 155    protected async ValueTask NotifyBackgroundFailureAsync(
 156        SqsTransportDelivery delivery,
 157        Exception exception,
 158        string queue,
 159        SqsSubscriberRole role)
 160    {
 2161        var callback = _subscriberOptions.OnBackgroundFailure;
 2162        if (callback is null)
 0163            return;
 164
 165        try
 166        {
 2167            var context = new SqsBackgroundFailureContext(
 2168                queue,
 2169                role.ToString(),
 2170                delivery.MessageId,
 2171                delivery.ReceiveCount,
 2172                TryReadCorrelationId(delivery),
 2173                exception);
 2174            await callback(context).ConfigureAwait(false);
 2175        }
 2176        catch (Exception callbackException)
 177        {
 2178            Logger.LogError(
 2179                callbackException,
 2180                "SQS background failure callback failed for already-deleted message {MessageId} on {Queue}.",
 2181                delivery.MessageId,
 2182                queue);
 2183        }
 2184    }
 185
 186    private string? TryReadCorrelationId(SqsTransportDelivery delivery)
 3187        => !string.IsNullOrWhiteSpace(_transportOptions.CorrelationIdAttribute)
 3188            && delivery.MessageAttributes.TryGetValue(_transportOptions.CorrelationIdAttribute, out var value)
 3189            && !string.IsNullOrWhiteSpace(value)
 3190                ? value
 3191                : null;
 192}
 193
 194internal sealed class AwaitingSqsMessageDispatcher(
 195    Func<SqsTransportDelivery, CancellationToken, Task> handler,
 196    SqsAsyncResponseOptions transportOptions,
 197    SqsSubscriberOptions subscriberOptions,
 198    ILogger logger,
 199    string queue,
 200    SqsSubscriberRole role)
 201    : SqsMessageDispatcher(handler, transportOptions, subscriberOptions, logger, queue, role)
 202{
 203    /// <summary>Handles the delivered message.</summary>
 204    public override async Task HandleAsync(
 205        SqsTransportDelivery delivery,
 206        CancellationToken subscriberCancellationToken)
 207    {
 208        try
 209        {
 210            await ExecuteHandlerAsync(delivery, subscriberCancellationToken).ConfigureAwait(false);
 211            await delivery.DeleteAsync().ConfigureAwait(false);
 212        }
 213        catch (Exception)
 214        {
 215            // SQS has no explicit NACK or dead-letter call: leaving the message undeleted lets it
 216            // reappear when its visibility timeout expires, ApproximateReceiveCount increments, and
 217            // the queue's redrive policy dead-letters it after maxReceiveCount receives.
 218            if (RedeliveryDelay is { } redeliveryDelay)
 219                await TryChangeVisibilityAsync(delivery, redeliveryDelay).ConfigureAwait(false);
 220        }
 221    }
 222}
 223
 224internal sealed class QueuedSqsMessageDispatcher : SqsMessageDispatcher
 225{
 226    private readonly Channel<SqsTransportDelivery> _queue;
 227    private readonly Task[] _workers;
 228    private readonly CancellationTokenSource _drainCancellation = new();
 229    private readonly TimeSpan _drainTimeout;
 230    private readonly int _capacity;
 231    private readonly string _queueName;
 232    private readonly SqsSubscriberRole _role;
 233    private int _pendingCount;
 234    private int _runningCount;
 235    private int _disposeStarted;
 236
 237    /// <summary>Creates an ACK-after-enqueue dispatcher with a bounded background queue.</summary>
 238    public QueuedSqsMessageDispatcher(
 239        Func<SqsTransportDelivery, CancellationToken, Task> handler,
 240        SqsAsyncResponseOptions transportOptions,
 241        SqsSubscriberOptions subscriberOptions,
 242        ILogger logger,
 243        string queue,
 244        SqsSubscriberRole role)
 245        : base(handler, transportOptions, subscriberOptions, logger, queue, role)
 246    {
 247        _drainTimeout = subscriberOptions.BackgroundDrainTimeout;
 248        _capacity = subscriberOptions.BackgroundQueueCapacity;
 249        _queueName = queue;
 250        _role = role;
 251        _queue = Channel.CreateBounded<SqsTransportDelivery>(new BoundedChannelOptions(subscriberOptions.BackgroundQueue
 252        {
 253            AllowSynchronousContinuations = false,
 254            FullMode = BoundedChannelFullMode.Wait,
 255            SingleReader = subscriberOptions.BackgroundWorkerCount == 1,
 256            SingleWriter = false
 257        });
 258
 259        _workers = Enumerable.Range(0, subscriberOptions.BackgroundWorkerCount)
 260            .Select(workerIndex => Task.Run(() => RunWorkerAsync(workerIndex)))
 261            .ToArray();
 262
 263        Logger.LogInformation(
 264            "Created SQS ACK-after-enqueue dispatcher for {Queue} with {WorkerCount} worker(s), queue capacity {QueueCap
 265            _queueName,
 266            subscriberOptions.BackgroundWorkerCount,
 267            subscriberOptions.BackgroundQueueCapacity,
 268            _drainTimeout);
 269    }
 270
 271    internal int PendingCount => Volatile.Read(ref _pendingCount);
 272    internal int RunningCount => Volatile.Read(ref _runningCount);
 273
 274    public override bool CanAcceptMore => Volatile.Read(ref _pendingCount) < _capacity;
 275
 276    public override int FreeCapacity => Math.Max(0, _capacity - Volatile.Read(ref _pendingCount));
 277
 278    public override async ValueTask WaitForCapacityAsync(CancellationToken cancellationToken)
 279    {
 280        // WaitToWriteAsync completes when the bounded channel has room (or the channel is completed
 281        // during dispose, in which case there is nothing left to gate).
 282        while (!CanAcceptMore)
 283        {
 284            if (!await _queue.Writer.WaitToWriteAsync(cancellationToken).ConfigureAwait(false))
 285                return;
 286        }
 287    }
 288
 289    /// <summary>Handles the delivered message.</summary>
 290    public override async Task HandleAsync(
 291        SqsTransportDelivery delivery,
 292        CancellationToken subscriberCancellationToken)
 293    {
 294        Interlocked.Increment(ref _pendingCount);
 295        if (!_queue.Writer.TryWrite(delivery))
 296        {
 297            Interlocked.Decrement(ref _pendingCount);
 298            Logger.LogWarning(
 299                "SQS background queue rejected message {MessageId} for {Queue}; leaving it to redeliver via its visibili
 300                delivery.MessageId,
 301                _queueName,
 302                PendingCount,
 303                RunningCount);
 304            // Do not release visibility to zero here: SQS counts every receive toward the queue's
 305            // redrive policy, so an instantly re-receivable message that keeps hitting a full queue
 306            // would cross maxReceiveCount and dead-letter without ever being processed. Let the
 307            // visibility timeout lapse naturally (or shorten it via RedeliveryDelay when configured)
 308            // so redelivery lands after capacity has had time to free.
 309            if (RedeliveryDelay is { } redeliveryDelay)
 310                await TryChangeVisibilityAsync(delivery, redeliveryDelay).ConfigureAwait(false);
 311            return;
 312        }
 313
 314        // The delivery now belongs to a background worker, which decrements _pendingCount when it
 315        // dequeues. Do not touch the counter or release visibility here, even if the delete below
 316        // fails — the message is already executing in-process and releasing it would trigger a
 317        // duplicate execution via redelivery.
 318        try
 319        {
 320            await delivery.DeleteAsync().ConfigureAwait(false);
 321        }
 322        catch (Exception ex)
 323        {
 324            Logger.LogError(
 325                ex,
 326                "Failed to delete SQS message {MessageId} for {Queue} after enqueue; it is being processed but SQS will 
 327                delivery.MessageId,
 328                _queueName);
 329        }
 330    }
 331
 332    /// <summary>Releases resources held by this instance.</summary>
 333    public override async ValueTask DisposeAsync()
 334    {
 335        if (Interlocked.Exchange(ref _disposeStarted, 1) != 0)
 336            return;
 337
 338        Logger.LogInformation(
 339            "Draining SQS ACK-after-enqueue dispatcher for {Queue}. Pending={PendingCount}, Running={RunningCount}.",
 340            _queueName,
 341            PendingCount,
 342            RunningCount);
 343        _queue.Writer.TryComplete();
 344
 345        try
 346        {
 347            await Task.WhenAll(_workers).WaitAsync(_drainTimeout).ConfigureAwait(false);
 348            _drainCancellation.Dispose();
 349        }
 350        catch (TimeoutException ex)
 351        {
 352            _drainCancellation.Cancel();
 353            Logger.LogWarning(
 354                ex,
 355                "Timed out while draining SQS ACK-after-enqueue dispatcher for {Queue}. Pending={PendingCount}, Running=
 356                _queueName,
 357                PendingCount,
 358                RunningCount);
 359
 360            _ = Task.WhenAll(_workers).ContinueWith(
 361                _ => _drainCancellation.Dispose(),
 362                CancellationToken.None,
 363                TaskContinuationOptions.ExecuteSynchronously,
 364                TaskScheduler.Default);
 365        }
 366    }
 367
 368    private async Task RunWorkerAsync(int workerIndex)
 369    {
 370        await foreach (var delivery in _queue.Reader.ReadAllAsync().ConfigureAwait(false))
 371        {
 372            Interlocked.Decrement(ref _pendingCount);
 373            Interlocked.Increment(ref _runningCount);
 374
 375            try
 376            {
 377                Logger.LogDebug(
 378                    "SQS background worker {WorkerIndex} handling message {MessageId} for {Queue}. Pending={PendingCount
 379                    workerIndex,
 380                    delivery.MessageId,
 381                    _queueName,
 382                    PendingCount,
 383                    RunningCount);
 384                await ExecuteHandlerAsync(
 385                    delivery,
 386                    _drainCancellation.Token,
 387                    logFailures: false).ConfigureAwait(false);
 388            }
 389            catch (Exception ex)
 390            {
 391                Logger.LogError(
 392                    ex,
 393                    "SQS background handler failed for already-deleted message {MessageId} on {Queue}.",
 394                    delivery.MessageId,
 395                    _queueName);
 396                await NotifyBackgroundFailureAsync(
 397                    delivery,
 398                    ex,
 399                    _queueName,
 400                    _role).ConfigureAwait(false);
 401            }
 402            finally
 403            {
 404                Interlocked.Decrement(ref _runningCount);
 405            }
 406        }
 407    }
 408}