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

Information
Class: Microsoft.Extensions.DependencyInjection.NatsAsyncResponseChannelServiceCollectionExtensions
Assembly: AsyncResponse.Channels.NATS
File(s): /_/src/Channels/AsyncResponse.Channels.NATS/ServiceCollectionExtensions.cs
Line coverage
100%
Covered lines: 36
Uncovered lines: 0
Coverable lines: 36
Total lines: 91
Line coverage: 100%
Branch coverage
75%
Covered branches: 3
Total branches: 4
Branch coverage: 75%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
WithNatsChannel(...)75%44100%

File(s)

/_/src/Channels/AsyncResponse.Channels.NATS/ServiceCollectionExtensions.cs

#LineLine coverage
 1using AsyncResponse;
 2using AsyncResponse.Channels.NATS;
 3using Microsoft.Extensions.DependencyInjection.Extensions;
 4using Microsoft.Extensions.Options;
 5using NATS.Client.Core;
 6using NATS.Net;
 7
 8namespace Microsoft.Extensions.DependencyInjection;
 9
 10/// <summary>
 11/// DI registration for the NATS-backed AsyncResponse channel package.
 12/// </summary>
 13public static class NatsAsyncResponseChannelServiceCollectionExtensions
 14{
 15    /// <summary>
 16    /// Registers NATS Core request/reply as the response channel — behind
 17    /// <see cref="IAsyncResponsePublisher"/>, <see cref="IAsyncResponseSubscriber"/>, and
 18    /// <see cref="IRecoverableAsyncResponseSubscriber"/> — together with the durable NATS JetStream
 19    /// Key-Value recovery-state store, so a response that arrives after the waiter died (e.g. a
 20    /// redeploy) can still resume or fail the owning flow. It also registers
 21    /// <see cref="IRecoverableAsyncResponseBuilder"/> for fluent durable recovery flows. The same
 22    /// channel instance serves the engine's recovery-state scanner and active-subscriber probe, so the
 23    /// watchdog works against NATS automatically.
 24    /// <para>
 25    /// Requires a <c>NATS.Client.Core.INatsConnection</c> singleton registered by the host — for
 26    /// example via <c>builder.Services.AddNatsClient(...)</c> from
 27    /// <c>NATS.Extensions.Microsoft.DependencyInjection</c>. Reuse the application's existing
 28    /// connection; don't create a second one. The recovery store requires JetStream to be enabled on
 29    /// the NATS server. Publisher/subscriber resolve to one shared instance: per-interface instances
 30    /// would split the channel's internal state.
 31    /// </para>
 32    /// </summary>
 33    public static AsyncResponseRegistrationBuilder WithNatsChannel(
 34        this AsyncResponseRegistrationBuilder builder,
 35        Action<NatsAsyncResponseChannelOptions>? configure = null)
 36    {
 37337        var services = builder.Services;
 37338        services.AddOptions();
 37339        if (configure is not null)
 36940            services.Configure(configure);
 41
 42        // NATS JetStream Key-Value client adapter (lazily creates the recovery bucket on first use).
 74243        services.TryAddSingleton<INatsKvStore>(provider => new NatsKvStoreAdapter(
 74244            provider.GetRequiredService<INatsConnection>().CreateKeyValueStoreContext(),
 74245            provider.GetRequiredService<IOptions<NatsAsyncResponseChannelOptions>>().Value));
 46
 47        // Durable recovery store; also the watchdog's recovery-state scanner.
 37348        services.TryAddSingleton<NatsRecoveryStateStore>();
 74249        services.Replace(ServiceDescriptor.Singleton<IRecoveryStateStore>(provider => provider.GetRequiredService<NatsRe
 70650        services.Replace(ServiceDescriptor.Singleton<IRecoveryStateScanner>(provider => provider.GetRequiredService<Nats
 51
 52        // NATS Core request/reply client adapter (raw extension calls isolated in NatsRawRequester).
 74253        services.TryAddSingleton<INatsResponseChannelClient>(provider => new NatsResponseChannelClient(
 74254            new NatsRawRequester(provider.GetRequiredService<INatsConnection>())));
 55
 56        // NATS response channel: publisher + subscriber + the watchdog's liveness probe, all one shared instance.
 37357        services.TryAddSingleton<NatsAsyncResponseChannel>();
 73158        services.Replace(ServiceDescriptor.Singleton<IAsyncResponsePublisher>(provider => provider.GetRequiredService<Na
 68559        services.Replace(ServiceDescriptor.Singleton<IRawAsyncResponsePublisher>(provider => provider.GetRequiredService
 73560        services.Replace(ServiceDescriptor.Singleton<IAsyncResponseSubscriber>(provider => provider.GetRequiredService<N
 70661        services.Replace(ServiceDescriptor.Singleton<IRecoverableAsyncResponseSubscriber>(provider => provider.GetRequir
 71062        services.Replace(ServiceDescriptor.Singleton<IActiveSubscriberProbe>(provider => provider.GetRequiredService<Nat
 63
 64        // Durable recovery capability: expose the recoverable fluent builder only when the NATS channel
 65        // package is the selected response channel. Plain IAsyncResponseBuilder still works and shares
 66        // the same implementation, but its static type does not offer recovery callbacks.
 70667        services.Replace(ServiceDescriptor.Singleton<IRecoverableAsyncResponseBuilder>(provider => new RecoverableAsyncR
 70668            provider.GetRequiredService<IRecoverableAsyncResponseSubscriber>(),
 70669            provider.GetService<IWorkerTransport>(),
 70670            provider.GetService<IAsyncResponseReplyTargetProvider>(),
 70671            provider.GetRequiredService<AsyncResponseContextPropagation>(),
 70672            provider.GetService<TimeProvider>(),
 70673            // The producer-side mirror of the ingress's inbound size budget (WorkerJobTooLargeException).
 70674            provider.GetService<IOptions<AsyncResponseOptions>>())));
 70675        services.Replace(ServiceDescriptor.Singleton<IAsyncResponseBuilder>(provider => provider.GetRequiredService<IRec
 76
 77        // The resolved default waiter timeout is declared through the marker so the startup
 78        // validator can require the durable-flow ledger TTL to out-live a timeout-less awaited
 79        // step, and the flow engine can extend a parked ledger by it — without either referencing
 80        // channel option types.
 37381        services.AddSingleton(provider =>
 37382        {
 33583            var options = provider.GetRequiredService<IOptions<NatsAsyncResponseChannelOptions>>().Value;
 33584            return new AsyncResponseChannelMarker(NatsAsyncResponseChannelOptions.ChannelName)
 33585            {
 33586                EffectiveDefaultWaitTimeout = options.DefaultTimeout ?? options.RecoveryStateExpiry
 33587            };
 37388        });
 37389        return builder;
 90    }
 91}