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

Information
Class: AsyncResponse.Transports.SqlServer.SqlServerWorkerSubscriber
Assembly: AsyncResponse.Transports.SqlServer
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.SqlServer/SqlServerSubscriberServices.cs
Line coverage
100%
Covered lines: 6
Uncovered lines: 0
Coverable lines: 6
Total lines: 176
Line coverage: 100%
Branch coverage
N/A
Covered branches: 0
Total branches: 0
Branch coverage: N/A
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
get_Queue()100%11100%
get_SubscriberOptions()100%11100%
get_Role()100%11100%
HandleMessageAsync(...)100%11100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.SqlServer/SqlServerSubscriberServices.cs

#LineLine coverage
 1using Microsoft.Extensions.Hosting;
 2using Microsoft.Extensions.Logging;
 3using Microsoft.Extensions.Options;
 4using System.Threading.Channels;
 5
 6namespace AsyncResponse.Transports.SqlServer;
 7
 8/// <summary>
 9/// Base hosted service that consumes one SQL Server queue and routes rows to AsyncResponse ingress
 10/// with configured acknowledgement, redelivery, and dead-letter behavior.
 11/// </summary>
 12internal abstract class SqlServerSubscriberService : BackgroundService
 13{
 14    private readonly SqlServerTransportStore _store;
 15    private readonly Channel<bool> _signals = Channel.CreateBounded<bool>(new BoundedChannelOptions(1)
 16    {
 17        SingleReader = true,
 18        SingleWriter = false,
 19        FullMode = BoundedChannelFullMode.DropWrite
 20    });
 21
 22    protected SqlServerSubscriberService(
 23        IOptions<SqlServerAsyncResponseTransportOptions> options,
 24        SqlServerTransportStore store,
 25        ILogger logger)
 26    {
 27        Options = options.Value;
 28        SqlServerTransportOptionsValidator.ValidateCommon(Options);
 29        _store = store;
 30        Logger = logger;
 31    }
 32
 33    protected SqlServerAsyncResponseTransportOptions Options { get; }
 34    protected ILogger Logger { get; }
 35
 36    protected abstract string Queue { get; }
 37    protected abstract SqlServerSubscriberOptions SubscriberOptions { get; }
 38    protected abstract SqlServerSubscriberRole Role { get; }
 39    protected abstract Task HandleMessageAsync(SqlServerTransportDelivery delivery, CancellationToken cancellationToken)
 40
 41    /// <inheritdoc />
 42    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
 43    {
 44        SqlServerTransportOptionsValidator.ValidateSubscriber(Options, SubscriberOptions, Role.ToString());
 45
 46        var failures = 0;
 47        while (!stoppingToken.IsCancellationRequested)
 48        {
 49            try
 50            {
 51                await RunSubscriberAsync(stoppingToken).ConfigureAwait(false);
 52                return;
 53            }
 54            catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 55            {
 56                return;
 57            }
 58            catch (Exception ex) when (!stoppingToken.IsCancellationRequested)
 59            {
 60                failures++;
 61                var delay = AsyncResponseRetry.Backoff(failures, Options.SubscriberRetryBaseDelay, Options.SubscriberRet
 62                Logger.LogWarning(ex, "SQL Server subscriber failed for queue {Queue} ({Role}); retrying in {RetryDelay}
 63                await Task.Delay(delay, stoppingToken).ConfigureAwait(false);
 64            }
 65        }
 66    }
 67
 68    private async Task RunSubscriberAsync(CancellationToken stoppingToken)
 69    {
 70        await _store.EnsureCreatedAsync(stoppingToken).ConfigureAwait(false);
 71
 72        // Same-process wake: publishes to this queue (or a NAK release, queue == null) signal the
 73        // loop directly since SQL Server has no LISTEN/NOTIFY; cross-process publishes are picked up
 74        // by the EmptyPollDelay poll below.
 75        Action<string?> onPublished = queue =>
 76        {
 77            if (queue is null || string.Equals(queue, Queue, StringComparison.Ordinal))
 78                _signals.Writer.TryWrite(true);
 79        };
 80        _store.MessagePublished += onPublished;
 81
 82        await using var dispatcher = new SqlServerMessageDispatcher(
 83            HandleMessageAsync,
 84            Options,
 85            SubscriberOptions,
 86            Logger,
 87            Role);
 88
 89        Logger.LogInformation(
 90            "SQL Server subscriber started. Queue: {Queue}. Role: {Role}. AckMode: {AckMode}.",
 91            Queue,
 92            Role,
 93            SubscriberOptions.AckMode);
 94
 95        try
 96        {
 97            while (!stoppingToken.IsCancellationRequested)
 98            {
 99                var claimed = 0;
 100                await foreach (var delivery in _store.ClaimBatchAsync(Queue, SubscriberOptions.BatchSize, Options.LockTi
 101                {
 102                    claimed++;
 103                    await dispatcher.HandleAsync(delivery, stoppingToken).ConfigureAwait(false);
 104                }
 105
 106                if (claimed > 0)
 107                    continue;
 108
 109                await WaitForSignalOrDelayAsync(stoppingToken).ConfigureAwait(false);
 110            }
 111        }
 112        finally
 113        {
 114            _store.MessagePublished -= onPublished;
 115        }
 116    }
 117
 118    private async Task WaitForSignalOrDelayAsync(CancellationToken cancellationToken)
 119    {
 120        var delay = Task.Delay(SubscriberOptions.EmptyPollDelay, cancellationToken);
 121        var signal = _signals.Reader.WaitToReadAsync(cancellationToken).AsTask();
 122        var completed = await Task.WhenAny(delay, signal).ConfigureAwait(false);
 123        if (completed == signal)
 124        {
 125            await signal.ConfigureAwait(false);
 126            while (_signals.Reader.TryRead(out _))
 127            {
 128            }
 129        }
 130    }
 131}
 132
 133/// <summary>Consumes worker-job rows and executes them through the AsyncResponse ingress.</summary>
 134internal sealed class SqlServerWorkerSubscriber : SqlServerSubscriberService
 135{
 136    private readonly IAsyncResponseIngress _ingress;
 137
 138    public SqlServerWorkerSubscriber(
 139        IOptions<SqlServerAsyncResponseTransportOptions> options,
 140        SqlServerTransportStore store,
 141        IAsyncResponseIngress ingress,
 142        ILogger<SqlServerWorkerSubscriber> logger)
 3143        : base(options, store, logger)
 3144        => _ingress = ingress;
 145
 1146    protected override string Queue => Options.WorkerQueue;
 1147    protected override SqlServerSubscriberOptions SubscriberOptions => Options.WorkerSubscriber;
 1148    protected override SqlServerSubscriberRole Role => SqlServerSubscriberRole.Worker;
 149
 150    protected override Task HandleMessageAsync(SqlServerTransportDelivery delivery, CancellationToken cancellationToken)
 3151        => _ingress.HandleWorkerMessageAsync(delivery.Payload);
 152}
 153
 154/// <summary>Consumes response rows and feeds them into the AsyncResponse ingress.</summary>
 155internal sealed class SqlServerResponseIngressSubscriber : SqlServerSubscriberService
 156{
 157    private readonly IAsyncResponseIngress _ingress;
 158
 159    public SqlServerResponseIngressSubscriber(
 160        IOptions<SqlServerAsyncResponseTransportOptions> options,
 161        SqlServerTransportStore store,
 162        IAsyncResponseIngress ingress,
 163        ILogger<SqlServerResponseIngressSubscriber> logger)
 164        : base(options, store, logger)
 165        => _ingress = ingress;
 166
 167    protected override string Queue => Options.ResponseQueue;
 168    protected override SqlServerSubscriberOptions SubscriberOptions => Options.ResponseSubscriber;
 169    protected override SqlServerSubscriberRole Role => SqlServerSubscriberRole.ResponseIngress;
 170
 171    protected override Task HandleMessageAsync(SqlServerTransportDelivery delivery, CancellationToken cancellationToken)
 172    {
 173        var correlationId = SqlServerCorrelationIdExtractor.Extract(delivery.Headers, delivery.Payload, Options);
 174        return _ingress.HandleResponseMessageAsync(delivery.Payload, correlationId);
 175    }
 176}