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

Information
Class: AsyncResponse.Transports.NATS.NatsWorkerTransport
Assembly: AsyncResponse.Transports.NATS
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.NATS/NatsWorkerTransport.cs
Line coverage
100%
Covered lines: 55
Uncovered lines: 0
Coverable lines: 55
Total lines: 111
Line coverage: 100%
Branch coverage
75%
Covered branches: 12
Total branches: 16
Branch coverage: 75%
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%
PublishAsync()66.67%1212100%
EnsureWorkerStreamOnceAsync()100%44100%

File(s)

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

#LineLine coverage
 1using Microsoft.Extensions.Options;
 2using NATS.Client.Core;
 3using NATS.Net;
 4using System.Diagnostics;
 5using System.Text.Json;
 6
 7namespace AsyncResponse.Transports.NATS;
 8
 9/// <summary>
 10/// Publishes <see cref="WorkerJobEnvelope"/> messages to a NATS JetStream subject.
 11/// </summary>
 12/// <remarks>
 13/// JetStream provides durable queueing and explicit acknowledgement. The correlation id travels as a
 14/// message header so the consuming side can correlate without parsing the body, and transient NATS
 15/// failures are retried with bounded exponential backoff before the exception is returned to the
 16/// caller.
 17/// </remarks>
 18public sealed class NatsWorkerTransport : IWorkerTransport
 19{
 20    private readonly NatsAsyncResponseTransportOptions _options;
 21    private readonly INatsJetStreamTransport _jetStream;
 22    private readonly NatsTransportSubjectSchema _schema;
 323    private readonly SemaphoreSlim _ensureStreamGate = new(1, 1);
 24    private bool _streamEnsured;
 25
 26    /// <summary>Runs the NatsWorkerTransport operation.</summary>
 27    public NatsWorkerTransport(
 28        IOptions<NatsAsyncResponseTransportOptions> options,
 29        INatsConnection connection)
 330        : this(options, new NatsJetStreamTransportAdapter(connection.CreateJetStreamContext()))
 31    {
 332    }
 33
 334    internal NatsWorkerTransport(
 335        IOptions<NatsAsyncResponseTransportOptions> options,
 336        INatsJetStreamTransport jetStream)
 37    {
 338        _options = options.Value;
 339        NatsTransportOptionsValidator.ValidateCommon(_options);
 340        _jetStream = jetStream;
 341        _schema = new NatsTransportSubjectSchema(_options);
 342    }
 43
 44    /// <summary>Publishes the supplied message.</summary>
 45    public async Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 46    {
 347        ArgumentNullException.ThrowIfNull(job);
 48
 349        using var activity = AsyncResponseDiagnostics.StartActivity(
 350            "asyncresponse.worker.publish",
 351            ActivityKind.Producer,
 352            job.CorrelationId);
 353        activity?.SetTag("asyncresponse.transport", "nats");
 354        activity?.SetTag("messaging.system", "nats");
 355        activity?.SetTag("messaging.destination.name", _schema.WorkerSubject);
 356        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 357        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 58
 59        try
 60        {
 361            if (_options.CreateStreams)
 362                await EnsureWorkerStreamOnceAsync(cancellationToken).ConfigureAwait(false);
 63
 364            var headers = string.IsNullOrWhiteSpace(job.CorrelationId)
 365                ? null
 366                : new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase)
 367                {
 368                    [_options.CorrelationIdHeader] = job.CorrelationId!
 369                };
 70
 371            var payload = AsyncResponseJson.Serialize(job);
 372            var sequence = await NatsTransportRetry.ExecuteAsync(
 373                token => _jetStream.PublishAsync(_schema.WorkerSubject, payload, headers, token),
 374                _options.PublishMaxAttempts,
 375                _options.PublishRetryBaseDelay,
 376                _options.PublishRetryMaxDelay,
 377                cancellationToken).ConfigureAwait(false);
 78
 379            activity?.SetTag("messaging.message.id", sequence);
 380        }
 281        catch (Exception ex)
 82        {
 283            AsyncResponseDiagnostics.SetError(activity, ex);
 384            throw;
 85        }
 386    }
 87
 88    private async Task EnsureWorkerStreamOnceAsync(CancellationToken cancellationToken)
 89    {
 390        if (_streamEnsured)
 391            return;
 92
 393        await _ensureStreamGate.WaitAsync(cancellationToken).ConfigureAwait(false);
 94        try
 95        {
 396            if (!_streamEnsured)
 97            {
 398                await _jetStream.EnsureStreamAsync(
 399                    _schema.WorkerStream,
 3100                    _schema.WorkerSubject,
 3101                    _options.StreamMaxMessages,
 3102                    cancellationToken).ConfigureAwait(false);
 3103                _streamEnsured = true;
 104            }
 3105        }
 106        finally
 107        {
 3108            _ensureStreamGate.Release();
 109        }
 3110    }
 111}