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

Information
Class: AsyncResponse.Transports.SqlServer.SqlServerWorkerTransport
Assembly: AsyncResponse.Transports.SqlServer
File(s): /_/src/Transports/AsyncResponse.Transports.SqlServer/SqlServerWorkerTransport.cs
Line coverage
95%
Covered lines: 46
Uncovered lines: 2
Coverable lines: 48
Total lines: 103
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.SqlServer/SqlServerWorkerTransport.cs

#LineLine coverage
 1using AsyncResponse.Internal;
 2using Microsoft.Extensions.Options;
 3using System.Diagnostics;
 4using System.Text.Json;
 5
 6namespace AsyncResponse.Transports.SqlServer;
 7
 8/// <summary>Publishes <see cref="WorkerJobEnvelope"/> messages to the SQL Server worker queue.</summary>
 9public sealed class SqlServerWorkerTransport : IWorkerTransport, IDelayedWorkerTransport
 10{
 11    private readonly SqlServerAsyncResponseTransportOptions _options;
 12    private readonly SqlServerTransportStore _store;
 13
 14    /// <summary>Creates a SQL Server worker transport over the configured connection string.</summary>
 15    public SqlServerWorkerTransport(IOptions<SqlServerAsyncResponseTransportOptions> options)
 016        : this(options, new SqlServerTransportStore(options))
 17    {
 018    }
 19
 19920    internal SqlServerWorkerTransport(
 19921        IOptions<SqlServerAsyncResponseTransportOptions> options,
 19922        SqlServerTransportStore store)
 23    {
 19924        _options = options.Value;
 19925        SqlServerTransportOptionsValidator.ValidateCommon(_options);
 19926        _store = store;
 19927    }
 28
 29    /// <inheritdoc />
 30    public Task PublishAsync(WorkerJobEnvelope job, CancellationToken cancellationToken = default)
 40931        => PublishCoreAsync(job, delay: null, cancellationToken);
 32
 33    /// <inheritdoc cref="IDelayedWorkerTransport.MaxPublishDelay"/>
 34    /// <remarks>The due time is a database timestamp; no per-hop cap applies.</remarks>
 535    public TimeSpan MaxPublishDelay => AsyncResponseChannelOptions.MaxPersistenceTtl;
 36
 37    /// <inheritdoc cref="IDelayedWorkerTransport.PublishAsync(WorkerJobEnvelope, TimeSpan, CancellationToken)"/>
 38    public Task PublishAsync(WorkerJobEnvelope job, TimeSpan delay, CancellationToken cancellationToken = default)
 39    {
 140        ArgumentOutOfRangeException.ThrowIfGreaterThan(delay, MaxPublishDelay);
 141        return PublishCoreAsync(job, delay > TimeSpan.Zero ? delay : null, cancellationToken);
 42    }
 43
 44    private async Task PublishCoreAsync(WorkerJobEnvelope job, TimeSpan? delay, CancellationToken cancellationToken)
 45    {
 41046        ArgumentNullException.ThrowIfNull(job);
 47
 41048        using var activity = AsyncResponseDiagnostics.StartActivity(
 41049            "asyncresponse.worker.publish",
 41050            ActivityKind.Producer,
 41051            job.CorrelationId);
 41052        activity?.SetTag("asyncresponse.transport", "sqlserver");
 41053        activity?.SetTag("messaging.system", "sqlserver");
 41054        activity?.SetTag("messaging.destination.name", _options.WorkerQueue);
 41055        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 41056        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 57
 58        try
 59        {
 41060            var headers = string.IsNullOrWhiteSpace(job.CorrelationId)
 41061                ? null
 41062                : new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase)
 41063                {
 41064                    [_options.CorrelationIdHeader] = job.CorrelationId!
 41065                };
 66
 41067            var payload = AsyncResponseJson.Serialize(job);
 68            // Stable id outside the retry loop so a retried publish is idempotent rather than enqueuing
 69            // the same worker job twice.
 41070            var messageId = Guid.NewGuid();
 41071            if (delay is { } delayTag)
 172                activity?.SetTag("asyncresponse.worker.delay_seconds", delayTag.TotalSeconds);
 41073            await SqlServerTransportRetry.ExecuteAsync(
 41074                async token =>
 41075                {
 41076                    await _store.PublishAsync(messageId, _options.WorkerQueue, payload, headers, token, delay).Configure
 40877                    return true;
 40878                },
 41079                _options.PublishMaxAttempts,
 41080                _options.PublishRetryBaseDelay,
 41081                _options.PublishRetryMaxDelay,
 41082                cancellationToken).ConfigureAwait(false);
 40883        }
 284        catch (Exception ex)
 85        {
 286            AsyncResponseDiagnostics.SetError(activity, ex);
 287            throw;
 88        }
 40889    }
 90}
 91
 92internal static class SqlServerTransportRetry
 93{
 94    public static Task<T> ExecuteAsync<T>(
 95        Func<CancellationToken, Task<T>> action,
 96        int maxAttempts,
 97        TimeSpan baseDelay,
 98        TimeSpan maxDelay,
 99        CancellationToken cancellationToken)
 100        => AsyncResponseRetry.ExecuteAsync(action, IsTransient, maxAttempts, baseDelay, maxDelay, cancellationToken);
 101
 102    public static bool IsTransient(Exception exception) => SqlServerTransientFaults.IsTransient(exception);
 103}