| | | 1 | | using AsiBackbone.Core.Outbox; |
| | | 2 | | using Microsoft.Extensions.DependencyInjection; |
| | | 3 | | using Microsoft.Extensions.Hosting; |
| | | 4 | | using Microsoft.Extensions.Logging; |
| | | 5 | | using Microsoft.Extensions.Options; |
| | | 6 | | |
| | | 7 | | namespace 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. |
| | 24 | 23 | | public sealed class GovernanceOutboxDrainHostedService( |
| | 24 | 24 | | IServiceScopeFactory scopeFactory, |
| | 24 | 25 | | IOptionsMonitor<GovernanceOutboxDrainWorkerOptions> optionsMonitor, |
| | 24 | 26 | | ILogger<GovernanceOutboxDrainHostedService> logger, |
| | 24 | 27 | | TimeProvider? timeProvider = null) : BackgroundService |
| | | 28 | | { |
| | 1 | 29 | | private static readonly Action<ILogger, Exception?> LogShutdownDrainCanceled = LoggerMessage.Define( |
| | 1 | 30 | | LogLevel.Debug, |
| | 1 | 31 | | new EventId(19801, nameof(LogShutdownDrainCanceled)), |
| | 1 | 32 | | "Governance outbox shutdown drain was canceled."); |
| | | 33 | | |
| | 1 | 34 | | private static readonly Action<ILogger, Exception?> LogShutdownDrainFailed = LoggerMessage.Define( |
| | 1 | 35 | | LogLevel.Warning, |
| | 1 | 36 | | new EventId(19802, nameof(LogShutdownDrainFailed)), |
| | 1 | 37 | | "Governance outbox shutdown drain failed."); |
| | | 38 | | |
| | 1 | 39 | | private static readonly Action<ILogger, Exception?> LogWorkerDisabled = LoggerMessage.Define( |
| | 1 | 40 | | LogLevel.Debug, |
| | 1 | 41 | | new EventId(19803, nameof(LogWorkerDisabled)), |
| | 1 | 42 | | "Governance outbox drain worker is disabled."); |
| | | 43 | | |
| | 1 | 44 | | private static readonly Action<ILogger, int, Exception?> LogDrainAttempted = LoggerMessage.Define<int>( |
| | 1 | 45 | | LogLevel.Debug, |
| | 1 | 46 | | new EventId(19804, nameof(LogDrainAttempted)), |
| | 1 | 47 | | "Governance outbox drain attempted {DrainedCount} entries."); |
| | | 48 | | |
| | 1 | 49 | | private static readonly Action<ILogger, Exception?> LogWorkerFailed = LoggerMessage.Define( |
| | 1 | 50 | | LogLevel.Warning, |
| | 1 | 51 | | new EventId(19805, nameof(LogWorkerFailed)), |
| | 1 | 52 | | "Governance outbox drain worker failed before the next polling interval."); |
| | | 53 | | |
| | 1 | 54 | | private static readonly Action<ILogger, Exception?> LogStartupValidationFailed = LoggerMessage.Define( |
| | 1 | 55 | | LogLevel.Critical, |
| | 1 | 56 | | new EventId(19806, nameof(LogStartupValidationFailed)), |
| | 1 | 57 | | "Governance outbox drain worker startup validation failed. Ensure an outbox store and governance emitter are reg |
| | | 58 | | |
| | 24 | 59 | | private readonly IServiceScopeFactory scopeFactory = scopeFactory ?? throw new ArgumentNullException(nameof(scopeFac |
| | 24 | 60 | | private readonly IOptionsMonitor<GovernanceOutboxDrainWorkerOptions> optionsMonitor = optionsMonitor ?? throw new Ar |
| | 24 | 61 | | private readonly ILogger<GovernanceOutboxDrainHostedService> logger = logger ?? throw new ArgumentNullException(name |
| | 24 | 62 | | private readonly TimeProvider timeProvider = timeProvider ?? TimeProvider.System; |
| | 24 | 63 | | private readonly Lock optionsChangedSync = new(); |
| | 24 | 64 | | private TaskCompletionSource optionsChanged = CreateOptionsChangedSource(); |
| | | 65 | | private long optionsVersion; |
| | | 66 | | |
| | | 67 | | /// <inheritdoc /> |
| | | 68 | | public override Task StartAsync(CancellationToken cancellationToken) |
| | | 69 | | { |
| | | 70 | | try |
| | | 71 | | { |
| | 16 | 72 | | cancellationToken.ThrowIfCancellationRequested(); |
| | 16 | 73 | | optionsMonitor.CurrentValue.Validate(); |
| | | 74 | | |
| | 16 | 75 | | using IServiceScope scope = scopeFactory.CreateScope(); |
| | 16 | 76 | | _ = scope.ServiceProvider.GetRequiredService<GovernanceOutboxDrain>(); |
| | 13 | 77 | | } |
| | 0 | 78 | | catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) |
| | | 79 | | { |
| | 0 | 80 | | throw; |
| | | 81 | | } |
| | 3 | 82 | | catch (Exception exception) |
| | | 83 | | { |
| | 3 | 84 | | LogStartupValidationFailed(logger, exception); |
| | 3 | 85 | | throw; |
| | | 86 | | } |
| | | 87 | | |
| | 13 | 88 | | return base.StartAsync(cancellationToken); |
| | | 89 | | } |
| | | 90 | | |
| | | 91 | | /// <inheritdoc /> |
| | | 92 | | public override async Task StopAsync(CancellationToken cancellationToken) |
| | | 93 | | { |
| | 20 | 94 | | await base.StopAsync(cancellationToken).ConfigureAwait(false); |
| | | 95 | | |
| | 20 | 96 | | GovernanceOutboxDrainWorkerOptions options = optionsMonitor.CurrentValue; |
| | | 97 | | |
| | 20 | 98 | | if (options.Enabled && options.DrainOnShutdown && !cancellationToken.IsCancellationRequested) |
| | | 99 | | { |
| | 5 | 100 | | using var shutdownDrainCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); |
| | 5 | 101 | | shutdownDrainCancellation.CancelAfter(options.ShutdownDrainTimeout); |
| | | 102 | | |
| | | 103 | | try |
| | | 104 | | { |
| | 5 | 105 | | _ = await DrainOnceAsync(options, shutdownDrainCancellation.Token).ConfigureAwait(false); |
| | 2 | 106 | | } |
| | 2 | 107 | | catch (OperationCanceledException) when (shutdownDrainCancellation.IsCancellationRequested) |
| | | 108 | | { |
| | 2 | 109 | | LogShutdownDrainCanceled(logger, null); |
| | 2 | 110 | | } |
| | 1 | 111 | | catch (Exception ex) |
| | | 112 | | { |
| | 1 | 113 | | LogShutdownDrainFailed(logger, ex); |
| | 1 | 114 | | } |
| | 5 | 115 | | } |
| | 20 | 116 | | } |
| | | 117 | | |
| | | 118 | | /// <inheritdoc /> |
| | | 119 | | protected override async Task ExecuteAsync(CancellationToken stoppingToken) |
| | | 120 | | { |
| | | 121 | | using IDisposable? optionsChangeRegistration = optionsMonitor.OnChange((_, _) => SignalOptionsChanged()); |
| | 12 | 122 | | bool disabledLogged = false; |
| | | 123 | | |
| | 21 | 124 | | while (!stoppingToken.IsCancellationRequested) |
| | | 125 | | { |
| | 21 | 126 | | (long observedVersion, Task optionsChangedTask) = CaptureOptionsChangeState(); |
| | 21 | 127 | | GovernanceOutboxDrainWorkerOptions options = optionsMonitor.CurrentValue; |
| | | 128 | | |
| | 21 | 129 | | if (!options.Enabled) |
| | | 130 | | { |
| | 9 | 131 | | if (!disabledLogged) |
| | | 132 | | { |
| | 5 | 133 | | LogWorkerDisabled(logger, null); |
| | 5 | 134 | | disabledLogged = true; |
| | | 135 | | } |
| | | 136 | | |
| | 9 | 137 | | await WaitForDelayOrOptionsChangeAsync( |
| | 9 | 138 | | options.PollingInterval, |
| | 9 | 139 | | observedVersion, |
| | 9 | 140 | | optionsChangedTask, |
| | 9 | 141 | | stoppingToken).ConfigureAwait(false); |
| | 7 | 142 | | continue; |
| | | 143 | | } |
| | | 144 | | |
| | 12 | 145 | | disabledLogged = false; |
| | | 146 | | |
| | | 147 | | try |
| | | 148 | | { |
| | 12 | 149 | | int drainedCount = await DrainOnceAsync(options, stoppingToken).ConfigureAwait(false); |
| | 11 | 150 | | LogDrainAttempted(logger, drainedCount, null); |
| | 11 | 151 | | await WaitForDelayOrOptionsChangeAsync( |
| | 11 | 152 | | options.PollingInterval, |
| | 11 | 153 | | observedVersion, |
| | 11 | 154 | | optionsChangedTask, |
| | 11 | 155 | | stoppingToken).ConfigureAwait(false); |
| | 2 | 156 | | } |
| | 9 | 157 | | catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) |
| | | 158 | | { |
| | 9 | 159 | | break; |
| | | 160 | | } |
| | 1 | 161 | | catch (Exception ex) |
| | | 162 | | { |
| | 1 | 163 | | LogWorkerFailed(logger, ex); |
| | 1 | 164 | | await WaitForDelayOrOptionsChangeAsync( |
| | 1 | 165 | | options.FailureDelay, |
| | 1 | 166 | | observedVersion, |
| | 1 | 167 | | optionsChangedTask, |
| | 1 | 168 | | stoppingToken).ConfigureAwait(false); |
| | | 169 | | } |
| | 2 | 170 | | } |
| | 9 | 171 | | } |
| | | 172 | | |
| | | 173 | | private async ValueTask<int> DrainOnceAsync( |
| | | 174 | | GovernanceOutboxDrainWorkerOptions options, |
| | | 175 | | CancellationToken cancellationToken) |
| | | 176 | | { |
| | 17 | 177 | | options.Validate(); |
| | 17 | 178 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 179 | | |
| | 17 | 180 | | using IServiceScope scope = scopeFactory.CreateScope(); |
| | 17 | 181 | | GovernanceOutboxDrain drain = scope.ServiceProvider.GetRequiredService<GovernanceOutboxDrain>(); |
| | 17 | 182 | | DateTimeOffset retryUtc = ResolveDrainUtc(options); |
| | 17 | 183 | | IReadOnlyList<GovernanceOutboxEntry> drainedEntries = await drain.DrainAsync( |
| | 17 | 184 | | retryUtc, |
| | 17 | 185 | | options.BatchSize, |
| | 17 | 186 | | cancellationToken) |
| | 17 | 187 | | .ConfigureAwait(false); |
| | | 188 | | |
| | 13 | 189 | | return drainedEntries.Count; |
| | 13 | 190 | | } |
| | | 191 | | |
| | | 192 | | private async ValueTask WaitForDelayOrOptionsChangeAsync( |
| | | 193 | | TimeSpan delay, |
| | | 194 | | long observedVersion, |
| | | 195 | | Task optionsChangedTask, |
| | | 196 | | CancellationToken cancellationToken) |
| | 21 | 197 | | { |
| | | 198 | | lock (optionsChangedSync) |
| | | 199 | | { |
| | 21 | 200 | | if (optionsVersion != observedVersion) |
| | | 201 | | { |
| | 0 | 202 | | return; |
| | | 203 | | } |
| | 21 | 204 | | } |
| | | 205 | | |
| | 21 | 206 | | using var delayCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); |
| | 21 | 207 | | var delayTask = Task.Delay(delay, delayCancellation.Token); |
| | 21 | 208 | | Task completedTask = await Task.WhenAny(delayTask, optionsChangedTask).ConfigureAwait(false); |
| | | 209 | | |
| | 21 | 210 | | if (completedTask == optionsChangedTask) |
| | | 211 | | { |
| | 3 | 212 | | await delayCancellation.CancelAsync().ConfigureAwait(false); |
| | 3 | 213 | | return; |
| | | 214 | | } |
| | | 215 | | |
| | 18 | 216 | | await delayTask.ConfigureAwait(false); |
| | 9 | 217 | | } |
| | | 218 | | |
| | | 219 | | private (long Version, Task ChangedTask) CaptureOptionsChangeState() |
| | 21 | 220 | | { |
| | | 221 | | lock (optionsChangedSync) |
| | | 222 | | { |
| | 21 | 223 | | return (optionsVersion, optionsChanged.Task); |
| | | 224 | | } |
| | 21 | 225 | | } |
| | | 226 | | |
| | | 227 | | private void SignalOptionsChanged() |
| | 3 | 228 | | { |
| | | 229 | | TaskCompletionSource completedSource; |
| | | 230 | | lock (optionsChangedSync) |
| | | 231 | | { |
| | 3 | 232 | | completedSource = optionsChanged; |
| | 3 | 233 | | optionsChanged = CreateOptionsChangedSource(); |
| | 3 | 234 | | optionsVersion++; |
| | 3 | 235 | | } |
| | | 236 | | |
| | 3 | 237 | | _ = completedSource.TrySetResult(); |
| | 3 | 238 | | } |
| | | 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 |
| | 17 | 251 | | DateTimeOffset utcNow = ReferenceEquals(options.RetryClock, GovernanceOutboxDrainWorkerOptions.DefaultRetryClock |
| | 17 | 252 | | ? timeProvider.GetUtcNow() |
| | 17 | 253 | | : options.RetryClock(); |
| | | 254 | | #pragma warning restore ASIB903 |
| | | 255 | | |
| | 17 | 256 | | return utcNow.ToUniversalTime(); |
| | | 257 | | } |
| | | 258 | | |
| | | 259 | | private static TaskCompletionSource CreateOptionsChangedSource() |
| | | 260 | | { |
| | 27 | 261 | | return new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); |
| | | 262 | | } |
| | | 263 | | } |