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

Information
Class: Microsoft.Extensions.DependencyInjection.KafkaAsyncResponseTransportServiceCollectionExtensions
Assembly: AsyncResponse.Transports.Kafka
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.Kafka/ServiceCollectionExtensions.cs
Line coverage
100%
Covered lines: 32
Uncovered lines: 0
Coverable lines: 32
Total lines: 71
Line coverage: 100%
Branch coverage
100%
Covered branches: 12
Total branches: 12
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
WithKafkaTransport(...)100%1212100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.Kafka/ServiceCollectionExtensions.cs

#LineLine coverage
 1using AsyncResponse;
 2using AsyncResponse.Transports.Kafka;
 3using Microsoft.Extensions.DependencyInjection.Extensions;
 4using Microsoft.Extensions.Options;
 5
 6namespace Microsoft.Extensions.DependencyInjection;
 7
 8/// <summary>
 9/// DI registration for the Apache Kafka AsyncResponse transport.
 10/// </summary>
 11public static class KafkaAsyncResponseTransportServiceCollectionExtensions
 12{
 13    /// <summary>
 14    /// Registers Apache Kafka as the worker transport and response ingress in one call:
 15    /// <list type="bullet">
 16    /// <item><description>worker jobs are produced to the configured worker topic;</description></item>
 17    /// <item><description>a hosted consumer-group subscriber reads worker messages and executes jobs;</description></it
 18    /// <item><description>a hosted consumer-group subscriber reads response messages and feeds them into <see cref="IAs
 19    /// </list>
 20    /// Kafka is a transport, not the waiter/recovery channel — its partitioned log has no targeted
 21    /// waiter wake and no per-key TTL store. Pair it with <c>.WithRedisChannel()</c>,
 22    /// <c>.WithNatsChannel()</c>, or <c>.WithPostgreSqlChannel()</c> when late responses must
 23    /// survive redeploys, or with <c>.WithInMemoryChannel()</c> for simple single-process waits.
 24    /// <para>
 25    /// The package speaks the Kafka protocol via <c>Confluent.Kafka</c>, so it also works against
 26    /// Redpanda, Amazon MSK, WarpStream, Aiven, and Confluent Cloud.
 27    /// </para>
 28    /// </summary>
 29    public static AsyncResponseRegistrationBuilder WithKafkaTransport(
 30        this AsyncResponseRegistrationBuilder builder,
 31        Action<KafkaAsyncResponseTransportOptions> configure)
 32    {
 333        ArgumentNullException.ThrowIfNull(configure);
 34
 335        var services = builder.Services;
 336        services.AddOptions();
 337        services.Configure(configure);
 38
 339        services.TryAddSingleton<IKafkaProducerClient>(provider => new KafkaProducerClientAdapter(
 340            provider.GetRequiredService<IOptions<KafkaAsyncResponseTransportOptions>>().Value));
 341        services.TryAddSingleton<IKafkaConsumerClientFactory>(provider => new KafkaConsumerClientFactory(
 342            provider.GetRequiredService<IOptions<KafkaAsyncResponseTransportOptions>>().Value));
 343        services.TryAddSingleton<IKafkaAdminClient>(provider => new KafkaAdminClientAdapter(
 344            provider.GetRequiredService<IOptions<KafkaAsyncResponseTransportOptions>>().Value));
 45
 346        services.TryAddSingleton(provider => new KafkaWorkerTransport(
 347            provider.GetRequiredService<IOptions<KafkaAsyncResponseTransportOptions>>(),
 348            provider.GetRequiredService<IKafkaProducerClient>()));
 349        services.Replace(ServiceDescriptor.Singleton<IWorkerTransport>(provider =>
 350            provider.GetRequiredService<KafkaWorkerTransport>()));
 351        services.Replace(ServiceDescriptor.Singleton<IAsyncResponseReplyTargetProvider, KafkaReplyTargetProvider>());
 352        services.AddSingleton(provider =>
 353        {
 354            // Resolved ack modes declared to the Core startup validator, which vetoes early ACK on
 355            // the worker queue durable-flow wake-ups ride (see AsyncResponseStartupValidator).
 356            var options = provider.GetRequiredService<Microsoft.Extensions.Options.IOptions<KafkaAsyncResponseTransportO
 357            return new AsyncResponseTransportMarker("Kafka")
 358            {
 359                WorkerSubscriberUsesEarlyAck = options.WorkerSubscriber.AckMode == KafkaAckMode.AckAfterEnqueue,
 360                WorkerAckModePath = $"{nameof(KafkaAsyncResponseTransportOptions)}.{nameof(options.WorkerSubscriber)}.{n
 361                ResponseSubscriberUsesEarlyAck = options.ResponseSubscriber.AckMode == KafkaAckMode.AckAfterEnqueue,
 362                ResponseAckModePath = $"{nameof(KafkaAsyncResponseTransportOptions)}.{nameof(options.ResponseSubscriber)
 363            };
 364        });
 65
 366        services.AddHostedService<KafkaWorkerSubscriber>();
 367        services.AddHostedService<KafkaResponseIngressSubscriber>();
 68
 369        return builder;
 70    }
 71}