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

Information
Class: AsyncResponse.Transports.MongoDB.MongoDbTransportRetry
Assembly: AsyncResponse.Transports.MongoDB
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.MongoDB/MongoDbWorkerTransport.cs
Line coverage
100%
Covered lines: 8
Uncovered lines: 0
Coverable lines: 8
Total lines: 96
Line coverage: 100%
Branch coverage
100%
Covered branches: 18
Total branches: 18
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
ExecuteAsync<T>(...)100%22100%
IsTransient(...)100%1616100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.MongoDB/MongoDbWorkerTransport.cs

#LineLine coverage
 1using Microsoft.Extensions.Options;
 2using MongoDB.Driver;
 3using System.Diagnostics;
 4using System.Text.Json;
 5
 6namespace AsyncResponse.Transports.MongoDB;
 7
 8/// <summary>Publishes <see cref="WorkerJobEnvelope"/> messages to the MongoDB worker queue.</summary>
 9public sealed class MongoDbWorkerTransport : IWorkerTransport
 10{
 11    private readonly MongoDbAsyncResponseTransportOptions _options;
 12    private readonly MongoDbTransportStore _store;
 13
 14    /// <summary>Creates a MongoDB worker transport over the host's shared database.</summary>
 15    public MongoDbWorkerTransport(
 16        IOptions<MongoDbAsyncResponseTransportOptions> options,
 17        IMongoDatabase database)
 18        : this(options, new MongoDbTransportStore(database, options))
 19    {
 20    }
 21
 22    internal MongoDbWorkerTransport(
 23        IOptions<MongoDbAsyncResponseTransportOptions> options,
 24        MongoDbTransportStore store)
 25    {
 26        _options = options.Value;
 27        MongoDbTransportOptionsValidator.ValidateCommon(_options);
 28        _store = store;
 29    }
 30
 31    /// <inheritdoc />
 32    public async Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 33    {
 34        ArgumentNullException.ThrowIfNull(job);
 35
 36        using var activity = AsyncResponseDiagnostics.StartActivity(
 37            "asyncresponse.worker.publish",
 38            ActivityKind.Producer,
 39            job.CorrelationId);
 40        activity?.SetTag("asyncresponse.transport", "mongodb");
 41        activity?.SetTag("messaging.system", "mongodb");
 42        activity?.SetTag("messaging.destination.name", _options.WorkerQueue);
 43        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 44        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 45
 46        try
 47        {
 48            var headers = string.IsNullOrWhiteSpace(job.CorrelationId)
 49                ? null
 50                : new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase)
 51                {
 52                    [_options.CorrelationIdHeader] = job.CorrelationId!
 53                };
 54
 55            var payload = AsyncResponseJson.Serialize(job);
 56            // Stable id outside the retry loop so a retried publish is idempotent rather than enqueuing
 57            // the same worker job twice.
 58            var messageId = Guid.NewGuid();
 59            await MongoDbTransportRetry.ExecuteAsync(
 60                async token =>
 61                {
 62                    await _store.PublishAsync(messageId, _options.WorkerQueue, payload, headers, token).ConfigureAwait(f
 63                    return true;
 64                },
 65                _options.PublishMaxAttempts,
 66                _options.PublishRetryBaseDelay,
 67                _options.PublishRetryMaxDelay,
 68                cancellationToken).ConfigureAwait(false);
 69        }
 70        catch (Exception ex)
 71        {
 72            AsyncResponseDiagnostics.SetError(activity, ex);
 73            throw;
 74        }
 75    }
 76}
 77
 78internal static class MongoDbTransportRetry
 79{
 80    public static Task<T> ExecuteAsync<T>(
 81        Func<CancellationToken, Task<T>> action,
 82        int maxAttempts,
 83        TimeSpan baseDelay,
 84        TimeSpan maxDelay,
 85        CancellationToken cancellationToken)
 386        => AsyncResponseRetry.ExecuteAsync(action, IsTransient, maxAttempts, baseDelay, maxDelay, cancellationToken);
 87
 88    public static bool IsTransient(Exception exception)
 389        => exception is not OperationCanceledException
 390           && (exception is MongoConnectionException
 391               or MongoNotPrimaryException
 392               or MongoNodeIsRecoveringException
 393               or MongoExecutionTimeoutException
 394               or TimeoutException
 395               || (exception is MongoException mongoException && mongoException.HasErrorLabel("RetryableWriteError")));
 96}