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

Information
Class: AsyncResponse.Transports.NATS.NatsMessageDispatcher
Assembly: AsyncResponse.Transports.NATS
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.NATS/NatsMessageDispatcher.cs
Line coverage
100%
Covered lines: 151
Uncovered lines: 0
Coverable lines: 151
Total lines: 291
Line coverage: 100%
Branch coverage
100%
Covered branches: 34
Total branches: 34
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%44100%
HandleAsync()100%22100%
ExecuteHandlerAsync()100%1414100%
HandleEarlyAckAsync()100%22100%
BackgroundWorkerLoopAsync()100%11100%
HandleFailureAsync()100%66100%
DeadLetterAsync()100%22100%
InvokeBackgroundFailureAsync()100%22100%
SanitizeHeaderValue(...)100%11100%
DisposeAsync()100%22100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.NATS/NatsMessageDispatcher.cs

#LineLine coverage
 1using Microsoft.Extensions.Logging;
 2using System.Diagnostics;
 3using System.Threading.Channels;
 4
 5namespace AsyncResponse.Transports.NATS;
 6
 7internal enum NatsSubscriberRole
 8{
 9    Worker,
 10    ResponseIngress
 11}
 12
 13/// <summary>
 14/// Applies the acknowledgement, redelivery, and dead-letter policy to JetStream deliveries.
 15/// <list type="bullet">
 16/// <item><description><see cref="NatsAckMode.AckAfterHandlerCompletes"/>: run the handler, then ACK;
 17/// on failure NAK for redelivery until <see cref="NatsSubscriberOptions.MaxDeliveryAttempts"/>, then
 18/// dead-letter and terminate.</description></item>
 19/// <item><description><see cref="NatsAckMode.AckAfterEnqueue"/>: enqueue to a bounded background
 20/// queue and ACK immediately; background handler failures are dead-lettered and reported.</description></item>
 21/// </list>
 22/// </summary>
 23internal sealed class NatsMessageDispatcher : IAsyncDisposable
 24{
 25    private readonly Func<NatsJobDelivery, CancellationToken, Task> _handler;
 26    private readonly INatsJetStreamTransport _jetStream;
 27    private readonly NatsAsyncResponseTransportOptions _options;
 28    private readonly NatsSubscriberOptions _subscriberOptions;
 29    private readonly NatsTransportSubjectSchema _schema;
 30    private readonly ILogger _logger;
 31    private readonly NatsSubscriberRole _role;
 32    private readonly string _consumer;
 33
 34    private readonly Channel<NatsJobDelivery>? _backgroundQueue;
 35    private readonly Task[]? _backgroundWorkers;
 36    private readonly CancellationTokenSource? _backgroundCts;
 37
 38    /// <summary>Runs the NatsMessageDispatcher operation.</summary>
 339    public NatsMessageDispatcher(
 340        Func<NatsJobDelivery, CancellationToken, Task> handler,
 341        INatsJetStreamTransport jetStream,
 342        NatsAsyncResponseTransportOptions options,
 343        NatsSubscriberOptions subscriberOptions,
 344        NatsTransportSubjectSchema schema,
 345        ILogger logger,
 346        NatsSubscriberRole role,
 347        string consumer)
 48    {
 349        NatsTransportOptionsValidator.ValidateSubscriber(options, subscriberOptions, role.ToString());
 50
 351        _handler = handler;
 352        _jetStream = jetStream;
 353        _options = options;
 354        _subscriberOptions = subscriberOptions;
 355        _schema = schema;
 356        _logger = logger;
 357        _role = role;
 358        _consumer = consumer;
 59
 360        if (subscriberOptions.AckMode is NatsAckMode.AckAfterEnqueue)
 61        {
 362            _backgroundQueue = Channel.CreateBounded<NatsJobDelivery>(new BoundedChannelOptions(subscriberOptions.Backgr
 363            {
 364                SingleReader = false,
 365                SingleWriter = true,
 366                FullMode = BoundedChannelFullMode.Wait
 367            });
 368            _backgroundCts = new CancellationTokenSource();
 369            _backgroundWorkers = new Task[subscriberOptions.BackgroundWorkerCount];
 370            for (var i = 0; i < _backgroundWorkers.Length; i++)
 371                _backgroundWorkers[i] = Task.Run(() => BackgroundWorkerLoopAsync(_backgroundCts.Token));
 72        }
 373    }
 74
 75    /// <summary>Handles the delivered message.</summary>
 76    public async Task HandleAsync(NatsJobDelivery delivery, CancellationToken cancellationToken)
 77    {
 378        if (_subscriberOptions.AckMode is NatsAckMode.AckAfterEnqueue)
 79        {
 380            await HandleEarlyAckAsync(delivery, cancellationToken).ConfigureAwait(false);
 381            return;
 82        }
 83
 84        try
 85        {
 386            await ExecuteHandlerAsync(delivery, cancellationToken).ConfigureAwait(false);
 387            await delivery.AckAsync().ConfigureAwait(false);
 388        }
 389        catch (Exception ex)
 90        {
 291            await HandleFailureAsync(delivery, ex, cancellationToken).ConfigureAwait(false);
 92        }
 393    }
 94
 95    // Single choke point for handler execution so both ACK modes emit the consumer receive span.
 96    private async Task ExecuteHandlerAsync(NatsJobDelivery delivery, CancellationToken cancellationToken)
 97    {
 398        using var activity = AsyncResponseDiagnostics.StartActivity(
 399            "asyncresponse.nats.receive",
 3100            ActivityKind.Consumer);
 3101        activity?.SetTag("asyncresponse.transport", "nats");
 3102        activity?.SetTag("asyncresponse.nats.role", _role.ToString());
 3103        activity?.SetTag("asyncresponse.nats.ack_mode", _subscriberOptions.AckMode.ToString());
 3104        activity?.SetTag("messaging.system", "nats");
 3105        activity?.SetTag("messaging.destination.name", delivery.Subject);
 3106        activity?.SetTag("messaging.nats.num_delivered", delivery.NumDelivered);
 107
 3108        if (delivery.Headers.TryGetValue(_options.CorrelationIdHeader, out var correlationId))
 3109            AsyncResponseDiagnostics.SetCorrelationId(activity, correlationId);
 110
 111        try
 112        {
 3113            await _handler(delivery, cancellationToken).ConfigureAwait(false);
 3114        }
 2115        catch (Exception ex)
 116        {
 2117            AsyncResponseDiagnostics.SetError(activity, ex);
 3118            throw;
 119        }
 3120    }
 121
 122    private async Task HandleEarlyAckAsync(NatsJobDelivery delivery, CancellationToken cancellationToken)
 123    {
 124        // Accept into the background queue and ACK. If the queue is saturated, wait for a worker to
 125        // free a slot instead of NAKing: the wait blocks the consume loop, so the subscriber stops
 126        // pulling new messages until capacity frees rather than churning NAK/redeliver cycles.
 3127        if (_backgroundQueue!.Writer.TryWrite(delivery))
 128        {
 3129            await delivery.AckAsync().ConfigureAwait(false);
 3130            return;
 131        }
 132
 133        try
 134        {
 2135            _logger.LogDebug("Background queue full for {Role}; pausing the consume loop until capacity frees.", _role);
 2136            await _backgroundQueue.Writer.WriteAsync(delivery, cancellationToken).ConfigureAwait(false);
 2137            await delivery.AckAsync().ConfigureAwait(false);
 2138        }
 2139        catch (Exception ex) when (ex is OperationCanceledException or ChannelClosedException)
 140        {
 141            // Subscriber stopping or dispatcher disposing: NAK so JetStream redelivers elsewhere.
 2142            _logger.LogDebug("Background queue unavailable for {Role} during shutdown; NAKing message for redelivery.", 
 2143            await delivery.NakAsync(_subscriberOptions.RedeliveryDelay).ConfigureAwait(false);
 144        }
 3145    }
 146
 147    private async Task BackgroundWorkerLoopAsync(CancellationToken cancellationToken)
 148    {
 149        // Token-less ReadAllAsync: on shutdown the queue is completed and fully drained, so every
 150        // already-ACKed delivery is attempted (with the drain token once the drain budget lapses)
 151        // instead of being silently dropped; each failure is dead-lettered and surfaced below.
 3152        await foreach (var delivery in _backgroundQueue!.Reader.ReadAllAsync().ConfigureAwait(false))
 153        {
 154            try
 155            {
 3156                await ExecuteHandlerAsync(delivery, cancellationToken).ConfigureAwait(false);
 3157            }
 3158            catch (Exception ex)
 159            {
 2160                _logger.LogError(ex, "Background handler failed for {Role} on subject {Subject} after early ACK.", _role
 2161                await DeadLetterAsync(delivery, ex, CancellationToken.None).ConfigureAwait(false);
 2162                await InvokeBackgroundFailureAsync(delivery, ex).ConfigureAwait(false);
 3163            }
 3164        }
 3165    }
 166
 167    private async Task HandleFailureAsync(NatsJobDelivery delivery, Exception exception, CancellationToken cancellationT
 168    {
 2169        var maxAttempts = _subscriberOptions.MaxDeliveryAttempts;
 2170        if (maxAttempts > 0 && delivery.NumDelivered >= maxAttempts)
 171        {
 2172            _logger.LogError(
 2173                exception,
 2174                "Message on subject {Subject} ({Role}) failed after {Attempts} attempts; dead-lettering.",
 2175                delivery.Subject,
 2176                _role,
 2177                delivery.NumDelivered);
 178
 2179            var shouldTerminate = await DeadLetterAsync(delivery, exception, cancellationToken).ConfigureAwait(false);
 2180            if (shouldTerminate)
 181            {
 2182                await delivery.TermAsync().ConfigureAwait(false);
 183            }
 184            else
 185            {
 2186                _logger.LogWarning(
 2187                    exception,
 2188                    "Dead-letter publish failed for subject {Subject} ({Role}); NAKing so the message can be retried.",
 2189                    delivery.Subject,
 2190                    _role);
 2191                await delivery.NakAsync(_subscriberOptions.RedeliveryDelay).ConfigureAwait(false);
 192            }
 193        }
 194        else
 195        {
 2196            _logger.LogWarning(
 2197                exception,
 2198                "Message on subject {Subject} ({Role}) failed on attempt {Attempt}; NAKing for redelivery.",
 2199                delivery.Subject,
 2200                _role,
 2201                delivery.NumDelivered);
 2202            await delivery.NakAsync(_subscriberOptions.RedeliveryDelay).ConfigureAwait(false);
 203        }
 2204    }
 205
 206    private async Task<bool> DeadLetterAsync(NatsJobDelivery delivery, Exception exception, CancellationToken cancellati
 207    {
 2208        if (!_options.DeadLetterEnabled)
 209        {
 2210            _logger.LogError(
 2211                exception,
 2212                "Message on subject {Subject} ({Role}) is unprocessable and dead-lettering is disabled; it will be dropp
 2213                delivery.Subject,
 2214                _role);
 2215            return true;
 216        }
 217
 2218        var headers = new Dictionary<string, string>(delivery.Headers, StringComparer.OrdinalIgnoreCase)
 2219        {
 2220            ["AR-DeadLetter-Reason"] = SanitizeHeaderValue(exception.Message),
 2221            ["AR-DeadLetter-Source-Subject"] = delivery.Subject,
 2222            ["AR-DeadLetter-Role"] = _role.ToString()
 2223        };
 224
 225        try
 226        {
 2227            await _jetStream.PublishAsync(_schema.DeadLetterSubject, delivery.Payload, headers, cancellationToken).Confi
 2228            _logger.LogInformation("Dead-lettered message from subject {Subject} ({Role}) to {DeadLetterSubject}.", deli
 2229            return true;
 230        }
 2231        catch (Exception ex)
 232        {
 2233            _logger.LogError(ex, "Failed to dead-letter message from subject {Subject} ({Role}).", delivery.Subject, _ro
 2234            return false;
 235        }
 2236    }
 237
 238    private async Task InvokeBackgroundFailureAsync(NatsJobDelivery delivery, Exception exception)
 239    {
 2240        if (_subscriberOptions.OnBackgroundFailure is null)
 2241            return;
 242
 243        try
 244        {
 2245            delivery.Headers.TryGetValue(_options.CorrelationIdHeader, out var correlationId);
 2246            var context = new NatsBackgroundFailureContext(delivery.Subject, _consumer, _role.ToString(), delivery.NumDe
 2247            await _subscriberOptions.OnBackgroundFailure(context).ConfigureAwait(false);
 2248        }
 2249        catch (Exception ex)
 250        {
 2251            _logger.LogError(ex, "OnBackgroundFailure callback threw for {Role}.", _role);
 2252        }
 2253    }
 254
 255    private static string SanitizeHeaderValue(string value)
 3256        => value.Replace('\r', ' ').Replace('\n', ' ');
 257
 258    /// <summary>Releases resources held by this instance.</summary>
 259    public async ValueTask DisposeAsync()
 260    {
 3261        if (_backgroundQueue is null)
 3262            return;
 263
 3264        _backgroundQueue.Writer.TryComplete();
 265        try
 266        {
 3267            await Task.WhenAll(_backgroundWorkers!).WaitAsync(_subscriberOptions.BackgroundDrainTimeout).ConfigureAwait(
 3268            _backgroundCts!.Dispose();
 3269        }
 270        catch (TimeoutException)
 271        {
 2272            _logger.LogWarning("Background handlers for {Role} did not drain within {Timeout}.", _role, _subscriberOptio
 2273            await _backgroundCts!.CancelAsync().ConfigureAwait(false);
 274
 275            // The workers are still running and observe _backgroundCts.Token inside ReadAllAsync, so disposing
 276            // it now would throw ObjectDisposedException inside them. Dispose once they actually finish, off
 277            // the shutdown path, so the source is not leaked either.
 3278            _ = Task.WhenAll(_backgroundWorkers!).ContinueWith(
 3279                _ => _backgroundCts.Dispose(),
 3280                CancellationToken.None,
 3281                TaskContinuationOptions.ExecuteSynchronously,
 3282                TaskScheduler.Default);
 283        }
 2284        catch (Exception ex)
 285        {
 286            // WhenAll only completes once every worker has finished, so the source is safe to dispose here.
 2287            _logger.LogDebug(ex, "Background worker drain for {Role} ended with an error.", _role);
 2288            _backgroundCts!.Dispose();
 3289        }
 3290    }
 291}