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

Information
Class: AsyncResponse.Transports.Redis.RedisWorkerTransport
Assembly: AsyncResponse.Transports.Redis
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.Redis/RedisWorkerTransport.cs
Line coverage
100%
Covered lines: 61
Uncovered lines: 0
Coverable lines: 61
Total lines: 117
Line coverage: 100%
Branch coverage
100%
Covered branches: 10
Total branches: 10
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%
RedisTransportOptionsValidatorWithValue(...)100%11100%
ValidatePublishOptions(...)100%11100%
PublishAsync()100%88100%
CreateMessageFields(...)100%22100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.Redis/RedisWorkerTransport.cs

#LineLine coverage
 1using Microsoft.Extensions.Options;
 2using StackExchange.Redis;
 3using System.Diagnostics;
 4using System.Text.Json;
 5
 6namespace AsyncResponse.Transports.Redis;
 7
 8/// <summary>
 9/// Publishes <see cref="WorkerJobEnvelope"/> messages to a Redis stream.
 10/// </summary>
 11/// <remarks>
 12/// Redis Streams provide durable queueing and consumer-group acknowledgement. Publishing uses XADD
 13/// with optional approximate trimming, and transient Redis failures are retried with bounded
 14/// exponential backoff before the exception is returned to the caller.
 15/// </remarks>
 16public sealed class RedisWorkerTransport : IWorkerTransport
 17{
 18    private readonly RedisAsyncResponseTransportOptions _options;
 19    private readonly IRedisStreamDatabase _database;
 20    private readonly RedisTransportKeySchema _keys;
 21
 22    /// <summary>Runs the RedisWorkerTransport operation.</summary>
 23    public RedisWorkerTransport(
 24        IOptions<RedisAsyncResponseTransportOptions> options,
 25        IConnectionMultiplexer multiplexer)
 326        : this(
 327            options,
 328            new RedisStreamDatabaseAdapter(
 329                multiplexer.GetDatabase(),
 330                RedisTransportOptionsValidatorWithValue(options).OperationTimeout))
 31    {
 332    }
 33
 334    internal RedisWorkerTransport(
 335        IOptions<RedisAsyncResponseTransportOptions> options,
 336        IRedisStreamDatabase database)
 37    {
 338        _options = options.Value;
 339        RedisTransportOptionsValidator.ValidateCommon(_options);
 340        ValidatePublishOptions(_options);
 341        _database = database;
 342        _keys = new RedisTransportKeySchema(_options);
 343    }
 44
 45    private static RedisAsyncResponseTransportOptions RedisTransportOptionsValidatorWithValue(
 46        IOptions<RedisAsyncResponseTransportOptions> options)
 47    {
 348        RedisTransportOptionsValidator.ValidateCommon(options.Value);
 349        ValidatePublishOptions(options.Value);
 350        return options.Value;
 51    }
 52
 53    private static void ValidatePublishOptions(RedisAsyncResponseTransportOptions options)
 54    {
 355        _ = new RedisTransportKeySchema(options).WorkerStream;
 356        _ = RedisTransportOptionsValidator.Required(options.PayloadField, nameof(options.PayloadField));
 357    }
 58
 59    /// <summary>Publishes the supplied message.</summary>
 60    public async Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 61    {
 362        ArgumentNullException.ThrowIfNull(job);
 63
 364        using var activity = AsyncResponseDiagnostics.StartActivity(
 365            "asyncresponse.worker.publish",
 366            ActivityKind.Producer,
 367            job.CorrelationId);
 368        activity?.SetTag("asyncresponse.transport", "redis");
 369        activity?.SetTag("messaging.system", "redis");
 370        activity?.SetTag("messaging.destination.name", _keys.WorkerStream.ToString());
 371        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 372        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 73
 74        try
 75        {
 376            var fields = CreateMessageFields(
 377                AsyncResponseJson.Serialize(job),
 378                job.CorrelationId,
 379                _options);
 380            var messageId = await RedisTransportRetry.ExecuteAsync(
 381                token => _database.StreamAddAsync(
 382                    _keys.WorkerStream,
 383                    fields,
 384                    _options.StreamMaxLength,
 385                    _options.UseApproximateStreamTrimming,
 386                    token),
 387                _options.PublishMaxAttempts,
 388                _options.PublishRetryBaseDelay,
 389                _options.PublishRetryMaxDelay,
 390                cancellationToken).ConfigureAwait(false);
 91
 392            activity?.SetTag("messaging.message.id", messageId.ToString());
 393        }
 294        catch (Exception ex)
 95        {
 296            AsyncResponseDiagnostics.SetError(activity, ex);
 397            throw;
 98        }
 399    }
 100
 101    internal static NameValueEntry[] CreateMessageFields(
 102        string payloadJson,
 103        string? correlationId,
 104        RedisAsyncResponseTransportOptions options)
 105    {
 3106        var payloadField = RedisTransportOptionsValidator.Required(options.PayloadField, nameof(options.PayloadField));
 3107        var correlationField = RedisTransportOptionsValidator.Required(options.CorrelationIdField, nameof(options.Correl
 108
 3109        return string.IsNullOrWhiteSpace(correlationId)
 3110            ? [new NameValueEntry(payloadField, payloadJson)]
 3111            :
 3112            [
 3113                new NameValueEntry(payloadField, payloadJson),
 3114                new NameValueEntry(correlationField, correlationId)
 3115            ];
 116    }
 117}