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

Information
Class: AsyncResponse.Transports.PostgreSQL.PostgreSqlWorkerTransport
Assembly: AsyncResponse.Transports.PostgreSQL
File(s): /_/src/Transports/AsyncResponse.Transports.PostgreSQL/PostgreSqlWorkerTransport.cs
Line coverage
95%
Covered lines: 46
Uncovered lines: 2
Coverable lines: 48
Total lines: 106
Line coverage: 95.8%
Branch coverage
64%
Covered branches: 9
Total branches: 14
Branch coverage: 64.2%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%210%
.ctor(...)100%11100%
PublishAsync(...)100%11100%
get_MaxPublishDelay()100%11100%
PublishAsync(...)50%22100%
PublishCoreAsync()66.66%121290.62%
<PublishCoreAsync()100%11100%

File(s)

/_/src/Transports/AsyncResponse.Transports.PostgreSQL/PostgreSqlWorkerTransport.cs

#LineLine coverage
 1using AsyncResponse.Internal;
 2using Microsoft.Extensions.Options;
 3using Npgsql;
 4using System.Diagnostics;
 5using System.Text.Json;
 6
 7namespace AsyncResponse.Transports.PostgreSQL;
 8
 9/// <summary>Publishes <see cref="WorkerJobEnvelope"/> messages to the PostgreSQL worker queue.</summary>
 10public sealed class PostgreSqlWorkerTransport : IWorkerTransport, IDelayedWorkerTransport
 11{
 12    private readonly PostgreSqlAsyncResponseTransportOptions _options;
 13    private readonly PostgreSqlTransportStore _store;
 14
 15    /// <summary>Creates a PostgreSQL worker transport over the host's shared data source.</summary>
 16    public PostgreSqlWorkerTransport(
 17        IOptions<PostgreSqlAsyncResponseTransportOptions> options,
 18        NpgsqlDataSource dataSource)
 019        : this(options, new PostgreSqlTransportStore(dataSource, options))
 20    {
 021    }
 22
 19923    internal PostgreSqlWorkerTransport(
 19924        IOptions<PostgreSqlAsyncResponseTransportOptions> options,
 19925        PostgreSqlTransportStore store)
 26    {
 19927        _options = options.Value;
 19928        PostgreSqlTransportOptionsValidator.ValidateCommon(_options);
 19929        _store = store;
 19930    }
 31
 32    /// <inheritdoc />
 33    public Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 40934        => PublishCoreAsync(job, delay: null, cancellationToken);
 35
 36    /// <inheritdoc cref="IDelayedWorkerTransport.MaxPublishDelay"/>
 37    /// <remarks>The due time is a database timestamp; no per-hop cap applies.</remarks>
 538    public TimeSpan MaxPublishDelay => AsyncResponseChannelOptions.MaxPersistenceTtl;
 39
 40    /// <inheritdoc cref="IDelayedWorkerTransport.PublishAsync(WorkerJobEnvelope, TimeSpan, CancellationToken)"/>
 41    public Task PublishAsync(WorkerJobEnvelope job, TimeSpan delay, CancellationToken cancellationToken = default)
 42    {
 143        ArgumentOutOfRangeException.ThrowIfGreaterThan(delay, MaxPublishDelay);
 144        return PublishCoreAsync(job, delay > TimeSpan.Zero ? delay : null, cancellationToken);
 45    }
 46
 47    private async Task PublishCoreAsync(WorkerJobEnvelope job, TimeSpan? delay, CancellationToken cancellationToken)
 48    {
 41049        ArgumentNullException.ThrowIfNull(job);
 50
 41051        using var activity = AsyncResponseDiagnostics.StartActivity(
 41052            "asyncresponse.worker.publish",
 41053            ActivityKind.Producer,
 41054            job.CorrelationId);
 41055        activity?.SetTag("asyncresponse.transport", "postgresql");
 41056        activity?.SetTag("messaging.system", "postgresql");
 41057        activity?.SetTag("messaging.destination.name", _options.WorkerQueue);
 41058        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 41059        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 60
 61        try
 62        {
 41063            var headers = string.IsNullOrWhiteSpace(job.CorrelationId)
 41064                ? null
 41065                : new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase)
 41066                {
 41067                    [_options.CorrelationIdHeader] = job.CorrelationId!
 41068                };
 69
 41070            var payload = AsyncResponseJson.Serialize(job);
 71            // Stable id outside the retry loop so a retried publish is idempotent rather than enqueuing
 72            // the same worker job twice.
 41073            var messageId = Guid.NewGuid();
 41074            if (delay is { } delayTag)
 175                activity?.SetTag("asyncresponse.worker.delay_seconds", delayTag.TotalSeconds);
 41076            await PostgreSqlTransportRetry.ExecuteAsync(
 41077                async token =>
 41078                {
 41079                    await _store.PublishAsync(messageId, _options.WorkerQueue, payload, headers, token, delay).Configure
 40880                    return true;
 40881                },
 41082                _options.PublishMaxAttempts,
 41083                _options.PublishRetryBaseDelay,
 41084                _options.PublishRetryMaxDelay,
 41085                cancellationToken).ConfigureAwait(false);
 40886        }
 287        catch (Exception ex)
 88        {
 289            AsyncResponseDiagnostics.SetError(activity, ex);
 290            throw;
 91        }
 40892    }
 93}
 94
 95internal static class PostgreSqlTransportRetry
 96{
 97    public static Task<T> ExecuteAsync<T>(
 98        Func<CancellationToken, Task<T>> action,
 99        int maxAttempts,
 100        TimeSpan baseDelay,
 101        TimeSpan maxDelay,
 102        CancellationToken cancellationToken)
 103        => AsyncResponseRetry.ExecuteAsync(action, IsTransient, maxAttempts, baseDelay, maxDelay, cancellationToken);
 104
 105    public static bool IsTransient(Exception exception) => PostgreSqlTransientFaults.IsTransient(exception);
 106}