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

Information
Class: AsyncResponse.Transports.Kafka.KafkaWorkerTransport
Assembly: AsyncResponse.Transports.Kafka
File(s): /home/runner/work/AsyncResponse/AsyncResponse/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)

/home/runner/work/AsyncResponse/AsyncResponse/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    {
 327    }
 28
 329    internal KafkaWorkerTransport(
 330        IOptions<KafkaAsyncResponseTransportOptions> options,
 331        IKafkaProducerClient producer)
 32    {
 333        _options = options.Value;
 334        KafkaTransportOptionsValidator.ValidateCommon(_options);
 335        _producer = producer;
 336        _topics = new KafkaTransportTopicSchema(_options);
 337    }
 38
 39    private static KafkaAsyncResponseTransportOptions KafkaTransportOptionsValidatorWithValue(
 40        IOptions<KafkaAsyncResponseTransportOptions> options)
 41    {
 242        KafkaTransportOptionsValidator.ValidateCommon(options.Value);
 343        return options.Value;
 44    }
 45
 46    /// <summary>Publishes the supplied message.</summary>
 47    public async Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 48    {
 349        ArgumentNullException.ThrowIfNull(job);
 50
 351        using var activity = AsyncResponseDiagnostics.StartActivity(
 352            "asyncresponse.worker.publish",
 353            ActivityKind.Producer,
 354            job.CorrelationId);
 355        activity?.SetTag("asyncresponse.transport", "kafka");
 356        activity?.SetTag("messaging.system", "kafka");
 357        activity?.SetTag("messaging.destination.name", _topics.WorkerTopic);
 358        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 359        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 60
 61        try
 62        {
 363            var payload = Encoding.UTF8.GetBytes(AsyncResponseJson.Serialize(job));
 364            var headers = CreateMessageHeaders(job.CorrelationId, _options);
 365            var result = await KafkaTransportRetry.ExecuteAsync(
 366                token => _producer.PublishAsync(
 367                    _topics.WorkerTopic,
 368                    string.IsNullOrWhiteSpace(job.CorrelationId) ? null : job.CorrelationId,
 369                    payload,
 370                    headers,
 371                    token),
 372                _options.PublishMaxAttempts,
 373                _options.PublishRetryBaseDelay,
 374                _options.PublishRetryMaxDelay,
 375                cancellationToken).ConfigureAwait(false);
 76
 377            activity?.SetTag("messaging.kafka.destination.partition", result.Partition);
 378            activity?.SetTag("messaging.kafka.message.offset", result.Offset);
 379        }
 280        catch (Exception ex)
 81        {
 282            AsyncResponseDiagnostics.SetError(activity, ex);
 383            throw;
 84        }
 385    }
 86
 87    internal static IReadOnlyList<KafkaTransportHeader> CreateMessageHeaders(
 88        string? correlationId,
 89        KafkaAsyncResponseTransportOptions options)
 90    {
 391        var correlationHeader = KafkaTransportOptionsValidator.Required(
 392            options.CorrelationIdHeader,
 393            nameof(options.CorrelationIdHeader));
 94
 395        return string.IsNullOrWhiteSpace(correlationId)
 396            ? []
 397            : [KafkaTransportHeader.Utf8(correlationHeader, correlationId)];
 98    }
 99}