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

Information
Class: AsyncResponse.Transports.NATS.NatsResponseIngressSubscriber
Assembly: AsyncResponse.Transports.NATS
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.NATS/NatsSubscriberServices.cs
Line coverage
100%
Covered lines: 11
Uncovered lines: 0
Coverable lines: 11
Total lines: 179
Line coverage: 100%
Branch coverage
N/A
Covered branches: 0
Total branches: 0
Branch coverage: N/A
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
.ctor(...)100%11100%
get_Subject()100%11100%
get_Stream()100%11100%
get_Consumer()100%11100%
get_SubscriberOptions()100%11100%
get_Role()100%11100%
HandleMessageAsync(...)100%11100%

File(s)

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

#LineLine coverage
 1using Microsoft.Extensions.Hosting;
 2using Microsoft.Extensions.Logging;
 3using Microsoft.Extensions.Options;
 4using NATS.Client.Core;
 5using NATS.Net;
 6
 7namespace AsyncResponse.Transports.NATS;
 8
 9/// <summary>
 10/// Base hosted service that consumes a JetStream subject through a durable consumer and routes each
 11/// message to the AsyncResponse ingress with the configured acknowledgement/redelivery/dead-letter
 12/// policy. A failed consume loop is retried with bounded backoff so a transient NATS outage does not
 13/// kill the subscriber.
 14/// </summary>
 15internal abstract class NatsSubscriberService : BackgroundService
 16{
 17    private readonly INatsJetStreamTransport _jetStream;
 18
 19    /// <summary>Runs the NatsSubscriberService operation.</summary>
 20    protected NatsSubscriberService(
 21        IOptions<NatsAsyncResponseTransportOptions> options,
 22        INatsConnection connection,
 23        ILogger logger)
 24        : this(options, new NatsJetStreamTransportAdapter(connection.CreateJetStreamContext()), logger)
 25    {
 26    }
 27
 28    /// <summary>Runs the NatsSubscriberService operation.</summary>
 29    protected NatsSubscriberService(
 30        IOptions<NatsAsyncResponseTransportOptions> options,
 31        INatsJetStreamTransport jetStream,
 32        ILogger logger)
 33    {
 34        Options = options.Value;
 35        NatsTransportOptionsValidator.ValidateCommon(Options);
 36        _jetStream = jetStream;
 37        Logger = logger;
 38        Schema = new NatsTransportSubjectSchema(Options);
 39    }
 40
 41    protected NatsAsyncResponseTransportOptions Options { get; }
 42    protected ILogger Logger { get; }
 43    protected NatsTransportSubjectSchema Schema { get; }
 44
 45    protected abstract string Subject { get; }
 46    protected abstract string Stream { get; }
 47    protected abstract string Consumer { get; }
 48    protected abstract NatsSubscriberOptions SubscriberOptions { get; }
 49    protected abstract NatsSubscriberRole Role { get; }
 50    /// <summary>Handles the delivered message.</summary>
 51    protected abstract Task HandleMessageAsync(NatsJobDelivery delivery, CancellationToken cancellationToken);
 52
 53    /// <summary>Runs this background operation until cancellation is requested.</summary>
 54    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
 55    {
 56        NatsTransportOptionsValidator.ValidateSubscriber(Options, SubscriberOptions, Role.ToString());
 57
 58        var failures = 0;
 59        while (!stoppingToken.IsCancellationRequested)
 60        {
 61            try
 62            {
 63                await RunSubscriberAsync(stoppingToken).ConfigureAwait(false);
 64                return;
 65            }
 66            catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 67            {
 68                return;
 69            }
 70            catch (Exception ex) when (!stoppingToken.IsCancellationRequested)
 71            {
 72                failures++;
 73                var retryDelay = AsyncResponseRetry.Backoff(failures, Options.SubscriberRetryBaseDelay, Options.Subscrib
 74                Logger.LogWarning(ex, "NATS subscriber failed for subject {Subject} ({Role}); retrying in {RetryDelay}."
 75                await Task.Delay(retryDelay, stoppingToken).ConfigureAwait(false);
 76            }
 77        }
 78    }
 79
 80    private async Task RunSubscriberAsync(CancellationToken stoppingToken)
 81    {
 82        if (Options.CreateStreams)
 83        {
 84            await _jetStream.EnsureStreamAsync(Stream, Subject, Options.StreamMaxMessages, stoppingToken).ConfigureAwait
 85            if (Options.DeadLetterEnabled)
 86                await _jetStream.EnsureStreamAsync(Schema.DeadLetterStream, Schema.DeadLetterSubject, Options.DeadLetter
 87        }
 88
 89        await _jetStream.EnsureConsumerAsync(Stream, Consumer, Options.AckWait, stoppingToken).ConfigureAwait(false);
 90
 91        await using var dispatcher = new NatsMessageDispatcher(
 92            HandleMessageAsync,
 93            _jetStream,
 94            Options,
 95            SubscriberOptions,
 96            Schema,
 97            Logger,
 98            Role,
 99            Consumer);
 100
 101        Logger.LogInformation(
 102            "NATS subscriber started. Subject: {Subject}. Stream: {Stream}. Consumer: {Consumer}. Role: {Role}. AckMode:
 103            Subject, Stream, Consumer, Role, SubscriberOptions.AckMode);
 104
 105        await foreach (var delivery in _jetStream.ConsumeAsync(Stream, Consumer, SubscriberOptions.BatchSize, stoppingTo
 106        {
 107            await dispatcher.HandleAsync(delivery, stoppingToken).ConfigureAwait(false);
 108        }
 109    }
 110}
 111
 112/// <summary>Consumes worker-job messages and executes them through the AsyncResponse ingress.</summary>
 113internal sealed class NatsWorkerSubscriber : NatsSubscriberService
 114{
 115    private readonly IAsyncResponseIngress _ingress;
 116
 117    /// <summary>Runs the NatsWorkerSubscriber operation.</summary>
 118    public NatsWorkerSubscriber(
 119        IOptions<NatsAsyncResponseTransportOptions> options,
 120        INatsConnection connection,
 121        IAsyncResponseIngress ingress,
 122        ILogger<NatsWorkerSubscriber> logger)
 123        : base(options, connection, logger)
 124        => _ingress = ingress;
 125
 126    internal NatsWorkerSubscriber(
 127        IOptions<NatsAsyncResponseTransportOptions> options,
 128        INatsJetStreamTransport jetStream,
 129        IAsyncResponseIngress ingress,
 130        ILogger<NatsWorkerSubscriber> logger)
 131        : base(options, jetStream, logger)
 132        => _ingress = ingress;
 133
 134    protected override string Subject => Schema.WorkerSubject;
 135    protected override string Stream => Schema.WorkerStream;
 136    protected override string Consumer => Options.WorkerConsumer;
 137    protected override NatsSubscriberOptions SubscriberOptions => Options.WorkerSubscriber;
 138    protected override NatsSubscriberRole Role => NatsSubscriberRole.Worker;
 139
 140    /// <summary>Handles the delivered message.</summary>
 141    protected override Task HandleMessageAsync(NatsJobDelivery delivery, CancellationToken cancellationToken)
 142        => _ingress.HandleWorkerMessageAsync(delivery.Payload);
 143}
 144
 145/// <summary>Consumes response messages and feeds them into the AsyncResponse ingress, correlated by header or JSON body
 146internal sealed class NatsResponseIngressSubscriber : NatsSubscriberService
 147{
 148    private readonly IAsyncResponseIngress _ingress;
 149
 150    /// <summary>Runs the NatsResponseIngressSubscriber operation.</summary>
 151    public NatsResponseIngressSubscriber(
 152        IOptions<NatsAsyncResponseTransportOptions> options,
 153        INatsConnection connection,
 154        IAsyncResponseIngress ingress,
 155        ILogger<NatsResponseIngressSubscriber> logger)
 3156        : base(options, connection, logger)
 3157        => _ingress = ingress;
 158
 159    internal NatsResponseIngressSubscriber(
 160        IOptions<NatsAsyncResponseTransportOptions> options,
 161        INatsJetStreamTransport jetStream,
 162        IAsyncResponseIngress ingress,
 163        ILogger<NatsResponseIngressSubscriber> logger)
 2164        : base(options, jetStream, logger)
 3165        => _ingress = ingress;
 166
 3167    protected override string Subject => Schema.ResponseSubject;
 3168    protected override string Stream => Schema.ResponseStream;
 3169    protected override string Consumer => Options.ResponseConsumer;
 3170    protected override NatsSubscriberOptions SubscriberOptions => Options.ResponseSubscriber;
 3171    protected override NatsSubscriberRole Role => NatsSubscriberRole.ResponseIngress;
 172
 173    /// <summary>Handles the delivered message.</summary>
 174    protected override Task HandleMessageAsync(NatsJobDelivery delivery, CancellationToken cancellationToken)
 175    {
 3176        var correlationId = NatsCorrelationIdExtractor.Extract(delivery.Headers, delivery.Payload, Options);
 3177        return _ingress.HandleResponseMessageAsync(delivery.Payload, correlationId);
 178    }
 179}