| | | 1 | | using System.Collections.Concurrent; |
| | | 2 | | using System.Diagnostics; |
| | | 3 | | using Elsa.Common; |
| | | 4 | | using Elsa.Workflows.Management; |
| | | 5 | | using Elsa.Workflows.Management.Entities; |
| | | 6 | | using Elsa.Workflows.Management.Filters; |
| | | 7 | | using Elsa.Workflows.Runtime.Options; |
| | | 8 | | using Microsoft.Extensions.DependencyInjection; |
| | | 9 | | using Microsoft.Extensions.Hosting; |
| | | 10 | | using Microsoft.Extensions.Logging; |
| | | 11 | | using Microsoft.Extensions.Options; |
| | | 12 | | |
| | | 13 | | namespace Elsa.Workflows.Runtime.Services; |
| | | 14 | | |
| | | 15 | | /// <summary> |
| | | 16 | | /// Default implementation of <see cref="IDrainOrchestrator"/>. |
| | | 17 | | /// </summary> |
| | | 18 | | /// <remarks> |
| | | 19 | | /// <para> |
| | | 20 | | /// Protocol per <c>contracts/drain-orchestrator.md</c>: |
| | | 21 | | /// 1. Enter drain via <see cref="IQuiescenceSignal.BeginDrainAsync"/>. |
| | | 22 | | /// 2. Pause every registered ingress source in parallel, each bounded by its own per-source timeout. Sources that |
| | | 23 | | /// fail to pause are recorded as <see cref="IngressSourceState.PauseFailed"/> and, if they implement |
| | | 24 | | /// <see cref="IForceStoppable"/>, are escalated. |
| | | 25 | | /// 3. Wait for <see cref="IExecutionCycleRegistry.ActiveCount"/> to reach zero, polling on a short interval. |
| | | 26 | | /// 4. On deadline breach (or on operator force, where deadline is zero), iterate live cycles, cancel each, |
| | | 27 | | /// persist still-Running instances as <see cref="WorkflowSubStatus.Interrupted"/>, leave |
| | | 28 | | /// Finished/Cancelled rows as they are, and write a <c>WorkflowInterrupted</c> entry in the |
| | | 29 | | /// per-instance execution log. |
| | | 30 | | /// 5. Return a <see cref="DrainOutcome"/>. |
| | | 31 | | /// </para> |
| | | 32 | | /// <para> |
| | | 33 | | /// All exceptions inside the protocol are captured into the outcome — only an <see cref="InvalidOperationException"/> |
| | | 34 | | /// thrown by a second non-force invocation propagates. |
| | | 35 | | /// </para> |
| | | 36 | | /// </remarks> |
| | | 37 | | public sealed class DrainOrchestrator : IDrainOrchestrator |
| | | 38 | | { |
| | 3 | 39 | | private static readonly TimeSpan SafetyEpsilon = TimeSpan.FromMilliseconds(500); |
| | 3 | 40 | | private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(10); |
| | | 41 | | |
| | | 42 | | private readonly IQuiescenceSignal _signal; |
| | | 43 | | private readonly IIngressSourceRegistry _registry; |
| | | 44 | | private readonly IExecutionCycleRegistry _cycles; |
| | | 45 | | private readonly IServiceScopeFactory _scopeFactory; |
| | | 46 | | private readonly IOptions<GracefulShutdownOptions> _options; |
| | | 47 | | private readonly IOptions<HostOptions> _hostOptions; |
| | | 48 | | private readonly ISystemClock _clock; |
| | | 49 | | private readonly IIdentityGenerator _identityGenerator; |
| | | 50 | | private readonly ILogger<DrainOrchestrator> _logger; |
| | | 51 | | |
| | 45 | 52 | | private readonly object _sync = new(); |
| | | 53 | | private DrainOutcome? _previousOutcome; |
| | | 54 | | private bool _drainInProgress; |
| | | 55 | | |
| | 45 | 56 | | public DrainOrchestrator( |
| | 45 | 57 | | IQuiescenceSignal signal, |
| | 45 | 58 | | IIngressSourceRegistry registry, |
| | 45 | 59 | | IExecutionCycleRegistry cycles, |
| | 45 | 60 | | IServiceScopeFactory scopeFactory, |
| | 45 | 61 | | IOptions<GracefulShutdownOptions> options, |
| | 45 | 62 | | IOptions<HostOptions> hostOptions, |
| | 45 | 63 | | ISystemClock clock, |
| | 45 | 64 | | IIdentityGenerator identityGenerator, |
| | 45 | 65 | | ILogger<DrainOrchestrator> logger) |
| | | 66 | | { |
| | 45 | 67 | | _signal = signal; |
| | 45 | 68 | | _registry = registry; |
| | 45 | 69 | | _cycles = cycles; |
| | 45 | 70 | | _scopeFactory = scopeFactory; |
| | 45 | 71 | | _options = options; |
| | 45 | 72 | | _hostOptions = hostOptions; |
| | 45 | 73 | | _clock = clock; |
| | 45 | 74 | | _identityGenerator = identityGenerator; |
| | 45 | 75 | | _logger = logger; |
| | 45 | 76 | | } |
| | | 77 | | |
| | | 78 | | /// <inheritdoc /> |
| | | 79 | | public async ValueTask<DrainOutcome> DrainAsync(DrainTrigger trigger, CancellationToken cancellationToken = default) |
| | | 80 | | { |
| | 41 | 81 | | lock (_sync) |
| | | 82 | | { |
| | 41 | 83 | | if (_drainInProgress) |
| | | 84 | | { |
| | 3 | 85 | | if (trigger == DrainTrigger.OperatorForce && _previousOutcome is not null) |
| | 0 | 86 | | return _previousOutcome with { WasCached = true }; |
| | 3 | 87 | | throw new InvalidOperationException($"Drain already in progress; second invocation rejected (trigger={tr |
| | | 88 | | } |
| | 38 | 89 | | if (_previousOutcome is not null) |
| | | 90 | | { |
| | 3 | 91 | | if (trigger == DrainTrigger.OperatorForce) return _previousOutcome with { WasCached = true }; |
| | 1 | 92 | | throw new InvalidOperationException("Drain already completed in this generation; subsequent non-force in |
| | | 93 | | } |
| | 36 | 94 | | _drainInProgress = true; |
| | 36 | 95 | | } |
| | | 96 | | |
| | 36 | 97 | | var startedAt = _clock.UtcNow; |
| | 36 | 98 | | var deadline = ComputeEffectiveDeadline(trigger); |
| | 36 | 99 | | var sw = Stopwatch.StartNew(); |
| | 36 | 100 | | TimeSpan pausePhase = TimeSpan.Zero; |
| | 36 | 101 | | TimeSpan waitPhase = TimeSpan.Zero; |
| | | 102 | | |
| | | 103 | | try |
| | | 104 | | { |
| | 36 | 105 | | await _signal.BeginDrainAsync(cancellationToken); |
| | 36 | 106 | | _logger.LogInformation("Drain initiated (trigger={Trigger}, deadline={Deadline}).", trigger, deadline); |
| | | 107 | | |
| | 36 | 108 | | var deadlineAt = startedAt + deadline; |
| | | 109 | | |
| | | 110 | | // Phase 1: parallel pause. Each source has its own timeout independent of the overall deadline, |
| | | 111 | | // but force-stop fallbacks are clamped to the *remaining* overall budget (deadlineAt - now), not |
| | | 112 | | // the full TimeSpan, so a slow per-source pause can't extend the total shutdown window. |
| | 36 | 113 | | await PauseAllSourcesAsync(deadlineAt, cancellationToken); |
| | 36 | 114 | | pausePhase = sw.Elapsed; |
| | 36 | 115 | | _logger.LogInformation("Ingress pause phase complete in {Elapsed}.", pausePhase); |
| | | 116 | | |
| | | 117 | | // Phase 2: wait for active execution cycles to drain. |
| | 36 | 118 | | var waitStart = sw.Elapsed; |
| | 36 | 119 | | var deadlineBreach = !await WaitForCyclesAsync(deadlineAt, cancellationToken); |
| | 36 | 120 | | waitPhase = sw.Elapsed - waitStart; |
| | | 121 | | |
| | | 122 | | // Phase 3: outcome assembly. |
| | 36 | 123 | | int forceCancelled = 0; |
| | 36 | 124 | | IReadOnlyList<string> forceCancelledIds = Array.Empty<string>(); |
| | | 125 | | |
| | 36 | 126 | | if (deadlineBreach || trigger == DrainTrigger.OperatorForce) |
| | | 127 | | { |
| | 19 | 128 | | _logger.LogWarning("Drain {What}; force-cancelling {Count} active execution cycle(s).", |
| | 19 | 129 | | deadlineBreach ? "deadline exceeded" : "operator-forced", _cycles.ActiveCount); |
| | 19 | 130 | | (forceCancelled, forceCancelledIds) = await ForceCancelActiveCyclesAsync(trigger, _signal.CurrentState.G |
| | | 131 | | } |
| | | 132 | | |
| | | 133 | | // OperatorForce always reports Forced regardless of zero-deadline breach mechanics. |
| | 34 | 134 | | var result = trigger switch |
| | 34 | 135 | | { |
| | 14 | 136 | | DrainTrigger.OperatorForce => DrainResult.Forced, |
| | 23 | 137 | | _ when deadlineBreach => DrainResult.DeadlineExceeded, |
| | 17 | 138 | | _ => DrainResult.CompletedWithinDeadline, |
| | 34 | 139 | | }; |
| | | 140 | | |
| | 34 | 141 | | var outcome = new DrainOutcome( |
| | 34 | 142 | | OverallResult: result, |
| | 34 | 143 | | StartedAt: startedAt, |
| | 34 | 144 | | CompletedAt: _clock.UtcNow, |
| | 34 | 145 | | PausePhaseDuration: pausePhase, |
| | 34 | 146 | | WaitPhaseDuration: waitPhase, |
| | 34 | 147 | | Sources: BuildSourceFinalStates(), |
| | 34 | 148 | | ExecutionCyclesForceCancelledCount: forceCancelled, |
| | 34 | 149 | | ForceCancelledInstanceIds: forceCancelledIds); |
| | | 150 | | |
| | 68 | 151 | | lock (_sync) _previousOutcome = outcome; |
| | 34 | 152 | | _logger.LogInformation("Drain completed: {Result} (paused={Paused}, waited={Waited}, forceCancelled={ForceCa |
| | 34 | 153 | | outcome.OverallResult, outcome.PausePhaseDuration, outcome.WaitPhaseDuration, outcome.ExecutionCyclesFor |
| | 34 | 154 | | return outcome; |
| | | 155 | | } |
| | 2 | 156 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 157 | | { |
| | | 158 | | // The "drain already in progress / completed" InvalidOperationExceptions are thrown OUTSIDE this |
| | | 159 | | // try block (during the lock-protected setup), so they bubble out without triggering this handler. |
| | | 160 | | // Any IOE that lands here is incidental — from a store/dispatcher inside the protocol — and should |
| | | 161 | | // be captured into the outcome rather than escaping the whole drain. |
| | 1 | 162 | | _logger.LogError(ex, "Drain aborted by unhandled exception."); |
| | 1 | 163 | | var outcome = new DrainOutcome( |
| | 1 | 164 | | OverallResult: DrainResult.AbortedByUnhandledException, |
| | 1 | 165 | | StartedAt: startedAt, |
| | 1 | 166 | | CompletedAt: _clock.UtcNow, |
| | 1 | 167 | | PausePhaseDuration: pausePhase, |
| | 1 | 168 | | WaitPhaseDuration: waitPhase, |
| | 1 | 169 | | Sources: BuildSourceFinalStates(), |
| | 1 | 170 | | ExecutionCyclesForceCancelledCount: 0, |
| | 1 | 171 | | ForceCancelledInstanceIds: Array.Empty<string>()); |
| | 2 | 172 | | lock (_sync) _previousOutcome = outcome; |
| | 1 | 173 | | return outcome; |
| | | 174 | | } |
| | | 175 | | finally |
| | | 176 | | { |
| | 72 | 177 | | lock (_sync) _drainInProgress = false; |
| | | 178 | | } |
| | 36 | 179 | | } |
| | | 180 | | |
| | | 181 | | private TimeSpan ComputeEffectiveDeadline(DrainTrigger trigger) |
| | | 182 | | { |
| | 51 | 183 | | if (trigger == DrainTrigger.OperatorForce) return TimeSpan.Zero; |
| | | 184 | | |
| | 21 | 185 | | var configured = _options.Value.DrainDeadline; |
| | 21 | 186 | | var hostShutdown = _hostOptions.Value.ShutdownTimeout; |
| | | 187 | | |
| | | 188 | | // Clamp to host's own shutdown budget minus a safety epsilon, so persistence can finish before the process is k |
| | 21 | 189 | | if (hostShutdown > SafetyEpsilon) |
| | | 190 | | { |
| | 21 | 191 | | var hostBudget = hostShutdown - SafetyEpsilon; |
| | 28 | 192 | | if (hostBudget < configured) return hostBudget; |
| | | 193 | | } |
| | | 194 | | |
| | 14 | 195 | | return configured; |
| | | 196 | | } |
| | | 197 | | |
| | | 198 | | private async Task PauseAllSourcesAsync(DateTimeOffset deadlineAt, CancellationToken cancellationToken) |
| | | 199 | | { |
| | 36 | 200 | | var sources = _registry.Sources; |
| | 55 | 201 | | if (sources.Count == 0) return; |
| | | 202 | | |
| | 47 | 203 | | var pauseTasks = sources.Select(source => PauseOneSourceAsync(source, deadlineAt, cancellationToken)).ToArray(); |
| | 17 | 204 | | await Task.WhenAll(pauseTasks); |
| | 36 | 205 | | } |
| | | 206 | | |
| | | 207 | | private async Task PauseOneSourceAsync(IIngressSource source, DateTimeOffset deadlineAt, CancellationToken cancellat |
| | | 208 | | { |
| | | 209 | | // Precedence: a positive per-source PauseTimeout wins; Zero (or negative) defers to the |
| | | 210 | | // configured GracefulShutdownOptions.IngressPauseTimeout. The resolved value is then |
| | | 211 | | // clamped to the overall remaining budget, with a 1 ms floor so a misconfigured zero |
| | | 212 | | // configured-default still produces a non-zero CancelAfter. |
| | 30 | 213 | | var configured = _options.Value.IngressPauseTimeout; |
| | 30 | 214 | | var sourceDeadline = source.PauseTimeout > TimeSpan.Zero ? source.PauseTimeout : configured; |
| | 30 | 215 | | var overallRemaining = deadlineAt - _clock.UtcNow; |
| | 38 | 216 | | if (sourceDeadline > overallRemaining) sourceDeadline = overallRemaining; |
| | 30 | 217 | | if (sourceDeadline <= TimeSpan.Zero) sourceDeadline = TimeSpan.FromMilliseconds(1); |
| | | 218 | | |
| | 30 | 219 | | _registry.RecordTransition(source.Name, IngressSourceState.Pausing); |
| | 30 | 220 | | using var perSourceCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); |
| | 30 | 221 | | perSourceCts.CancelAfter(sourceDeadline); |
| | | 222 | | |
| | | 223 | | try |
| | | 224 | | { |
| | 30 | 225 | | await source.PauseAsync(perSourceCts.Token).AsTask().WaitAsync(perSourceCts.Token); |
| | 26 | 226 | | _registry.RecordTransition(source.Name, IngressSourceState.Paused); |
| | 26 | 227 | | } |
| | 2 | 228 | | catch (OperationCanceledException) when (perSourceCts.IsCancellationRequested && !cancellationToken.IsCancellati |
| | | 229 | | { |
| | 2 | 230 | | await _registry.MarkPauseFailedAsync(source.Name, "timeout", new TimeoutException($"Source '{source.Name}' d |
| | 2 | 231 | | await TryForceStopAsync(source, deadlineAt, cancellationToken); |
| | 2 | 232 | | } |
| | 2 | 233 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 234 | | { |
| | 2 | 235 | | await _registry.MarkPauseFailedAsync(source.Name, "exception", ex); |
| | 2 | 236 | | await TryForceStopAsync(source, deadlineAt, cancellationToken); |
| | | 237 | | } |
| | 30 | 238 | | } |
| | | 239 | | |
| | | 240 | | private async Task TryForceStopAsync(IIngressSource source, DateTimeOffset deadlineAt, CancellationToken cancellatio |
| | | 241 | | { |
| | 7 | 242 | | if (source is not IForceStoppable forceStoppable) return; |
| | | 243 | | |
| | | 244 | | // Use the *remaining* budget — not the full overall deadline — so a per-source pause that already |
| | | 245 | | // burned most of the window can't get another full deadline's worth of force-stop runway. |
| | 1 | 246 | | var remaining = deadlineAt - _clock.UtcNow; |
| | 1 | 247 | | if (remaining <= TimeSpan.Zero) |
| | | 248 | | { |
| | 0 | 249 | | _logger.LogWarning("Skipping force-stop of ingress source '{Name}': overall drain deadline already exceeded. |
| | 0 | 250 | | return; |
| | | 251 | | } |
| | | 252 | | |
| | | 253 | | try |
| | | 254 | | { |
| | 1 | 255 | | using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); |
| | 1 | 256 | | cts.CancelAfter(remaining); |
| | 1 | 257 | | await forceStoppable.ForceStopAsync(cts.Token); |
| | 1 | 258 | | _logger.LogInformation("Force-stopped ingress source '{Name}'.", source.Name); |
| | 1 | 259 | | } |
| | 0 | 260 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 261 | | { |
| | 0 | 262 | | _logger.LogWarning(ex, "Force-stop of ingress source '{Name}' failed; drain continues.", source.Name); |
| | 0 | 263 | | } |
| | 4 | 264 | | } |
| | | 265 | | |
| | | 266 | | private async Task<bool> WaitForCyclesAsync(DateTimeOffset deadlineAt, CancellationToken cancellationToken) |
| | | 267 | | { |
| | 2945 | 268 | | while (_cycles.ActiveCount > 0) |
| | | 269 | | { |
| | 2947 | 270 | | if (_clock.UtcNow >= deadlineAt) return false; |
| | 5818 | 271 | | try { await Task.Delay(PollInterval, cancellationToken); } |
| | 0 | 272 | | catch (OperationCanceledException) { return _cycles.ActiveCount == 0; } |
| | | 273 | | } |
| | 17 | 274 | | return true; |
| | 36 | 275 | | } |
| | | 276 | | |
| | | 277 | | /// <summary>Bound on how long the orchestrator waits for a runner to settle after invoking execution cycle cancella |
| | 3 | 278 | | private static readonly TimeSpan ForceCancelSettleTimeout = TimeSpan.FromSeconds(2); |
| | | 279 | | |
| | | 280 | | /// <summary> |
| | | 281 | | /// Per-handle bound on the forensic Interrupted persist write in Phase C. Decoupled from the drain |
| | | 282 | | /// cancellation token: when the host's drain CT has already fired, the persist still has up to this long |
| | | 283 | | /// to land the row + log entry. Keeps "proceed to persist anyway" honest while preventing a stuck DB |
| | | 284 | | /// from hanging shutdown indefinitely. |
| | | 285 | | /// </summary> |
| | 3 | 286 | | private static readonly TimeSpan PersistInterruptedTimeout = TimeSpan.FromSeconds(5); |
| | | 287 | | |
| | | 288 | | /// <summary> |
| | | 289 | | /// Short shutdown budget for the pre-cancel instance snapshot. Force-cancel runs after the drain |
| | | 290 | | /// deadline, so remaining deadline is already zero. WaitAsync unblocks even if the store ignores |
| | | 291 | | /// its cancellation token. The budget starts when a snapshot slot is acquired — not while |
| | | 292 | | /// waiting for <see cref="MaxConcurrentPreCancelSnapshotFinds"/>. |
| | | 293 | | /// </summary> |
| | 3 | 294 | | private static readonly TimeSpan PreCancelSnapshotTimeout = TimeSpan.FromMilliseconds(250); |
| | | 295 | | |
| | | 296 | | /// <summary> |
| | | 297 | | /// Cap on concurrent pre-cancel snapshot Finds. Unbounded <c>Task.WhenAll</c> self-contends |
| | | 298 | | /// under large live-cycle N: more 250ms timeouts, more <c>drainInduced</c> excludes, more |
| | | 299 | | /// Interrupted misses. Queue wait uses the drain token only; each Find still gets its own |
| | | 300 | | /// 250ms after it acquires a slot. |
| | | 301 | | /// </summary> |
| | | 302 | | internal const int MaxConcurrentPreCancelSnapshotFinds = 16; |
| | | 303 | | |
| | | 304 | | private async Task<(int Count, IReadOnlyList<string> Ids)> ForceCancelActiveCyclesAsync(DrainTrigger trigger, string |
| | | 305 | | { |
| | 19 | 306 | | var cap = _options.Value.MaxForceCancelledInstanceIdsReported; |
| | 19 | 307 | | var reportedIds = new List<string>(capacity: Math.Min(cap, 16)); |
| | 19 | 308 | | var totalCancelled = 0; |
| | | 309 | | |
| | 19 | 310 | | var live = _cycles.ListActiveCycles(); |
| | 19 | 311 | | if (live.Count == 0) return (0, Array.Empty<string>()); |
| | | 312 | | |
| | 19 | 313 | | using var scope = _scopeFactory.CreateScope(); |
| | 18 | 314 | | var instanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>(); |
| | 18 | 315 | | var logStore = scope.ServiceProvider.GetRequiredService<IWorkflowExecutionLogStore>(); |
| | | 316 | | |
| | 18 | 317 | | var reason = trigger == DrainTrigger.OperatorForce |
| | 18 | 318 | | ? WorkflowInterruptedPayload.ReasonOperatorForce |
| | 18 | 319 | | : WorkflowInterruptedPayload.ReasonDeadlineBreach; |
| | | 320 | | |
| | | 321 | | // Snapshot before Phase A. Each Find has its own 250ms budget and its own |
| | | 322 | | // DI scope so concurrent reads do not share an EF DbContext (Phase C stays |
| | | 323 | | // sequential on the outer scope for the same reason). Finds run concurrently |
| | | 324 | | // under a semaphore so a large live-cycle set cannot self-contend into |
| | | 325 | | // 250ms timeouts. Timeout/error: exclude that id (prefer preserving |
| | | 326 | | // user-cancel / #8052). A successful null Find is not drain-induced yet — |
| | | 327 | | // no persisted user-cancel exists, but we only write the forensic log if |
| | | 328 | | // Phase A actually cancels that live handle. Drain never promotes a |
| | | 329 | | // Finished/Cancelled row to Running/Interrupted (#8419). |
| | 18 | 330 | | var drainInducedInstanceIds = new HashSet<string>(StringComparer.Ordinal); |
| | 18 | 331 | | var drainInducedCandidateIds = new ConcurrentDictionary<Guid, string>(); |
| | 18 | 332 | | var activeSnapshotHandleIds = new ConcurrentDictionary<Guid, byte>(); |
| | 18 | 333 | | var missingPersistedRowIds = new ConcurrentBag<string>(); |
| | 18 | 334 | | using var snapshotGate = new SemaphoreSlim(MaxConcurrentPreCancelSnapshotFinds); |
| | 18 | 335 | | var snapshotTasks = live.Select(async handle => |
| | 18 | 336 | | { |
| | 18 | 337 | | try |
| | 18 | 338 | | { |
| | 64 | 339 | | await snapshotGate.WaitAsync(cancellationToken).ConfigureAwait(false); |
| | 18 | 340 | | try |
| | 18 | 341 | | { |
| | 64 | 342 | | using var findScope = _scopeFactory.CreateScope(); |
| | 64 | 343 | | var findStore = findScope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>(); |
| | 64 | 344 | | using var perFindCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); |
| | 64 | 345 | | perFindCts.CancelAfter(PreCancelSnapshotTimeout); |
| | 64 | 346 | | var snapshot = await findStore |
| | 64 | 347 | | .FindAsync(new WorkflowInstanceFilter { Id = handle.WorkflowInstanceId }, perFindCts.Token) |
| | 64 | 348 | | .AsTask() |
| | 64 | 349 | | .WaitAsync(perFindCts.Token) |
| | 64 | 350 | | .ConfigureAwait(false); |
| | 48 | 351 | | if (snapshot is null) |
| | 18 | 352 | | { |
| | 6 | 353 | | missingPersistedRowIds.Add(handle.WorkflowInstanceId); |
| | 6 | 354 | | return null; |
| | 18 | 355 | | } |
| | 18 | 356 | | |
| | 42 | 357 | | if (snapshot.SubStatus != WorkflowSubStatus.Cancelled) |
| | 18 | 358 | | { |
| | 39 | 359 | | drainInducedCandidateIds.TryAdd(handle.Id, handle.WorkflowInstanceId); |
| | 39 | 360 | | if (snapshot.IsExecuting) |
| | 38 | 361 | | activeSnapshotHandleIds.TryAdd(handle.Id, 0); |
| | 39 | 362 | | return handle.WorkflowInstanceId; |
| | 18 | 363 | | } |
| | 3 | 364 | | } |
| | 18 | 365 | | finally |
| | 18 | 366 | | { |
| | 18 | 367 | | // Release after the await budget, not after the store call returns. |
| | 18 | 368 | | // A cancel-ignoring Find can outlive this slot (documented secondary; |
| | 18 | 369 | | // holding the slot until it completes would stall WhenAll / Phase A). |
| | 64 | 370 | | snapshotGate.Release(); |
| | 18 | 371 | | } |
| | 3 | 372 | | } |
| | 16 | 373 | | catch (Exception ex) when (!ex.IsFatal()) |
| | 18 | 374 | | { |
| | 16 | 375 | | _logger.LogWarning(ex, "Pre-cancel snapshot for instance {InstanceId} timed out or failed; excluding fro |
| | 16 | 376 | | } |
| | 18 | 377 | | |
| | 19 | 378 | | return null; |
| | 82 | 379 | | }); |
| | | 380 | | |
| | 18 | 381 | | await Task.WhenAll(snapshotTasks).ConfigureAwait(false); |
| | | 382 | | |
| | | 383 | | // Force-cancel proceeds in three phases. The split exists because cancelling and |
| | | 384 | | // awaiting in the same loop made every execution cycle after the first run at full speed |
| | | 385 | | // through the prior execution cycle's settle window — total wall time was O(N × settle |
| | | 386 | | // timeout), defeating the intent of force-cancel under concurrency. |
| | | 387 | | // |
| | | 388 | | // Phase A — cancel every handle synchronously. CancellationTokenSource.Cancel is |
| | | 389 | | // cheap and we want every runner to observe cancellation simultaneously rather |
| | | 390 | | // than serialized behind preceding settle waits. |
| | 18 | 391 | | var cancelledInstanceIds = new HashSet<string>(StringComparer.Ordinal); |
| | 18 | 392 | | var cancelledHandles = new List<ExecutionCycleHandle>(live.Count); |
| | 18 | 393 | | var handlesToPersist = new List<ExecutionCycleHandle>(live.Count); |
| | 18 | 394 | | var activeSnapshotHandlesToRecover = new List<ExecutionCycleHandle>(live.Count); |
| | 18 | 395 | | var disposedActiveSnapshotHandleIds = new HashSet<Guid>(); |
| | 163 | 396 | | foreach (var handle in live) |
| | | 397 | | { |
| | | 398 | | try |
| | | 399 | | { |
| | | 400 | | // Only a real transition counts. Cancel() is a no-op on an already-disposed handle |
| | | 401 | | // (cycle finished during snapshot); treating that as drain-induced would rewrite a |
| | | 402 | | // later Finished/Cancelled row the runner already committed. |
| | 64 | 403 | | if (!handle.TryCancel()) |
| | | 404 | | { |
| | 8 | 405 | | if (activeSnapshotHandleIds.ContainsKey(handle.Id)) |
| | 4 | 406 | | activeSnapshotHandlesToRecover.Add(handle); |
| | 8 | 407 | | continue; |
| | | 408 | | } |
| | | 409 | | |
| | 55 | 410 | | totalCancelled++; |
| | 55 | 411 | | cancelledHandles.Add(handle); |
| | 55 | 412 | | handlesToPersist.Add(handle); |
| | 55 | 413 | | cancelledInstanceIds.Add(handle.WorkflowInstanceId); |
| | 110 | 414 | | if (reportedIds.Count < cap) reportedIds.Add(handle.WorkflowInstanceId); |
| | 55 | 415 | | } |
| | 1 | 416 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 417 | | { |
| | 0 | 418 | | _logger.LogError(ex, "Failed to cancel execution cycle {ExecutionCycleId} (instance={InstanceId}).", han |
| | 0 | 419 | | } |
| | | 420 | | } |
| | | 421 | | |
| | | 422 | | var recoveryHandleIds = activeSnapshotHandlesToRecover.Select(handle => handle.Id).ToHashSet(); |
| | | 423 | | |
| | | 424 | | // A live handle we ourselves cancelled whose snapshot found no row is drain-induced: |
| | | 425 | | // there was no persisted user-cancel to preserve. Timeout/error and disposed no-ops stay excluded. |
| | | 426 | | // Do not use reportedIds here — that list is capped by MaxForceCancelledInstanceIdsReported. |
| | 144 | 427 | | foreach (var handle in cancelledHandles) |
| | | 428 | | { |
| | 55 | 429 | | if (drainInducedCandidateIds.TryGetValue(handle.Id, out var instanceId)) |
| | 33 | 430 | | drainInducedInstanceIds.Add(instanceId); |
| | | 431 | | } |
| | | 432 | | |
| | 46 | 433 | | foreach (var instanceId in missingPersistedRowIds) |
| | | 434 | | { |
| | 6 | 435 | | if (cancelledInstanceIds.Contains(instanceId)) |
| | 3 | 436 | | drainInducedInstanceIds.Add(instanceId); |
| | | 437 | | } |
| | | 438 | | |
| | | 439 | | // Phase B — wait for every runner to settle in parallel under a shared deadline. |
| | | 440 | | // The handle disposes when ExecutionCycleTrackingMiddleware exits its `using` block, which |
| | | 441 | | // is after the workflow runner has finished its commit. We want runners' terminal |
| | | 442 | | // commits to land before we overwrite the sub-status with Interrupted — but we |
| | | 443 | | // bound the wait so a non-cancellable activity cannot block drain, accepting the |
| | | 444 | | // runner-clobber race for that one instance (the recovery scan picks it up). |
| | | 445 | | // Total wall time for this phase is at most ForceCancelSettleTimeout regardless |
| | | 446 | | // of N. Retained active-snapshot recovery candidates use an independent token so |
| | | 447 | | // a canceled drain still observes their deferred disposal within that same bound. |
| | 17 | 448 | | var settleTasks = live.Select(async handle => |
| | 17 | 449 | | { |
| | 17 | 450 | | try |
| | 17 | 451 | | { |
| | 63 | 452 | | var settleCancellationToken = recoveryHandleIds.Contains(handle.Id) ? CancellationToken.None : cancellat |
| | 63 | 453 | | await handle.Disposed.WaitAsync(ForceCancelSettleTimeout, settleCancellationToken).ConfigureAwait(false) |
| | 10 | 454 | | } |
| | 52 | 455 | | catch (TimeoutException) |
| | 17 | 456 | | { |
| | 52 | 457 | | if (recoveryHandleIds.Contains(handle.Id)) |
| | 17 | 458 | | { |
| | 0 | 459 | | _logger.LogWarning("Retained active-snapshot recovery candidate {ExecutionCycleId} (instance={Instan |
| | 17 | 460 | | } |
| | 17 | 461 | | else |
| | 17 | 462 | | { |
| | 52 | 463 | | _logger.LogWarning("Execution cycle {ExecutionCycleId} (instance={InstanceId}) did not settle within |
| | 17 | 464 | | } |
| | 52 | 465 | | } |
| | 2 | 466 | | catch (OperationCanceledException) { /* drain CT fired — proceed to persist anyway */ } |
| | 80 | 467 | | }); |
| | 17 | 468 | | await Task.WhenAll(settleTasks).ConfigureAwait(false); |
| | | 469 | | |
| | 42 | 470 | | foreach (var handle in activeSnapshotHandlesToRecover) |
| | | 471 | | { |
| | 4 | 472 | | if (!handle.Disposed.IsCompleted) |
| | | 473 | | continue; |
| | | 474 | | |
| | 4 | 475 | | disposedActiveSnapshotHandleIds.Add(handle.Id); |
| | 4 | 476 | | handlesToPersist.Add(handle); |
| | | 477 | | } |
| | | 478 | | |
| | | 479 | | // Phase C — persist Interrupted for every cancelled handle and disposed active checkpoint. |
| | | 480 | | // Sequential to keep DbContext usage single-threaded; per-handle persistence is small. |
| | | 481 | | // |
| | | 482 | | // Each persist runs under its own bounded token that is NOT linked to the drain CT. |
| | | 483 | | // Phase B's catch on OperationCanceledException explicitly comments "drain CT fired — |
| | | 484 | | // proceed to persist anyway"; reusing the cancelled token here would have made that |
| | | 485 | | // a lie — the very first await inside PersistInterruptedAsync would observe the |
| | | 486 | | // cancelled token and throw. Result: every execution cycle would be left in an unrecovered |
| | | 487 | | // executing state on host shutdown. The bounded non-drain token preserves the |
| | | 488 | | // forensic write while preventing a stuck DB from hanging shutdown indefinitely. |
| | 152 | 489 | | foreach (var handle in handlesToPersist) |
| | | 490 | | { |
| | | 491 | | try |
| | | 492 | | { |
| | 59 | 493 | | using var persistCts = new CancellationTokenSource(PersistInterruptedTimeout); |
| | 59 | 494 | | await PersistInterruptedAsync(instanceStore, logStore, handle, generationId, reason, drainInducedInstanc |
| | 59 | 495 | | } |
| | 0 | 496 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 497 | | { |
| | 0 | 498 | | _logger.LogError(ex, "Failed to persist Interrupted for execution cycle {ExecutionCycleId} (instance={In |
| | 0 | 499 | | } |
| | 59 | 500 | | } |
| | | 501 | | |
| | 17 | 502 | | return (totalCancelled, reportedIds); |
| | 17 | 503 | | } |
| | | 504 | | |
| | | 505 | | private async Task PersistInterruptedAsync( |
| | | 506 | | IWorkflowInstanceStore instanceStore, |
| | | 507 | | IWorkflowExecutionLogStore logStore, |
| | | 508 | | ExecutionCycleHandle handle, |
| | | 509 | | string generationId, |
| | | 510 | | string reason, |
| | | 511 | | HashSet<string> drainInducedInstanceIds, |
| | | 512 | | bool requireExecuting, |
| | | 513 | | CancellationToken cancellationToken) |
| | | 514 | | { |
| | 59 | 515 | | var instance = await instanceStore.FindAsync(new WorkflowInstanceFilter { Id = handle.WorkflowInstanceId }, canc |
| | | 516 | | |
| | | 517 | | // A disposed checkpoint is only recoverable while its row still shows execution; a later |
| | | 518 | | // natural suspension or completion must remain untouched. |
| | 59 | 519 | | if (requireExecuting && (instance is null || instance.Status == WorkflowStatus.Finished || !instance.IsExecuting |
| | 1 | 520 | | return; |
| | | 521 | | |
| | 58 | 522 | | if (instance is null) |
| | | 523 | | { |
| | | 524 | | // The instance row was never persisted (e.g., a execution cycle whose runner never reached commitStateHandl |
| | | 525 | | // We cannot update a row that doesn't exist, but we MUST still write the forensic trail so postmortem |
| | | 526 | | // recovery can audit what happened. Write a synthetic WorkflowInterrupted log entry with the |
| | | 527 | | // PersistenceFailure discriminator — the workflow-definition fields are unknown so they remain empty. |
| | 1 | 528 | | var orphanPayload = new WorkflowInterruptedPayload( |
| | 1 | 529 | | InterruptedAt: _clock.UtcNow, |
| | 1 | 530 | | Reason: WorkflowInterruptedPayload.ReasonPersistenceFailure, |
| | 1 | 531 | | GenerationId: generationId, |
| | 1 | 532 | | LastActivityId: null, |
| | 1 | 533 | | LastActivityNodeId: null, |
| | 1 | 534 | | IngressSourceName: handle.IngressSourceName, |
| | 1 | 535 | | ExecutionCycleDuration: _clock.UtcNow - handle.StartedAt); |
| | | 536 | | try |
| | | 537 | | { |
| | 1 | 538 | | await logStore.AddAsync(new Entities.WorkflowExecutionLogRecord |
| | 1 | 539 | | { |
| | 1 | 540 | | Id = _identityGenerator.GenerateId(), |
| | 1 | 541 | | WorkflowInstanceId = handle.WorkflowInstanceId, |
| | 1 | 542 | | WorkflowDefinitionId = string.Empty, |
| | 1 | 543 | | WorkflowDefinitionVersionId = string.Empty, |
| | 1 | 544 | | WorkflowVersion = 0, |
| | 1 | 545 | | ActivityInstanceId = string.Empty, |
| | 1 | 546 | | ActivityId = string.Empty, |
| | 1 | 547 | | ActivityType = string.Empty, |
| | 1 | 548 | | ActivityNodeId = string.Empty, |
| | 1 | 549 | | Timestamp = orphanPayload.InterruptedAt, |
| | 1 | 550 | | EventName = WorkflowInterruptedPayload.WorkflowInterruptedEventName, |
| | 1 | 551 | | Source = "Elsa.Workflows.Runtime.GracefulShutdown", |
| | 1 | 552 | | Message = $"Workflow execution cycle force-cancelled by the runtime ({orphanPayload.Reason}); no ins |
| | 1 | 553 | | Payload = orphanPayload, |
| | 1 | 554 | | }, cancellationToken); |
| | 1 | 555 | | } |
| | 0 | 556 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 557 | | { |
| | 0 | 558 | | _logger.LogWarning(ex, "Failed to write WorkflowInterrupted log entry for orphan execution cycle {Execut |
| | 0 | 559 | | } |
| | 1 | 560 | | return; |
| | | 561 | | } |
| | | 562 | | |
| | 57 | 563 | | if (ShouldSkipInterruptedPersist(instance, drainInducedInstanceIds)) |
| | | 564 | | { |
| | 19 | 565 | | _logger.LogInformation( |
| | 19 | 566 | | "Skipping Interrupted persist for instance {InstanceId}: already in terminal status {Status}/{SubStatus} |
| | 19 | 567 | | instance.Id, instance.Status, instance.SubStatus); |
| | 19 | 568 | | return; |
| | | 569 | | } |
| | | 570 | | |
| | 38 | 571 | | var actualReason = reason; |
| | | 572 | | |
| | | 573 | | try |
| | | 574 | | { |
| | | 575 | | // Conditional write: do not SaveAsync the Find snapshot. Never promote a |
| | | 576 | | // Finished/Cancelled row (#8419) — that left a Running/Interrupted row whose |
| | | 577 | | // serialized WorkflowState still said Finished with no scheduled work, and the |
| | | 578 | | // recovery scanner then requeued a workflow that can never resume. |
| | | 579 | | // allowFinishedCancelled remains on the store signature for 3.8.4 compatibility; |
| | | 580 | | // drain no longer passes true. Drain-induced Finished/Cancelled rows skip the |
| | | 581 | | // mark and still receive the WorkflowInterrupted forensic log below. |
| | 38 | 582 | | if (instance.Status != WorkflowStatus.Finished) |
| | | 583 | | { |
| | 32 | 584 | | var marked = await instanceStore.TryMarkInterruptedAsync(instance.Id, cancellationToken); |
| | 31 | 585 | | if (!marked) |
| | | 586 | | { |
| | 0 | 587 | | _logger.LogInformation( |
| | 0 | 588 | | "Skipping Interrupted persist for instance {InstanceId}: a concurrent persist already left it in |
| | 0 | 589 | | instance.Id); |
| | 0 | 590 | | return; |
| | | 591 | | } |
| | | 592 | | } |
| | 37 | 593 | | } |
| | 1 | 594 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 595 | | { |
| | 1 | 596 | | _logger.LogWarning(ex, "Failed to persist Interrupted sub-status for instance {InstanceId}; falling back to |
| | 1 | 597 | | actualReason = WorkflowInterruptedPayload.ReasonPersistenceFailure; |
| | 1 | 598 | | } |
| | | 599 | | |
| | 38 | 600 | | var payload = new WorkflowInterruptedPayload( |
| | 38 | 601 | | InterruptedAt: _clock.UtcNow, |
| | 38 | 602 | | Reason: actualReason, |
| | 38 | 603 | | GenerationId: generationId, |
| | 38 | 604 | | LastActivityId: null, |
| | 38 | 605 | | LastActivityNodeId: null, |
| | 38 | 606 | | IngressSourceName: handle.IngressSourceName, |
| | 38 | 607 | | ExecutionCycleDuration: _clock.UtcNow - handle.StartedAt); |
| | | 608 | | |
| | | 609 | | try |
| | | 610 | | { |
| | 38 | 611 | | await logStore.LogWorkflowInterruptedAsync(_identityGenerator, instance, payload, cancellationToken); |
| | 38 | 612 | | } |
| | 0 | 613 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 614 | | { |
| | 0 | 615 | | _logger.LogWarning(ex, "Failed to write WorkflowInterrupted log entry for instance {InstanceId}.", instance. |
| | 0 | 616 | | } |
| | 59 | 617 | | } |
| | | 618 | | |
| | | 619 | | private IReadOnlyList<IngressSourceFinalState> BuildSourceFinalStates() |
| | | 620 | | { |
| | 35 | 621 | | return _registry.Snapshot() |
| | 23 | 622 | | .Select(s => new IngressSourceFinalState(s.Name, s.State, s.LastError, WasForceStopped: s.State == IngressSo |
| | 35 | 623 | | .ToArray(); |
| | | 624 | | } |
| | | 625 | | |
| | | 626 | | /// <summary> |
| | | 627 | | /// Skip persist and forensic log when the row is already a real terminal outcome: |
| | | 628 | | /// natural completion, fault, user-Cancelled, or unknown pre-state (snapshot |
| | | 629 | | /// timeout/error). Drain-induced Finished/Cancelled rows still get the |
| | | 630 | | /// <c>WorkflowInterrupted</c> log but are not promoted to Running/Interrupted (#8419). |
| | | 631 | | /// </summary> |
| | | 632 | | private static bool ShouldSkipInterruptedPersist(WorkflowInstance instance, HashSet<string> drainInducedInstanceIds) |
| | | 633 | | { |
| | 57 | 634 | | if (instance.Status != WorkflowStatus.Finished) |
| | 32 | 635 | | return false; |
| | | 636 | | |
| | 25 | 637 | | if (instance.SubStatus == WorkflowSubStatus.Cancelled) |
| | 25 | 638 | | return !drainInducedInstanceIds.Contains(instance.Id); |
| | | 639 | | |
| | 0 | 640 | | return true; |
| | | 641 | | } |
| | | 642 | | } |