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

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

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
ExecuteAsync()100%22100%
RunSubscriberAsync()100%44100%
WaitForSignalOrDelayAsync()100%44100%

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;
 315    private readonly Channel<bool> _signals = Channel.CreateBounded<bool>(new BoundedChannelOptions(1)
 316    {
 317        SingleReader = true,
 318        SingleWriter = false,
 319        FullMode = BoundedChannelFullMode.DropWrite
 320    });
 21
 322    protected SqlServerSubscriberService(
 323        IOptions<SqlServerAsyncResponseTransportOptions> options,
 324        SqlServerTransportStore store,
 325        ILogger logger)
 26    {
 327        Options = options.Value;
 328        SqlServerTransportOptionsValidator.ValidateCommon(Options);
 329        _store = store;
 330        Logger = logger;
 331    }
 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    {
 144        SqlServerTransportOptionsValidator.ValidateSubscriber(Options, SubscriberOptions, Role.ToString());
 45
 146        var failures = 0;
 147        while (!stoppingToken.IsCancellationRequested)
 48        {
 49            try
 50            {
 151                await RunSubscriberAsync(stoppingToken).ConfigureAwait(false);
 152                return;
 53            }
 154            catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 55            {
 156                return;
 57            }
 158            catch (Exception ex) when (!stoppingToken.IsCancellationRequested)
 59            {
 160                failures++;
 161                var delay = AsyncResponseRetry.Backoff(failures, Options.SubscriberRetryBaseDelay, Options.SubscriberRet
 162                Logger.LogWarning(ex, "SQL Server subscriber failed for queue {Queue} ({Role}); retrying in {RetryDelay}
 163                await Task.Delay(delay, stoppingToken).ConfigureAwait(false);
 64            }
 65        }
 166    }
 67
 68    private async Task RunSubscriberAsync(CancellationToken stoppingToken)
 69    {
 170        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.
 175        Action<string?> onPublished = queue =>
 176        {
 177            if (queue is null || string.Equals(queue, Queue, StringComparison.Ordinal))
 178                _signals.Writer.TryWrite(true);
 179        };
 180        _store.MessagePublished += onPublished;
 81
 182        await using var dispatcher = new SqlServerMessageDispatcher(
 183            HandleMessageAsync,
 184            Options,
 185            SubscriberOptions,
 186            Logger,
 187            Role);
 88
 189        Logger.LogInformation(
 190            "SQL Server subscriber started. Queue: {Queue}. Role: {Role}. AckMode: {AckMode}.",
 191            Queue,
 192            Role,
 193            SubscriberOptions.AckMode);
 94
 95        try
 96        {
 197            while (!stoppingToken.IsCancellationRequested)
 98            {
 199                var claimed = 0;
 1100                await foreach (var delivery in _store.ClaimBatchAsync(Queue, SubscriberOptions.BatchSize, Options.LockTi
 101                {
 1102                    claimed++;
 1103                    await dispatcher.HandleAsync(delivery, stoppingToken).ConfigureAwait(false);
 104                }
 105
 1106                if (claimed > 0)
 107                    continue;
 108
 1109                await WaitForSignalOrDelayAsync(stoppingToken).ConfigureAwait(false);
 110            }
 1111        }
 112        finally
 113        {
 1114            _store.MessagePublished -= onPublished;
 115        }
 1116    }
 117
 118    private async Task WaitForSignalOrDelayAsync(CancellationToken cancellationToken)
 119    {
 1120        var delay = Task.Delay(SubscriberOptions.EmptyPollDelay, cancellationToken);
 1121        var signal = _signals.Reader.WaitToReadAsync(cancellationToken).AsTask();
 1122        var completed = await Task.WhenAny(delay, signal).ConfigureAwait(false);
 1123        if (completed == signal)
 124        {
 1125            await signal.ConfigureAwait(false);
 1126            while (_signals.Reader.TryRead(out _))
 127            {
 128            }
 129        }
 1130    }
 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)
 143        : base(options, store, logger)
 144        => _ingress = ingress;
 145
 146    protected override string Queue => Options.WorkerQueue;
 147    protected override SqlServerSubscriberOptions SubscriberOptions => Options.WorkerSubscriber;
 148    protected override SqlServerSubscriberRole Role => SqlServerSubscriberRole.Worker;
 149
 150    protected override Task HandleMessageAsync(SqlServerTransportDelivery delivery, CancellationToken cancellationToken)
 151        => _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}