< Summary

Information
Class: Elsa.Workflows.Runtime.Services.DrainOrchestrator
Assembly: Elsa.Workflows.Runtime
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/DrainOrchestrator.cs
Line coverage
92%
Covered lines: 313
Uncovered lines: 24
Coverable lines: 337
Total lines: 642
Line coverage: 92.8%
Branch coverage
88%
Covered branches: 90
Total branches: 102
Branch coverage: 88.2%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.cctor()100%11100%
.ctor(...)100%11100%
DrainAsync()65%202093.75%
ComputeEffectiveDeadline(...)83.33%66100%
PauseAllSourcesAsync()100%22100%
PauseOneSourceAsync()62.5%88100%
TryForceStopAsync()75%5464.28%
WaitForCyclesAsync()100%4483.33%
ForceCancelActiveCyclesAsync()92.3%262694.18%
<ForceCancelActiveCyclesAsync()100%66100%
<ForceCancelActiveCyclesAsync()75%4490%
PersistInterruptedAsync()87.5%191677.27%
BuildSourceFinalStates()50%22100%
ShouldSkipInterruptedPersist(...)75%4480%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/DrainOrchestrator.cs

#LineLine coverage
 1using System.Collections.Concurrent;
 2using System.Diagnostics;
 3using Elsa.Common;
 4using Elsa.Workflows.Management;
 5using Elsa.Workflows.Management.Entities;
 6using Elsa.Workflows.Management.Filters;
 7using Elsa.Workflows.Runtime.Options;
 8using Microsoft.Extensions.DependencyInjection;
 9using Microsoft.Extensions.Hosting;
 10using Microsoft.Extensions.Logging;
 11using Microsoft.Extensions.Options;
 12
 13namespace 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>
 37public sealed class DrainOrchestrator : IDrainOrchestrator
 38{
 339    private static readonly TimeSpan SafetyEpsilon = TimeSpan.FromMilliseconds(500);
 340    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
 4552    private readonly object _sync = new();
 53    private DrainOutcome? _previousOutcome;
 54    private bool _drainInProgress;
 55
 4556    public DrainOrchestrator(
 4557        IQuiescenceSignal signal,
 4558        IIngressSourceRegistry registry,
 4559        IExecutionCycleRegistry cycles,
 4560        IServiceScopeFactory scopeFactory,
 4561        IOptions<GracefulShutdownOptions> options,
 4562        IOptions<HostOptions> hostOptions,
 4563        ISystemClock clock,
 4564        IIdentityGenerator identityGenerator,
 4565        ILogger<DrainOrchestrator> logger)
 66    {
 4567        _signal = signal;
 4568        _registry = registry;
 4569        _cycles = cycles;
 4570        _scopeFactory = scopeFactory;
 4571        _options = options;
 4572        _hostOptions = hostOptions;
 4573        _clock = clock;
 4574        _identityGenerator = identityGenerator;
 4575        _logger = logger;
 4576    }
 77
 78    /// <inheritdoc />
 79    public async ValueTask<DrainOutcome> DrainAsync(DrainTrigger trigger, CancellationToken cancellationToken = default)
 80    {
 4181        lock (_sync)
 82        {
 4183            if (_drainInProgress)
 84            {
 385                if (trigger == DrainTrigger.OperatorForce && _previousOutcome is not null)
 086                    return _previousOutcome with { WasCached = true };
 387                throw new InvalidOperationException($"Drain already in progress; second invocation rejected (trigger={tr
 88            }
 3889            if (_previousOutcome is not null)
 90            {
 391                if (trigger == DrainTrigger.OperatorForce) return _previousOutcome with { WasCached = true };
 192                throw new InvalidOperationException("Drain already completed in this generation; subsequent non-force in
 93            }
 3694            _drainInProgress = true;
 3695        }
 96
 3697        var startedAt = _clock.UtcNow;
 3698        var deadline = ComputeEffectiveDeadline(trigger);
 3699        var sw = Stopwatch.StartNew();
 36100        TimeSpan pausePhase = TimeSpan.Zero;
 36101        TimeSpan waitPhase = TimeSpan.Zero;
 102
 103        try
 104        {
 36105            await _signal.BeginDrainAsync(cancellationToken);
 36106            _logger.LogInformation("Drain initiated (trigger={Trigger}, deadline={Deadline}).", trigger, deadline);
 107
 36108            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.
 36113            await PauseAllSourcesAsync(deadlineAt, cancellationToken);
 36114            pausePhase = sw.Elapsed;
 36115            _logger.LogInformation("Ingress pause phase complete in {Elapsed}.", pausePhase);
 116
 117            // Phase 2: wait for active execution cycles to drain.
 36118            var waitStart = sw.Elapsed;
 36119            var deadlineBreach = !await WaitForCyclesAsync(deadlineAt, cancellationToken);
 36120            waitPhase = sw.Elapsed - waitStart;
 121
 122            // Phase 3: outcome assembly.
 36123            int forceCancelled = 0;
 36124            IReadOnlyList<string> forceCancelledIds = Array.Empty<string>();
 125
 36126            if (deadlineBreach || trigger == DrainTrigger.OperatorForce)
 127            {
 19128                _logger.LogWarning("Drain {What}; force-cancelling {Count} active execution cycle(s).",
 19129                    deadlineBreach ? "deadline exceeded" : "operator-forced", _cycles.ActiveCount);
 19130                (forceCancelled, forceCancelledIds) = await ForceCancelActiveCyclesAsync(trigger, _signal.CurrentState.G
 131            }
 132
 133            // OperatorForce always reports Forced regardless of zero-deadline breach mechanics.
 34134            var result = trigger switch
 34135            {
 14136                DrainTrigger.OperatorForce => DrainResult.Forced,
 23137                _ when deadlineBreach => DrainResult.DeadlineExceeded,
 17138                _ => DrainResult.CompletedWithinDeadline,
 34139            };
 140
 34141            var outcome = new DrainOutcome(
 34142                OverallResult: result,
 34143                StartedAt: startedAt,
 34144                CompletedAt: _clock.UtcNow,
 34145                PausePhaseDuration: pausePhase,
 34146                WaitPhaseDuration: waitPhase,
 34147                Sources: BuildSourceFinalStates(),
 34148                ExecutionCyclesForceCancelledCount: forceCancelled,
 34149                ForceCancelledInstanceIds: forceCancelledIds);
 150
 68151            lock (_sync) _previousOutcome = outcome;
 34152            _logger.LogInformation("Drain completed: {Result} (paused={Paused}, waited={Waited}, forceCancelled={ForceCa
 34153                outcome.OverallResult, outcome.PausePhaseDuration, outcome.WaitPhaseDuration, outcome.ExecutionCyclesFor
 34154            return outcome;
 155        }
 2156        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.
 1162            _logger.LogError(ex, "Drain aborted by unhandled exception.");
 1163            var outcome = new DrainOutcome(
 1164                OverallResult: DrainResult.AbortedByUnhandledException,
 1165                StartedAt: startedAt,
 1166                CompletedAt: _clock.UtcNow,
 1167                PausePhaseDuration: pausePhase,
 1168                WaitPhaseDuration: waitPhase,
 1169                Sources: BuildSourceFinalStates(),
 1170                ExecutionCyclesForceCancelledCount: 0,
 1171                ForceCancelledInstanceIds: Array.Empty<string>());
 2172            lock (_sync) _previousOutcome = outcome;
 1173            return outcome;
 174        }
 175        finally
 176        {
 72177            lock (_sync) _drainInProgress = false;
 178        }
 36179    }
 180
 181    private TimeSpan ComputeEffectiveDeadline(DrainTrigger trigger)
 182    {
 51183        if (trigger == DrainTrigger.OperatorForce) return TimeSpan.Zero;
 184
 21185        var configured = _options.Value.DrainDeadline;
 21186        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
 21189        if (hostShutdown > SafetyEpsilon)
 190        {
 21191            var hostBudget = hostShutdown - SafetyEpsilon;
 28192            if (hostBudget < configured) return hostBudget;
 193        }
 194
 14195        return configured;
 196    }
 197
 198    private async Task PauseAllSourcesAsync(DateTimeOffset deadlineAt, CancellationToken cancellationToken)
 199    {
 36200        var sources = _registry.Sources;
 55201        if (sources.Count == 0) return;
 202
 47203        var pauseTasks = sources.Select(source => PauseOneSourceAsync(source, deadlineAt, cancellationToken)).ToArray();
 17204        await Task.WhenAll(pauseTasks);
 36205    }
 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.
 30213        var configured = _options.Value.IngressPauseTimeout;
 30214        var sourceDeadline = source.PauseTimeout > TimeSpan.Zero ? source.PauseTimeout : configured;
 30215        var overallRemaining = deadlineAt - _clock.UtcNow;
 38216        if (sourceDeadline > overallRemaining) sourceDeadline = overallRemaining;
 30217        if (sourceDeadline <= TimeSpan.Zero) sourceDeadline = TimeSpan.FromMilliseconds(1);
 218
 30219        _registry.RecordTransition(source.Name, IngressSourceState.Pausing);
 30220        using var perSourceCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 30221        perSourceCts.CancelAfter(sourceDeadline);
 222
 223        try
 224        {
 30225            await source.PauseAsync(perSourceCts.Token).AsTask().WaitAsync(perSourceCts.Token);
 26226            _registry.RecordTransition(source.Name, IngressSourceState.Paused);
 26227        }
 2228        catch (OperationCanceledException) when (perSourceCts.IsCancellationRequested && !cancellationToken.IsCancellati
 229        {
 2230            await _registry.MarkPauseFailedAsync(source.Name, "timeout", new TimeoutException($"Source '{source.Name}' d
 2231            await TryForceStopAsync(source, deadlineAt, cancellationToken);
 2232        }
 2233        catch (Exception ex) when (!ex.IsFatal())
 234        {
 2235            await _registry.MarkPauseFailedAsync(source.Name, "exception", ex);
 2236            await TryForceStopAsync(source, deadlineAt, cancellationToken);
 237        }
 30238    }
 239
 240    private async Task TryForceStopAsync(IIngressSource source, DateTimeOffset deadlineAt, CancellationToken cancellatio
 241    {
 7242        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.
 1246        var remaining = deadlineAt - _clock.UtcNow;
 1247        if (remaining <= TimeSpan.Zero)
 248        {
 0249            _logger.LogWarning("Skipping force-stop of ingress source '{Name}': overall drain deadline already exceeded.
 0250            return;
 251        }
 252
 253        try
 254        {
 1255            using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 1256            cts.CancelAfter(remaining);
 1257            await forceStoppable.ForceStopAsync(cts.Token);
 1258            _logger.LogInformation("Force-stopped ingress source '{Name}'.", source.Name);
 1259        }
 0260        catch (Exception ex) when (!ex.IsFatal())
 261        {
 0262            _logger.LogWarning(ex, "Force-stop of ingress source '{Name}' failed; drain continues.", source.Name);
 0263        }
 4264    }
 265
 266    private async Task<bool> WaitForCyclesAsync(DateTimeOffset deadlineAt, CancellationToken cancellationToken)
 267    {
 2945268        while (_cycles.ActiveCount > 0)
 269        {
 2947270            if (_clock.UtcNow >= deadlineAt) return false;
 5818271            try { await Task.Delay(PollInterval, cancellationToken); }
 0272            catch (OperationCanceledException) { return _cycles.ActiveCount == 0; }
 273        }
 17274        return true;
 36275    }
 276
 277    /// <summary>Bound on how long the orchestrator waits for a runner to settle after invoking execution cycle cancella
 3278    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>
 3286    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>
 3294    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    {
 19306        var cap = _options.Value.MaxForceCancelledInstanceIdsReported;
 19307        var reportedIds = new List<string>(capacity: Math.Min(cap, 16));
 19308        var totalCancelled = 0;
 309
 19310        var live = _cycles.ListActiveCycles();
 19311        if (live.Count == 0) return (0, Array.Empty<string>());
 312
 19313        using var scope = _scopeFactory.CreateScope();
 18314        var instanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
 18315        var logStore = scope.ServiceProvider.GetRequiredService<IWorkflowExecutionLogStore>();
 316
 18317        var reason = trigger == DrainTrigger.OperatorForce
 18318            ? WorkflowInterruptedPayload.ReasonOperatorForce
 18319            : 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).
 18330        var drainInducedInstanceIds = new HashSet<string>(StringComparer.Ordinal);
 18331        var drainInducedCandidateIds = new ConcurrentDictionary<Guid, string>();
 18332        var activeSnapshotHandleIds = new ConcurrentDictionary<Guid, byte>();
 18333        var missingPersistedRowIds = new ConcurrentBag<string>();
 18334        using var snapshotGate = new SemaphoreSlim(MaxConcurrentPreCancelSnapshotFinds);
 18335        var snapshotTasks = live.Select(async handle =>
 18336        {
 18337            try
 18338            {
 64339                await snapshotGate.WaitAsync(cancellationToken).ConfigureAwait(false);
 18340                try
 18341                {
 64342                    using var findScope = _scopeFactory.CreateScope();
 64343                    var findStore = findScope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
 64344                    using var perFindCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 64345                    perFindCts.CancelAfter(PreCancelSnapshotTimeout);
 64346                    var snapshot = await findStore
 64347                        .FindAsync(new WorkflowInstanceFilter { Id = handle.WorkflowInstanceId }, perFindCts.Token)
 64348                        .AsTask()
 64349                        .WaitAsync(perFindCts.Token)
 64350                        .ConfigureAwait(false);
 48351                    if (snapshot is null)
 18352                    {
 6353                        missingPersistedRowIds.Add(handle.WorkflowInstanceId);
 6354                        return null;
 18355                    }
 18356
 42357                    if (snapshot.SubStatus != WorkflowSubStatus.Cancelled)
 18358                    {
 39359                        drainInducedCandidateIds.TryAdd(handle.Id, handle.WorkflowInstanceId);
 39360                        if (snapshot.IsExecuting)
 38361                            activeSnapshotHandleIds.TryAdd(handle.Id, 0);
 39362                        return handle.WorkflowInstanceId;
 18363                    }
 3364                }
 18365                finally
 18366                {
 18367                    // Release after the await budget, not after the store call returns.
 18368                    // A cancel-ignoring Find can outlive this slot (documented secondary;
 18369                    // holding the slot until it completes would stall WhenAll / Phase A).
 64370                    snapshotGate.Release();
 18371                }
 3372            }
 16373            catch (Exception ex) when (!ex.IsFatal())
 18374            {
 16375                _logger.LogWarning(ex, "Pre-cancel snapshot for instance {InstanceId} timed out or failed; excluding fro
 16376            }
 18377
 19378            return null;
 82379        });
 380
 18381        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.
 18391        var cancelledInstanceIds = new HashSet<string>(StringComparer.Ordinal);
 18392        var cancelledHandles = new List<ExecutionCycleHandle>(live.Count);
 18393        var handlesToPersist = new List<ExecutionCycleHandle>(live.Count);
 18394        var activeSnapshotHandlesToRecover = new List<ExecutionCycleHandle>(live.Count);
 18395        var disposedActiveSnapshotHandleIds = new HashSet<Guid>();
 163396        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.
 64403                if (!handle.TryCancel())
 404                {
 8405                    if (activeSnapshotHandleIds.ContainsKey(handle.Id))
 4406                        activeSnapshotHandlesToRecover.Add(handle);
 8407                    continue;
 408                }
 409
 55410                totalCancelled++;
 55411                cancelledHandles.Add(handle);
 55412                handlesToPersist.Add(handle);
 55413                cancelledInstanceIds.Add(handle.WorkflowInstanceId);
 110414                if (reportedIds.Count < cap) reportedIds.Add(handle.WorkflowInstanceId);
 55415            }
 1416            catch (Exception ex) when (!ex.IsFatal())
 417            {
 0418                _logger.LogError(ex, "Failed to cancel execution cycle {ExecutionCycleId} (instance={InstanceId}).", han
 0419            }
 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.
 144427        foreach (var handle in cancelledHandles)
 428        {
 55429            if (drainInducedCandidateIds.TryGetValue(handle.Id, out var instanceId))
 33430                drainInducedInstanceIds.Add(instanceId);
 431        }
 432
 46433        foreach (var instanceId in missingPersistedRowIds)
 434        {
 6435            if (cancelledInstanceIds.Contains(instanceId))
 3436                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.
 17448        var settleTasks = live.Select(async handle =>
 17449        {
 17450            try
 17451            {
 63452                var settleCancellationToken = recoveryHandleIds.Contains(handle.Id) ? CancellationToken.None : cancellat
 63453                await handle.Disposed.WaitAsync(ForceCancelSettleTimeout, settleCancellationToken).ConfigureAwait(false)
 10454            }
 52455            catch (TimeoutException)
 17456            {
 52457                if (recoveryHandleIds.Contains(handle.Id))
 17458                {
 0459                    _logger.LogWarning("Retained active-snapshot recovery candidate {ExecutionCycleId} (instance={Instan
 17460                }
 17461                else
 17462                {
 52463                    _logger.LogWarning("Execution cycle {ExecutionCycleId} (instance={InstanceId}) did not settle within
 17464                }
 52465            }
 2466            catch (OperationCanceledException) { /* drain CT fired — proceed to persist anyway */ }
 80467        });
 17468        await Task.WhenAll(settleTasks).ConfigureAwait(false);
 469
 42470        foreach (var handle in activeSnapshotHandlesToRecover)
 471        {
 4472            if (!handle.Disposed.IsCompleted)
 473                continue;
 474
 4475            disposedActiveSnapshotHandleIds.Add(handle.Id);
 4476            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.
 152489        foreach (var handle in handlesToPersist)
 490        {
 491            try
 492            {
 59493                using var persistCts = new CancellationTokenSource(PersistInterruptedTimeout);
 59494                await PersistInterruptedAsync(instanceStore, logStore, handle, generationId, reason, drainInducedInstanc
 59495            }
 0496            catch (Exception ex) when (!ex.IsFatal())
 497            {
 0498                _logger.LogError(ex, "Failed to persist Interrupted for execution cycle {ExecutionCycleId} (instance={In
 0499            }
 59500        }
 501
 17502        return (totalCancelled, reportedIds);
 17503    }
 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    {
 59515        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.
 59519        if (requireExecuting && (instance is null || instance.Status == WorkflowStatus.Finished || !instance.IsExecuting
 1520            return;
 521
 58522        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.
 1528            var orphanPayload = new WorkflowInterruptedPayload(
 1529                InterruptedAt: _clock.UtcNow,
 1530                Reason: WorkflowInterruptedPayload.ReasonPersistenceFailure,
 1531                GenerationId: generationId,
 1532                LastActivityId: null,
 1533                LastActivityNodeId: null,
 1534                IngressSourceName: handle.IngressSourceName,
 1535                ExecutionCycleDuration: _clock.UtcNow - handle.StartedAt);
 536            try
 537            {
 1538                await logStore.AddAsync(new Entities.WorkflowExecutionLogRecord
 1539                {
 1540                    Id = _identityGenerator.GenerateId(),
 1541                    WorkflowInstanceId = handle.WorkflowInstanceId,
 1542                    WorkflowDefinitionId = string.Empty,
 1543                    WorkflowDefinitionVersionId = string.Empty,
 1544                    WorkflowVersion = 0,
 1545                    ActivityInstanceId = string.Empty,
 1546                    ActivityId = string.Empty,
 1547                    ActivityType = string.Empty,
 1548                    ActivityNodeId = string.Empty,
 1549                    Timestamp = orphanPayload.InterruptedAt,
 1550                    EventName = WorkflowInterruptedPayload.WorkflowInterruptedEventName,
 1551                    Source = "Elsa.Workflows.Runtime.GracefulShutdown",
 1552                    Message = $"Workflow execution cycle force-cancelled by the runtime ({orphanPayload.Reason}); no ins
 1553                    Payload = orphanPayload,
 1554                }, cancellationToken);
 1555            }
 0556            catch (Exception ex) when (!ex.IsFatal())
 557            {
 0558                _logger.LogWarning(ex, "Failed to write WorkflowInterrupted log entry for orphan execution cycle {Execut
 0559            }
 1560            return;
 561        }
 562
 57563        if (ShouldSkipInterruptedPersist(instance, drainInducedInstanceIds))
 564        {
 19565            _logger.LogInformation(
 19566                "Skipping Interrupted persist for instance {InstanceId}: already in terminal status {Status}/{SubStatus}
 19567                instance.Id, instance.Status, instance.SubStatus);
 19568            return;
 569        }
 570
 38571        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.
 38582            if (instance.Status != WorkflowStatus.Finished)
 583            {
 32584                var marked = await instanceStore.TryMarkInterruptedAsync(instance.Id, cancellationToken);
 31585                if (!marked)
 586                {
 0587                    _logger.LogInformation(
 0588                        "Skipping Interrupted persist for instance {InstanceId}: a concurrent persist already left it in
 0589                        instance.Id);
 0590                    return;
 591                }
 592            }
 37593        }
 1594        catch (Exception ex) when (!ex.IsFatal())
 595        {
 1596            _logger.LogWarning(ex, "Failed to persist Interrupted sub-status for instance {InstanceId}; falling back to 
 1597            actualReason = WorkflowInterruptedPayload.ReasonPersistenceFailure;
 1598        }
 599
 38600        var payload = new WorkflowInterruptedPayload(
 38601            InterruptedAt: _clock.UtcNow,
 38602            Reason: actualReason,
 38603            GenerationId: generationId,
 38604            LastActivityId: null,
 38605            LastActivityNodeId: null,
 38606            IngressSourceName: handle.IngressSourceName,
 38607            ExecutionCycleDuration: _clock.UtcNow - handle.StartedAt);
 608
 609        try
 610        {
 38611            await logStore.LogWorkflowInterruptedAsync(_identityGenerator, instance, payload, cancellationToken);
 38612        }
 0613        catch (Exception ex) when (!ex.IsFatal())
 614        {
 0615            _logger.LogWarning(ex, "Failed to write WorkflowInterrupted log entry for instance {InstanceId}.", instance.
 0616        }
 59617    }
 618
 619    private IReadOnlyList<IngressSourceFinalState> BuildSourceFinalStates()
 620    {
 35621        return _registry.Snapshot()
 23622            .Select(s => new IngressSourceFinalState(s.Name, s.State, s.LastError, WasForceStopped: s.State == IngressSo
 35623            .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    {
 57634        if (instance.Status != WorkflowStatus.Finished)
 32635            return false;
 636
 25637        if (instance.SubStatus == WorkflowSubStatus.Cancelled)
 25638            return !drainInducedInstanceIds.Contains(instance.Id);
 639
 0640        return true;
 641    }
 642}