< Summary

Information
Class: AsiBackbone.AspNetCore.Outbox.GovernanceOutboxDrainHostedService
Assembly: AsiBackbone.AspNetCore
File(s): /home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.AspNetCore/Outbox/GovernanceOutboxDrainHostedService.cs
Line coverage
97%
Covered lines: 131
Uncovered lines: 3
Coverable lines: 134
Total lines: 263
Line coverage: 97.7%
Branch coverage
84%
Covered branches: 22
Total branches: 26
Branch coverage: 84.6%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)62.5%88100%
.cctor()100%11100%
StartAsync(...)100%1181.82%
StopAsync()100%66100%
ExecuteAsync()100%66100%
DrainOnceAsync()100%11100%
WaitForDelayOrOptionsChangeAsync()75%4491.67%
CaptureOptionsChangeState()100%11100%
SignalOptionsChanged()100%11100%
ResolveDrainUtc(...)100%22100%
CreateOptionsChangedSource()100%11100%

File(s)

/home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.AspNetCore/Outbox/GovernanceOutboxDrainHostedService.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 outbox drain from an ASP.NET Core or generic-host background worker.
 11/// </summary>
 12/// <remarks>
 13/// Hosting remains outside Core. Startup validates the scoped drain dependency graph and fails when the store or emitte
 14/// cannot be resolved. Each drain cycle then resolves the drain through a new scoped service provider so durable provid
 15/// that depend on scoped infrastructure, such as a host-owned EF Core <c>DbContext</c>, remain safe to use. Runtime cha
 16/// to <see cref="GovernanceOutboxDrainWorkerOptions.Enabled" /> pause or resume new drain cycles without
 17/// terminating the hosted service.
 18/// </remarks>
 19/// <param name="scopeFactory">The factory used to create a scope for each drain cycle.</param>
 20/// <param name="optionsMonitor">The monitored worker options.</param>
 21/// <param name="logger">The worker logger.</param>
 22/// <param name="timeProvider">The clock used for drain timestamps unless <see cref="GovernanceOutboxDrainWorkerOptions.
 2423public sealed class GovernanceOutboxDrainHostedService(
 2424    IServiceScopeFactory scopeFactory,
 2425    IOptionsMonitor<GovernanceOutboxDrainWorkerOptions> optionsMonitor,
 2426    ILogger<GovernanceOutboxDrainHostedService> logger,
 2427    TimeProvider? timeProvider = null) : BackgroundService
 28{
 129    private static readonly Action<ILogger, Exception?> LogShutdownDrainCanceled = LoggerMessage.Define(
 130        LogLevel.Debug,
 131        new EventId(19801, nameof(LogShutdownDrainCanceled)),
 132        "Governance outbox shutdown drain was canceled.");
 33
 134    private static readonly Action<ILogger, Exception?> LogShutdownDrainFailed = LoggerMessage.Define(
 135        LogLevel.Warning,
 136        new EventId(19802, nameof(LogShutdownDrainFailed)),
 137        "Governance outbox shutdown drain failed.");
 38
 139    private static readonly Action<ILogger, Exception?> LogWorkerDisabled = LoggerMessage.Define(
 140        LogLevel.Debug,
 141        new EventId(19803, nameof(LogWorkerDisabled)),
 142        "Governance outbox drain worker is disabled.");
 43
 144    private static readonly Action<ILogger, int, Exception?> LogDrainAttempted = LoggerMessage.Define<int>(
 145        LogLevel.Debug,
 146        new EventId(19804, nameof(LogDrainAttempted)),
 147        "Governance outbox drain attempted {DrainedCount} entries.");
 48
 149    private static readonly Action<ILogger, Exception?> LogWorkerFailed = LoggerMessage.Define(
 150        LogLevel.Warning,
 151        new EventId(19805, nameof(LogWorkerFailed)),
 152        "Governance outbox drain worker failed before the next polling interval.");
 53
 154    private static readonly Action<ILogger, Exception?> LogStartupValidationFailed = LoggerMessage.Define(
 155        LogLevel.Critical,
 156        new EventId(19806, nameof(LogStartupValidationFailed)),
 157        "Governance outbox drain worker startup validation failed. Ensure an outbox store and governance emitter are reg
 58
 2459    private readonly IServiceScopeFactory scopeFactory = scopeFactory ?? throw new ArgumentNullException(nameof(scopeFac
 2460    private readonly IOptionsMonitor<GovernanceOutboxDrainWorkerOptions> optionsMonitor = optionsMonitor ?? throw new Ar
 2461    private readonly ILogger<GovernanceOutboxDrainHostedService> logger = logger ?? throw new ArgumentNullException(name
 2462    private readonly TimeProvider timeProvider = timeProvider ?? TimeProvider.System;
 2463    private readonly Lock optionsChangedSync = new();
 2464    private TaskCompletionSource optionsChanged = CreateOptionsChangedSource();
 65    private long optionsVersion;
 66
 67    /// <inheritdoc />
 68    public override Task StartAsync(CancellationToken cancellationToken)
 69    {
 70        try
 71        {
 1672            cancellationToken.ThrowIfCancellationRequested();
 1673            optionsMonitor.CurrentValue.Validate();
 74
 1675            using IServiceScope scope = scopeFactory.CreateScope();
 1676            _ = scope.ServiceProvider.GetRequiredService<GovernanceOutboxDrain>();
 1377        }
 078        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 79        {
 080            throw;
 81        }
 382        catch (Exception exception)
 83        {
 384            LogStartupValidationFailed(logger, exception);
 385            throw;
 86        }
 87
 1388        return base.StartAsync(cancellationToken);
 89    }
 90
 91    /// <inheritdoc />
 92    public override async Task StopAsync(CancellationToken cancellationToken)
 93    {
 2094        await base.StopAsync(cancellationToken).ConfigureAwait(false);
 95
 2096        GovernanceOutboxDrainWorkerOptions options = optionsMonitor.CurrentValue;
 97
 2098        if (options.Enabled && options.DrainOnShutdown && !cancellationToken.IsCancellationRequested)
 99        {
 5100            using var shutdownDrainCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 5101            shutdownDrainCancellation.CancelAfter(options.ShutdownDrainTimeout);
 102
 103            try
 104            {
 5105                _ = await DrainOnceAsync(options, shutdownDrainCancellation.Token).ConfigureAwait(false);
 2106            }
 2107            catch (OperationCanceledException) when (shutdownDrainCancellation.IsCancellationRequested)
 108            {
 2109                LogShutdownDrainCanceled(logger, null);
 2110            }
 1111            catch (Exception ex)
 112            {
 1113                LogShutdownDrainFailed(logger, ex);
 1114            }
 5115        }
 20116    }
 117
 118    /// <inheritdoc />
 119    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
 120    {
 121        using IDisposable? optionsChangeRegistration = optionsMonitor.OnChange((_, _) => SignalOptionsChanged());
 12122        bool disabledLogged = false;
 123
 21124        while (!stoppingToken.IsCancellationRequested)
 125        {
 21126            (long observedVersion, Task optionsChangedTask) = CaptureOptionsChangeState();
 21127            GovernanceOutboxDrainWorkerOptions options = optionsMonitor.CurrentValue;
 128
 21129            if (!options.Enabled)
 130            {
 9131                if (!disabledLogged)
 132                {
 5133                    LogWorkerDisabled(logger, null);
 5134                    disabledLogged = true;
 135                }
 136
 9137                await WaitForDelayOrOptionsChangeAsync(
 9138                    options.PollingInterval,
 9139                    observedVersion,
 9140                    optionsChangedTask,
 9141                    stoppingToken).ConfigureAwait(false);
 7142                continue;
 143            }
 144
 12145            disabledLogged = false;
 146
 147            try
 148            {
 12149                int drainedCount = await DrainOnceAsync(options, stoppingToken).ConfigureAwait(false);
 11150                LogDrainAttempted(logger, drainedCount, null);
 11151                await WaitForDelayOrOptionsChangeAsync(
 11152                    options.PollingInterval,
 11153                    observedVersion,
 11154                    optionsChangedTask,
 11155                    stoppingToken).ConfigureAwait(false);
 2156            }
 9157            catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 158            {
 9159                break;
 160            }
 1161            catch (Exception ex)
 162            {
 1163                LogWorkerFailed(logger, ex);
 1164                await WaitForDelayOrOptionsChangeAsync(
 1165                    options.FailureDelay,
 1166                    observedVersion,
 1167                    optionsChangedTask,
 1168                    stoppingToken).ConfigureAwait(false);
 169            }
 2170        }
 9171    }
 172
 173    private async ValueTask<int> DrainOnceAsync(
 174        GovernanceOutboxDrainWorkerOptions options,
 175        CancellationToken cancellationToken)
 176    {
 17177        options.Validate();
 17178        cancellationToken.ThrowIfCancellationRequested();
 179
 17180        using IServiceScope scope = scopeFactory.CreateScope();
 17181        GovernanceOutboxDrain drain = scope.ServiceProvider.GetRequiredService<GovernanceOutboxDrain>();
 17182        DateTimeOffset retryUtc = ResolveDrainUtc(options);
 17183        IReadOnlyList<GovernanceOutboxEntry> drainedEntries = await drain.DrainAsync(
 17184            retryUtc,
 17185            options.BatchSize,
 17186            cancellationToken)
 17187            .ConfigureAwait(false);
 188
 13189        return drainedEntries.Count;
 13190    }
 191
 192    private async ValueTask WaitForDelayOrOptionsChangeAsync(
 193        TimeSpan delay,
 194        long observedVersion,
 195        Task optionsChangedTask,
 196        CancellationToken cancellationToken)
 21197    {
 198        lock (optionsChangedSync)
 199        {
 21200            if (optionsVersion != observedVersion)
 201            {
 0202                return;
 203            }
 21204        }
 205
 21206        using var delayCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 21207        var delayTask = Task.Delay(delay, delayCancellation.Token);
 21208        Task completedTask = await Task.WhenAny(delayTask, optionsChangedTask).ConfigureAwait(false);
 209
 21210        if (completedTask == optionsChangedTask)
 211        {
 3212            await delayCancellation.CancelAsync().ConfigureAwait(false);
 3213            return;
 214        }
 215
 18216        await delayTask.ConfigureAwait(false);
 9217    }
 218
 219    private (long Version, Task ChangedTask) CaptureOptionsChangeState()
 21220    {
 221        lock (optionsChangedSync)
 222        {
 21223            return (optionsVersion, optionsChanged.Task);
 224        }
 21225    }
 226
 227    private void SignalOptionsChanged()
 3228    {
 229        TaskCompletionSource completedSource;
 230        lock (optionsChangedSync)
 231        {
 3232            completedSource = optionsChanged;
 3233            optionsChanged = CreateOptionsChangedSource();
 3234            optionsVersion++;
 3235        }
 236
 3237        _ = completedSource.TrySetResult();
 3238    }
 239
 240    /// <summary>
 241    /// Resolves the drain timestamp for one drain cycle.
 242    /// </summary>
 243    /// <remarks>
 244    /// A custom <see cref="GovernanceOutboxDrainWorkerOptions.RetryClock" /> (obsolete, <c>ASIB903</c>) is honored for 
 245    /// its default value, the registered <see cref="TimeProvider" /> supplies the time, so the worker, the drain, and t
 246    /// signing providers read one clock.
 247    /// </remarks>
 248    private DateTimeOffset ResolveDrainUtc(GovernanceOutboxDrainWorkerOptions options)
 249    {
 250#pragma warning disable ASIB903 // A custom delegate is still honored for compatibility during the 7.x deprecation windo
 17251        DateTimeOffset utcNow = ReferenceEquals(options.RetryClock, GovernanceOutboxDrainWorkerOptions.DefaultRetryClock
 17252            ? timeProvider.GetUtcNow()
 17253            : options.RetryClock();
 254#pragma warning restore ASIB903
 255
 17256        return utcNow.ToUniversalTime();
 257    }
 258
 259    private static TaskCompletionSource CreateOptionsChangedSource()
 260    {
 27261        return new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
 262    }
 263}