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

Information
Class: AsyncResponse.Transports.MongoDB.MongoDbWorkerTransport
Assembly: AsyncResponse.Transports.MongoDB
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.MongoDB/MongoDbWorkerTransport.cs
Line coverage
100%
Covered lines: 42
Uncovered lines: 0
Coverable lines: 42
Total lines: 96
Line coverage: 100%
Branch coverage
62%
Covered branches: 5
Total branches: 8
Branch coverage: 62.5%
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%
PublishAsync()62.5%88100%
<PublishAsync()100%11100%

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)
 318        : this(options, new MongoDbTransportStore(database, options))
 19    {
 320    }
 21
 322    internal MongoDbWorkerTransport(
 323        IOptions<MongoDbAsyncResponseTransportOptions> options,
 324        MongoDbTransportStore store)
 25    {
 326        _options = options.Value;
 327        MongoDbTransportOptionsValidator.ValidateCommon(_options);
 328        _store = store;
 329    }
 30
 31    /// <inheritdoc />
 32    public async Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 33    {
 334        ArgumentNullException.ThrowIfNull(job);
 35
 336        using var activity = AsyncResponseDiagnostics.StartActivity(
 337            "asyncresponse.worker.publish",
 338            ActivityKind.Producer,
 339            job.CorrelationId);
 340        activity?.SetTag("asyncresponse.transport", "mongodb");
 341        activity?.SetTag("messaging.system", "mongodb");
 342        activity?.SetTag("messaging.destination.name", _options.WorkerQueue);
 343        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 344        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 45
 46        try
 47        {
 348            var headers = string.IsNullOrWhiteSpace(job.CorrelationId)
 349                ? null
 350                : new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase)
 351                {
 352                    [_options.CorrelationIdHeader] = job.CorrelationId!
 353                };
 54
 355            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.
 358            var messageId = Guid.NewGuid();
 359            await MongoDbTransportRetry.ExecuteAsync(
 360                async token =>
 361                {
 362                    await _store.PublishAsync(messageId, _options.WorkerQueue, payload, headers, token).ConfigureAwait(f
 363                    return true;
 364                },
 365                _options.PublishMaxAttempts,
 366                _options.PublishRetryBaseDelay,
 367                _options.PublishRetryMaxDelay,
 368                cancellationToken).ConfigureAwait(false);
 369        }
 270        catch (Exception ex)
 71        {
 272            AsyncResponseDiagnostics.SetError(activity, ex);
 373            throw;
 74        }
 375    }
 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)
 86        => AsyncResponseRetry.ExecuteAsync(action, IsTransient, maxAttempts, baseDelay, maxDelay, cancellationToken);
 87
 88    public static bool IsTransient(Exception exception)
 89        => exception is not OperationCanceledException
 90           && (exception is MongoConnectionException
 91               or MongoNotPrimaryException
 92               or MongoNodeIsRecoveringException
 93               or MongoExecutionTimeoutException
 94               or TimeoutException
 95               || (exception is MongoException mongoException && mongoException.HasErrorLabel("RetryableWriteError")));
 96}