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

Information
Class: AsyncResponse.Transports.PostgreSQL.PostgreSqlSubscriberService
Assembly: AsyncResponse.Transports.PostgreSQL
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.PostgreSQL/PostgreSqlSubscriberServices.cs
Line coverage
98%
Covered lines: 73
Uncovered lines: 1
Coverable lines: 74
Total lines: 194
Line coverage: 98.6%
Branch coverage
100%
Covered branches: 10
Total branches: 10
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%4496.3%
ListenLoopAsync()100%11100%
WaitForSignalOrDelayAsync()100%44100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/AsyncResponse.Transports.PostgreSQL/PostgreSqlSubscriberServices.cs

#LineLine coverage
 1using Microsoft.Extensions.Hosting;
 2using Microsoft.Extensions.Logging;
 3using Microsoft.Extensions.Options;
 4using System.Threading.Channels;
 5
 6namespace AsyncResponse.Transports.PostgreSQL;
 7
 8/// <summary>
 9/// Base hosted service that consumes one PostgreSQL queue and routes rows to AsyncResponse ingress
 10/// with configured acknowledgement, redelivery, and dead-letter behavior.
 11/// </summary>
 12internal abstract class PostgreSqlSubscriberService : BackgroundService
 13{
 14    private readonly PostgreSqlTransportStore _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 PostgreSqlSubscriberService(
 323        IOptions<PostgreSqlAsyncResponseTransportOptions> options,
 324        PostgreSqlTransportStore store,
 325        ILogger logger)
 26    {
 327        Options = options.Value;
 328        PostgreSqlTransportOptionsValidator.ValidateCommon(Options);
 329        _store = store;
 330        Logger = logger;
 331    }
 32
 33    protected PostgreSqlAsyncResponseTransportOptions Options { get; }
 34    protected ILogger Logger { get; }
 35
 36    protected abstract string Queue { get; }
 37    protected abstract PostgreSqlSubscriberOptions SubscriberOptions { get; }
 38    protected abstract PostgreSqlSubscriberRole Role { get; }
 39    protected abstract Task HandleMessageAsync(PostgreSqlTransportDelivery delivery, CancellationToken cancellationToken
 40
 41    /// <inheritdoc />
 42    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
 43    {
 144        PostgreSqlTransportOptionsValidator.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, "PostgreSQL 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);
 171        using var signalCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
 172        var listenTask = Task.Run(() => ListenLoopAsync(signalCts.Token), signalCts.Token);
 73
 174        await using var dispatcher = new PostgreSqlMessageDispatcher(
 175            HandleMessageAsync,
 176            Options,
 177            SubscriberOptions,
 178            Logger,
 179            Role);
 80
 181        Logger.LogInformation(
 182            "PostgreSQL subscriber started. Queue: {Queue}. Role: {Role}. AckMode: {AckMode}.",
 183            Queue,
 184            Role,
 185            SubscriberOptions.AckMode);
 86
 87        try
 88        {
 189            while (!stoppingToken.IsCancellationRequested)
 90            {
 191                var claimed = 0;
 192                await foreach (var delivery in _store.ClaimBatchAsync(Queue, SubscriberOptions.BatchSize, Options.LockTi
 93                {
 194                    claimed++;
 195                    await dispatcher.HandleAsync(delivery, stoppingToken).ConfigureAwait(false);
 96                }
 97
 198                if (claimed > 0)
 99                    continue;
 100
 1101                await WaitForSignalOrDelayAsync(stoppingToken).ConfigureAwait(false);
 102            }
 103        }
 104        finally
 105        {
 1106            await signalCts.CancelAsync().ConfigureAwait(false);
 107            try
 108            {
 1109                await listenTask.WaitAsync(Options.ShutdownTimeout).ConfigureAwait(false);
 1110            }
 0111            catch (Exception ex) when (ex is OperationCanceledException or TimeoutException)
 112            {
 1113            }
 114        }
 1115    }
 116
 117    private async Task ListenLoopAsync(CancellationToken cancellationToken)
 118    {
 119        try
 120        {
 1121            await _store.ExecuteListenAsync(() =>
 1122            {
 1123                _signals.Writer.TryWrite(true);
 1124                return Task.CompletedTask;
 1125            }, cancellationToken).ConfigureAwait(false);
 1126        }
 1127        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 128        {
 1129        }
 1130        catch (Exception ex)
 131        {
 1132            Logger.LogDebug(ex, "PostgreSQL LISTEN helper for queue {Queue} stopped; polling continues.", Queue);
 1133        }
 1134    }
 135
 136    private async Task WaitForSignalOrDelayAsync(CancellationToken cancellationToken)
 137    {
 1138        var delay = Task.Delay(SubscriberOptions.EmptyPollDelay, cancellationToken);
 1139        var signal = _signals.Reader.WaitToReadAsync(cancellationToken).AsTask();
 1140        var completed = await Task.WhenAny(delay, signal).ConfigureAwait(false);
 1141        if (completed == signal)
 142        {
 1143            await signal.ConfigureAwait(false);
 1144            while (_signals.Reader.TryRead(out _))
 145            {
 146            }
 147        }
 1148    }
 149}
 150
 151/// <summary>Consumes worker-job rows and executes them through the AsyncResponse ingress.</summary>
 152internal sealed class PostgreSqlWorkerSubscriber : PostgreSqlSubscriberService
 153{
 154    private readonly IAsyncResponseIngress _ingress;
 155
 156    public PostgreSqlWorkerSubscriber(
 157        IOptions<PostgreSqlAsyncResponseTransportOptions> options,
 158        PostgreSqlTransportStore store,
 159        IAsyncResponseIngress ingress,
 160        ILogger<PostgreSqlWorkerSubscriber> logger)
 161        : base(options, store, logger)
 162        => _ingress = ingress;
 163
 164    protected override string Queue => Options.WorkerQueue;
 165    protected override PostgreSqlSubscriberOptions SubscriberOptions => Options.WorkerSubscriber;
 166    protected override PostgreSqlSubscriberRole Role => PostgreSqlSubscriberRole.Worker;
 167
 168    protected override Task HandleMessageAsync(PostgreSqlTransportDelivery delivery, CancellationToken cancellationToken
 169        => _ingress.HandleWorkerMessageAsync(delivery.Payload);
 170}
 171
 172/// <summary>Consumes response rows and feeds them into the AsyncResponse ingress.</summary>
 173internal sealed class PostgreSqlResponseIngressSubscriber : PostgreSqlSubscriberService
 174{
 175    private readonly IAsyncResponseIngress _ingress;
 176
 177    public PostgreSqlResponseIngressSubscriber(
 178        IOptions<PostgreSqlAsyncResponseTransportOptions> options,
 179        PostgreSqlTransportStore store,
 180        IAsyncResponseIngress ingress,
 181        ILogger<PostgreSqlResponseIngressSubscriber> logger)
 182        : base(options, store, logger)
 183        => _ingress = ingress;
 184
 185    protected override string Queue => Options.ResponseQueue;
 186    protected override PostgreSqlSubscriberOptions SubscriberOptions => Options.ResponseSubscriber;
 187    protected override PostgreSqlSubscriberRole Role => PostgreSqlSubscriberRole.ResponseIngress;
 188
 189    protected override Task HandleMessageAsync(PostgreSqlTransportDelivery delivery, CancellationToken cancellationToken
 190    {
 191        var correlationId = PostgreSqlCorrelationIdExtractor.Extract(delivery.Headers, delivery.Payload, Options);
 192        return _ingress.HandleResponseMessageAsync(delivery.Payload, correlationId);
 193    }
 194}