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

Information
Class: AsyncResponse.Transports.Kafka.KafkaWorkerTransport
Assembly: AsyncResponse.Transports.Kafka
File(s): /_/src/Transports/AsyncResponse.Transports.Kafka/KafkaWorkerTransport.cs
Line coverage
100%
Covered lines: 48
Uncovered lines: 0
Coverable lines: 48
Total lines: 99
Line coverage: 100%
Branch coverage
100%
Covered branches: 14
Total branches: 14
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%11100%
.ctor(...)100%11100%
KafkaTransportOptionsValidatorWithValue(...)100%11100%
PublishAsync()100%1010100%
CreateMessageHeaders(...)100%22100%

File(s)

/_/src/Transports/AsyncResponse.Transports.Kafka/KafkaWorkerTransport.cs

#LineLine coverage
 1using Microsoft.Extensions.Options;
 2using System.Diagnostics;
 3using System.Text;
 4using System.Text.Json;
 5
 6namespace AsyncResponse.Transports.Kafka;
 7
 8/// <summary>
 9/// Publishes <see cref="WorkerJobEnvelope"/> messages to a Kafka topic.
 10/// </summary>
 11/// <remarks>
 12/// The correlation id is carried in a message header and doubles as the partition key, so all jobs
 13/// of one flow stay ordered within their partition; jobs without a correlation id are spread
 14/// round-robin. The producer is idempotent with <c>acks=all</c>, and transient broker failures are
 15/// retried with bounded exponential backoff before the exception is returned to the caller.
 16/// </remarks>
 17public sealed class KafkaWorkerTransport : IWorkerTransport
 18{
 19    private readonly KafkaAsyncResponseTransportOptions _options;
 20    private readonly IKafkaProducerClient _producer;
 21    private readonly KafkaTransportTopicSchema _topics;
 22
 23    /// <summary>Runs the KafkaWorkerTransport operation.</summary>
 24    public KafkaWorkerTransport(IOptions<KafkaAsyncResponseTransportOptions> options)
 225        : this(options, new KafkaProducerClientAdapter(KafkaTransportOptionsValidatorWithValue(options)))
 26    {
 227    }
 28
 21629    internal KafkaWorkerTransport(
 21630        IOptions<KafkaAsyncResponseTransportOptions> options,
 21631        IKafkaProducerClient producer)
 32    {
 21633        _options = options.Value;
 21634        KafkaTransportOptionsValidator.ValidateCommon(_options);
 21435        _producer = producer;
 21436        _topics = new KafkaTransportTopicSchema(_options);
 21437    }
 38
 39    private static KafkaAsyncResponseTransportOptions KafkaTransportOptionsValidatorWithValue(
 40        IOptions<KafkaAsyncResponseTransportOptions> options)
 41    {
 242        KafkaTransportOptionsValidator.ValidateCommon(options.Value);
 243        return options.Value;
 44    }
 45
 46    /// <summary>Publishes the supplied message.</summary>
 47    public async Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 48    {
 42549        ArgumentNullException.ThrowIfNull(job);
 50
 42351        using var activity = AsyncResponseDiagnostics.StartActivity(
 42352            "asyncresponse.worker.publish",
 42353            ActivityKind.Producer,
 42354            job.CorrelationId);
 42355        activity?.SetTag("asyncresponse.transport", "kafka");
 42356        activity?.SetTag("messaging.system", "kafka");
 42357        activity?.SetTag("messaging.destination.name", _topics.WorkerTopic);
 42358        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 42359        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 60
 61        try
 62        {
 42363            var payload = Encoding.UTF8.GetBytes(AsyncResponseJson.Serialize(job));
 42364            var headers = CreateMessageHeaders(job.CorrelationId, _options);
 42365            var result = await KafkaTransportRetry.ExecuteAsync(
 43166                token => _producer.PublishAsync(
 43167                    _topics.WorkerTopic,
 43168                    string.IsNullOrWhiteSpace(job.CorrelationId) ? null : job.CorrelationId,
 43169                    payload,
 43170                    headers,
 43171                    token),
 42372                _options.PublishMaxAttempts,
 42373                _options.PublishRetryBaseDelay,
 42374                _options.PublishRetryMaxDelay,
 42375                cancellationToken).ConfigureAwait(false);
 76
 41977            activity?.SetTag("messaging.kafka.destination.partition", result.Partition);
 41978            activity?.SetTag("messaging.kafka.message.offset", result.Offset);
 41979        }
 480        catch (Exception ex)
 81        {
 482            AsyncResponseDiagnostics.SetError(activity, ex);
 483            throw;
 84        }
 41985    }
 86
 87    internal static IReadOnlyList<KafkaTransportHeader> CreateMessageHeaders(
 88        string? correlationId,
 89        KafkaAsyncResponseTransportOptions options)
 90    {
 42391        var correlationHeader = KafkaTransportOptionsValidator.Required(
 42392            options.CorrelationIdHeader,
 42393            nameof(options.CorrelationIdHeader));
 94
 42395        return string.IsNullOrWhiteSpace(correlationId)
 42396            ? []
 42397            : [KafkaTransportHeader.Utf8(correlationHeader, correlationId)];
 98    }
 99}