< Summary

Information
Class: AsiBackbone.AspNetCore.Outbox.AsiBackboneGovernanceOutboxDrainHostedService
Assembly: AsiBackbone.AspNetCore
File(s): /home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.AspNetCore/Outbox/AsiBackboneGovernanceOutboxDrainHostedService.cs
Line coverage
100%
Covered lines: 114
Uncovered lines: 0
Coverable lines: 114
Total lines: 205
Line coverage: 100%
Branch coverage
86%
Covered branches: 19
Total branches: 22
Branch coverage: 86.3%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)50%66100%
.cctor()100%11100%
StopAsync()100%66100%
ExecuteAsync()100%66100%
DrainOnceAsync()100%11100%
WaitForDelayOrOptionsChangeAsync()100%44100%
CaptureOptionsChangeState()100%11100%
SignalOptionsChanged()100%11100%
CreateOptionsChangedSource()100%11100%

File(s)

/home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.AspNetCore/Outbox/AsiBackboneGovernanceOutboxDrainHostedService.cs

#LineLine coverage
 1using AsiBackbone.Core.Outbox;
 2using Microsoft.Extensions.DependencyInjection;
 3using Microsoft.Extensions.Hosting;
 4using Microsoft.Extensions.Logging;
 5using Microsoft.Extensions.Options;
 6
 7namespace AsiBackbone.AspNetCore.Outbox;
 8
 9/// <summary>
 10/// Runs the provider-neutral governance outbox drain from an ASP.NET Core or generic-host background worker.
 11/// </summary>
 12/// <remarks>
 13/// Hosting remains outside Core. The worker resolves the drain through a scoped service provider so durable providers t
 14/// </remarks>
 2015public sealed class AsiBackboneGovernanceOutboxDrainHostedService(
 2016    IServiceScopeFactory scopeFactory,
 2017    IOptionsMonitor<AsiBackboneGovernanceOutboxDrainWorkerOptions> optionsMonitor,
 2018    ILogger<AsiBackboneGovernanceOutboxDrainHostedService> logger) : BackgroundService
 19{
 120    private static readonly Action<ILogger, Exception?> LogShutdownDrainCanceled = LoggerMessage.Define(
 121        LogLevel.Debug,
 122        new EventId(19801, nameof(LogShutdownDrainCanceled)),
 123        "Governance outbox shutdown drain was canceled.");
 24
 125    private static readonly Action<ILogger, Exception?> LogShutdownDrainFailed = LoggerMessage.Define(
 126        LogLevel.Warning,
 127        new EventId(19802, nameof(LogShutdownDrainFailed)),
 128        "Governance outbox shutdown drain failed.");
 29
 130    private static readonly Action<ILogger, Exception?> LogWorkerDisabled = LoggerMessage.Define(
 131        LogLevel.Debug,
 132        new EventId(19803, nameof(LogWorkerDisabled)),
 133        "Governance outbox drain worker is disabled.");
 34
 135    private static readonly Action<ILogger, int, Exception?> LogDrainAttempted = LoggerMessage.Define<int>(
 136        LogLevel.Debug,
 137        new EventId(19804, nameof(LogDrainAttempted)),
 138        "Governance outbox drain attempted {DrainedCount} entries.");
 39
 140    private static readonly Action<ILogger, Exception?> LogWorkerFailed = LoggerMessage.Define(
 141        LogLevel.Warning,
 142        new EventId(19805, nameof(LogWorkerFailed)),
 143        "Governance outbox drain worker failed before the next polling interval.");
 44
 2045    private readonly IServiceScopeFactory scopeFactory = scopeFactory ?? throw new ArgumentNullException(nameof(scopeFac
 2046    private readonly IOptionsMonitor<AsiBackboneGovernanceOutboxDrainWorkerOptions> optionsMonitor = optionsMonitor ?? t
 2047    private readonly ILogger<AsiBackboneGovernanceOutboxDrainHostedService> logger = logger ?? throw new ArgumentNullExc
 2048    private readonly Lock optionsChangedSync = new();
 2049    private TaskCompletionSource optionsChanged = CreateOptionsChangedSource();
 50    private long optionsVersion;
 51
 52    /// <inheritdoc />
 53    public override async Task StopAsync(CancellationToken cancellationToken)
 54    {
 1955        await base.StopAsync(cancellationToken).ConfigureAwait(false);
 56
 1957        AsiBackboneGovernanceOutboxDrainWorkerOptions options = optionsMonitor.CurrentValue;
 58
 1959        if (options.Enabled && options.DrainOnShutdown && !cancellationToken.IsCancellationRequested)
 60        {
 561            using var shutdownDrainCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 562            shutdownDrainCancellation.CancelAfter(options.ShutdownDrainTimeout);
 63
 64            try
 65            {
 566                _ = await DrainOnceAsync(options, shutdownDrainCancellation.Token).ConfigureAwait(false);
 267            }
 268            catch (OperationCanceledException) when (shutdownDrainCancellation.IsCancellationRequested)
 69            {
 270                LogShutdownDrainCanceled(logger, null);
 271            }
 172            catch (Exception ex)
 73            {
 174                LogShutdownDrainFailed(logger, ex);
 175            }
 576        }
 1977    }
 78
 79    /// <inheritdoc />
 80    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
 81    {
 1482        using IDisposable? optionsChangeRegistration = optionsMonitor.OnChange((_, _) => SignalOptionsChanged());
 1183        bool disabledLogged = false;
 84
 2185        while (!stoppingToken.IsCancellationRequested)
 86        {
 2187            (long observedVersion, Task optionsChangedTask) = CaptureOptionsChangeState();
 2188            AsiBackboneGovernanceOutboxDrainWorkerOptions options = optionsMonitor.CurrentValue;
 89
 2190            if (!options.Enabled)
 91            {
 992                if (!disabledLogged)
 93                {
 594                    LogWorkerDisabled(logger, null);
 595                    disabledLogged = true;
 96                }
 97
 998                await WaitForDelayOrOptionsChangeAsync(
 999                    options.PollingInterval,
 9100                    observedVersion,
 9101                    optionsChangedTask,
 9102                    stoppingToken).ConfigureAwait(false);
 7103                continue;
 104            }
 105
 12106            disabledLogged = false;
 107
 108            try
 109            {
 12110                int drainedCount = await DrainOnceAsync(options, stoppingToken).ConfigureAwait(false);
 9111                LogDrainAttempted(logger, drainedCount, null);
 9112                await WaitForDelayOrOptionsChangeAsync(
 9113                    options.PollingInterval,
 9114                    observedVersion,
 9115                    optionsChangedTask,
 9116                    stoppingToken).ConfigureAwait(false);
 2117            }
 7118            catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 119            {
 7120                break;
 121            }
 3122            catch (Exception ex)
 123            {
 3124                LogWorkerFailed(logger, ex);
 3125                await WaitForDelayOrOptionsChangeAsync(
 3126                    options.FailureDelay,
 3127                    observedVersion,
 3128                    optionsChangedTask,
 3129                    stoppingToken).ConfigureAwait(false);
 130            }
 3131        }
 7132    }
 133
 134    private async ValueTask<int> DrainOnceAsync(
 135        AsiBackboneGovernanceOutboxDrainWorkerOptions options,
 136        CancellationToken cancellationToken)
 137    {
 17138        options.Validate();
 17139        cancellationToken.ThrowIfCancellationRequested();
 140
 17141        using IServiceScope scope = scopeFactory.CreateScope();
 17142        AsiBackboneGovernanceOutboxDrain drain = scope.ServiceProvider.GetRequiredService<AsiBackboneGovernanceOutboxDra
 14143        DateTimeOffset retryUtc = options.RetryClock().ToUniversalTime();
 14144        IReadOnlyList<GovernanceOutboxEntry> drainedEntries = await drain.DrainAsync(
 14145            retryUtc,
 14146            options.BatchSize,
 14147            cancellationToken)
 14148            .ConfigureAwait(false);
 149
 11150        return drainedEntries.Count;
 11151    }
 152
 153    private async ValueTask WaitForDelayOrOptionsChangeAsync(
 154        TimeSpan delay,
 155        long observedVersion,
 156        Task optionsChangedTask,
 157        CancellationToken cancellationToken)
 21158    {
 159        lock (optionsChangedSync)
 160        {
 21161            if (optionsVersion != observedVersion)
 162            {
 1163                return;
 164            }
 20165        }
 166
 20167        using var delayCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 20168        var delayTask = Task.Delay(delay, delayCancellation.Token);
 20169        Task completedTask = await Task.WhenAny(delayTask, optionsChangedTask).ConfigureAwait(false);
 170
 20171        if (completedTask == optionsChangedTask)
 172        {
 2173            await delayCancellation.CancelAsync().ConfigureAwait(false);
 2174            return;
 175        }
 176
 18177        await delayTask.ConfigureAwait(false);
 10178    }
 179
 180    private (long Version, Task ChangedTask) CaptureOptionsChangeState()
 21181    {
 182        lock (optionsChangedSync)
 183        {
 21184            return (optionsVersion, optionsChanged.Task);
 185        }
 21186    }
 187
 188    private void SignalOptionsChanged()
 3189    {
 190        TaskCompletionSource completedSource;
 191        lock (optionsChangedSync)
 192        {
 3193            completedSource = optionsChanged;
 3194            optionsChanged = CreateOptionsChangedSource();
 3195            optionsVersion++;
 3196        }
 197
 3198        _ = completedSource.TrySetResult();
 3199    }
 200
 201    private static TaskCompletionSource CreateOptionsChangedSource()
 202    {
 23203        return new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
 204    }
 205}