| | | 1 | | using Elsa.Common; |
| | | 2 | | |
| | | 3 | | namespace Elsa.Workflows.Runtime; |
| | | 4 | | |
| | | 5 | | /// <summary> |
| | | 6 | | /// Tracks a single in-flight execution cycle of the workflow runtime. Created when a cycle starts, disposed when |
| | | 7 | | /// it completes or is force-cancelled during drain. Active-cycle accounting is done through |
| | | 8 | | /// <see cref="IExecutionCycleRegistry"/>. |
| | | 9 | | /// </summary> |
| | | 10 | | public sealed class ExecutionCycleHandle : IDisposable |
| | | 11 | | { |
| | | 12 | | private readonly CancellationTokenSource _cycleCts; |
| | | 13 | | private readonly CancellationTokenRegistration _linkedTokenRegistration; |
| | | 14 | | private readonly Action<ExecutionCycleHandle>? _onDisposed; |
| | | 15 | | private readonly Action? _cancelCallback; |
| | 692 | 16 | | private readonly TaskCompletionSource _disposedTcs = new(TaskCreationOptions.RunContinuationsAsynchronously); |
| | 692 | 17 | | private readonly object _cycleCtsGate = new(); |
| | | 18 | | private int _cycleCtsCancellationInProgress; |
| | | 19 | | private bool _cycleCtsDisposeRequested; |
| | | 20 | | private bool _cycleCtsDisposed; |
| | | 21 | | private int _lifecycleState; |
| | | 22 | | |
| | | 23 | | private const int ActiveState = 0; |
| | | 24 | | private const int CancellingState = 1; |
| | | 25 | | private const int CancelledState = 2; |
| | | 26 | | private const int DisposedState = 3; |
| | | 27 | | |
| | | 28 | | /// <summary> |
| | | 29 | | /// Creates a new handle. The owning <see cref="IExecutionCycleRegistry"/> supplies <paramref name="onDisposed"/> |
| | | 30 | | /// so it can decrement its active count. The optional <paramref name="cancelCallback"/> is invoked by |
| | | 31 | | /// <see cref="Cancel"/> to propagate cancellation into the workflow execution itself (e.g., |
| | | 32 | | /// <c>WorkflowExecutionContext.Cancel()</c>) so that the running cycle stops scheduling new activities. Without it, |
| | | 33 | | /// <see cref="Cancel"/> only signals the cycle's own <see cref="CancellationToken"/>, which is consumed by the |
| | | 34 | | /// execution-cycle registry but does NOT propagate to the workflow runner's pipeline (the runner reads from |
| | | 35 | | /// <c>WorkflowExecutionContext.CancellationToken</c>, which is captured at context construction and is not part |
| | | 36 | | /// of this linked chain). |
| | | 37 | | /// </summary> |
| | 692 | 38 | | public ExecutionCycleHandle( |
| | 692 | 39 | | Guid id, |
| | 692 | 40 | | string workflowInstanceId, |
| | 692 | 41 | | string? ingressSourceName, |
| | 692 | 42 | | DateTimeOffset startedAt, |
| | 692 | 43 | | CancellationToken linkedToken, |
| | 692 | 44 | | Action<ExecutionCycleHandle>? onDisposed = null, |
| | 692 | 45 | | Action? cancelCallback = null) |
| | | 46 | | { |
| | 692 | 47 | | Id = id; |
| | 692 | 48 | | WorkflowInstanceId = workflowInstanceId; |
| | 692 | 49 | | IngressSourceName = ingressSourceName; |
| | 692 | 50 | | StartedAt = startedAt; |
| | 692 | 51 | | _cycleCts = new CancellationTokenSource(); |
| | 692 | 52 | | _linkedTokenRegistration = linkedToken.UnsafeRegister( |
| | 1 | 53 | | static state => ((ExecutionCycleHandle)state!).PropagateLinkedCancellation(), |
| | 692 | 54 | | this); |
| | 692 | 55 | | _onDisposed = onDisposed; |
| | 692 | 56 | | _cancelCallback = cancelCallback; |
| | 692 | 57 | | } |
| | | 58 | | |
| | | 59 | | /// <summary>Unique identifier for this execution cycle within the current runtime generation.</summary> |
| | 1619 | 60 | | public Guid Id { get; } |
| | | 61 | | |
| | | 62 | | /// <summary>Workflow instance whose execution this cycle is driving.</summary> |
| | 414 | 63 | | public string WorkflowInstanceId { get; } |
| | | 64 | | |
| | | 65 | | /// <summary> |
| | | 66 | | /// Name of the ingress source that initiated this cycle, when attribution is available (null for direct API invocat |
| | | 67 | | /// </summary> |
| | 39 | 68 | | public string? IngressSourceName { get; } |
| | | 69 | | |
| | | 70 | | /// <summary>When the cycle started.</summary> |
| | 39 | 71 | | public DateTimeOffset StartedAt { get; } |
| | | 72 | | |
| | | 73 | | /// <summary>Cancellation token passed into the workflow execution pipeline for this cycle.</summary> |
| | 10 | 74 | | public CancellationToken CancellationToken => _cycleCts.Token; |
| | | 75 | | |
| | | 76 | | /// <summary> |
| | | 77 | | /// Completes after <see cref="Dispose"/> logically releases the handle and physically cleans up its linked CTS — |
| | | 78 | | /// i.e., when the workflow runner finishes the cycle (cleanly or via cancellation) and the middleware exits its |
| | | 79 | | /// <c>using</c> block. If cancellation callbacks are in flight, <see cref="Dispose"/> may return before this |
| | | 80 | | /// cleanup completes. The drain orchestrator awaits this with a timeout before persisting |
| | | 81 | | /// <see cref="WorkflowSubStatus.Interrupted"/>, ensuring its write happens AFTER any commit the runner emits in |
| | | 82 | | /// response to <see cref="Cancel"/>. |
| | | 83 | | /// </summary> |
| | 78 | 84 | | public Task Disposed => _disposedTcs.Task; |
| | | 85 | | |
| | | 86 | | /// <summary> |
| | | 87 | | /// Cancels the cycle — used by the drain orchestrator on deadline breach or operator force. |
| | | 88 | | /// Invokes the cancel callback (when supplied at construction) to propagate cancellation into the workflow |
| | | 89 | | /// execution, then cancels the cycle's own linked CTS. Safe to call multiple times; idempotent. |
| | | 90 | | /// </summary> |
| | 5 | 91 | | public void Cancel() => TryCancel(); |
| | | 92 | | |
| | | 93 | | /// <summary> |
| | | 94 | | /// Attempts to cancel the cycle. Returns <c>true</c> only when this call transitioned the handle from |
| | | 95 | | /// active to cancelled. Returns <c>false</c> when the handle was already disposed or already cancelling/cancelled, |
| | | 96 | | /// so drain can avoid treating a finished cycle as a force-cancel. Disposal wins if it races with the cancellation |
| | | 97 | | /// callback, so a cycle that completes while cancellation is in flight is not reported as drain-cancelled. |
| | | 98 | | /// </summary> |
| | | 99 | | public bool TryCancel() |
| | | 100 | | { |
| | 80 | 101 | | if (Interlocked.CompareExchange(ref _lifecycleState, CancellingState, ActiveState) != ActiveState) |
| | 13 | 102 | | return false; |
| | | 103 | | |
| | | 104 | | // Propagate to the workflow execution first (this typically marks the workflow as Cancelled and clears its |
| | | 105 | | // schedule, so the runner stops scheduling new activities). The orchestrator's subsequent Interrupted |
| | | 106 | | // persistence then overrides the Cancelled sub-status — see Disposed-await sequencing in DrainOrchestrator. |
| | 133 | 107 | | try { _cancelCallback?.Invoke(); } |
| | 2 | 108 | | catch (Exception ex) when (!ex.IsFatal()) { /* Cancellation is best-effort; non-fatal failures here must not bre |
| | | 109 | | |
| | 67 | 110 | | PropagateCycleCtsCancellation(); |
| | | 111 | | |
| | | 112 | | // Publish cancellation only after its effects complete. Dispose can transition CancellingState directly to |
| | | 113 | | // DisposedState, making this CAS fail when the cycle completed during the callback or CTS cancellation. |
| | 65 | 114 | | return Interlocked.CompareExchange(ref _lifecycleState, CancelledState, CancellingState) == CancellingState; |
| | | 115 | | } |
| | | 116 | | |
| | | 117 | | /// <summary> |
| | | 118 | | /// Logically releases the handle and notifies the registry. If cancellation callbacks are in flight, this method |
| | | 119 | | /// may return before physical cleanup of the linked CTS completes; <see cref="Disposed"/> is signaled afterwards. |
| | | 120 | | /// </summary> |
| | | 121 | | public void Dispose() |
| | | 122 | | { |
| | | 123 | | while (true) |
| | | 124 | | { |
| | 635 | 125 | | var state = Volatile.Read(ref _lifecycleState); |
| | 639 | 126 | | if (state == DisposedState) return; |
| | 631 | 127 | | if (Interlocked.CompareExchange(ref _lifecycleState, DisposedState, state) == state) break; |
| | | 128 | | } |
| | | 129 | | |
| | 631 | 130 | | _onDisposed?.Invoke(this); |
| | 631 | 131 | | RequestCycleCtsDisposal(); |
| | 631 | 132 | | } |
| | | 133 | | |
| | | 134 | | private void PropagateLinkedCancellation() |
| | | 135 | | { |
| | 1 | 136 | | PropagateCycleCtsCancellation(); |
| | 1 | 137 | | } |
| | | 138 | | |
| | | 139 | | private void PropagateCycleCtsCancellation() |
| | | 140 | | { |
| | 68 | 141 | | lock (_cycleCtsGate) |
| | | 142 | | { |
| | 68 | 143 | | if (_cycleCtsDisposed) |
| | 3 | 144 | | return; |
| | | 145 | | |
| | 65 | 146 | | _cycleCtsCancellationInProgress++; |
| | 65 | 147 | | } |
| | | 148 | | |
| | | 149 | | try |
| | | 150 | | { |
| | 65 | 151 | | _cycleCts.Cancel(); |
| | 62 | 152 | | } |
| | 0 | 153 | | catch (ObjectDisposedException) |
| | | 154 | | { |
| | | 155 | | // Dispose may have won before cancellation propagation started. |
| | 0 | 156 | | } |
| | 3 | 157 | | catch (Exception ex) when (!ex.IsFatal()) |
| | | 158 | | { |
| | | 159 | | // CTS callbacks are best-effort; preserve the lifecycle transition even when one reports a non-fatal error. |
| | 1 | 160 | | } |
| | | 161 | | finally |
| | | 162 | | { |
| | 65 | 163 | | var dispose = false; |
| | 65 | 164 | | lock (_cycleCtsGate) |
| | | 165 | | { |
| | 65 | 166 | | _cycleCtsCancellationInProgress--; |
| | 65 | 167 | | if (_cycleCtsDisposeRequested && _cycleCtsCancellationInProgress == 0 && !_cycleCtsDisposed) |
| | | 168 | | { |
| | 3 | 169 | | _cycleCtsDisposed = true; |
| | 3 | 170 | | dispose = true; |
| | | 171 | | } |
| | 65 | 172 | | } |
| | | 173 | | |
| | 65 | 174 | | if (dispose) |
| | 3 | 175 | | DisposeCycleCts(); |
| | 65 | 176 | | } |
| | 66 | 177 | | } |
| | | 178 | | |
| | | 179 | | private void RequestCycleCtsDisposal() |
| | | 180 | | { |
| | 631 | 181 | | var dispose = false; |
| | 631 | 182 | | lock (_cycleCtsGate) |
| | | 183 | | { |
| | 631 | 184 | | _cycleCtsDisposeRequested = true; |
| | 631 | 185 | | if (_cycleCtsCancellationInProgress == 0 && !_cycleCtsDisposed) |
| | | 186 | | { |
| | 628 | 187 | | _cycleCtsDisposed = true; |
| | 628 | 188 | | dispose = true; |
| | | 189 | | } |
| | 631 | 190 | | } |
| | | 191 | | |
| | 631 | 192 | | if (dispose) |
| | 628 | 193 | | DisposeCycleCts(); |
| | 631 | 194 | | } |
| | | 195 | | |
| | | 196 | | private void DisposeCycleCts() |
| | | 197 | | { |
| | | 198 | | // CancellationTokenRegistration.Dispose is self-unregister-safe when this is called from the linked |
| | | 199 | | // token callback, and waits for a callback running on another thread before releasing the registration. |
| | 631 | 200 | | _linkedTokenRegistration.Dispose(); |
| | 631 | 201 | | _cycleCts.Dispose(); |
| | 631 | 202 | | _disposedTcs.TrySetResult(); |
| | 631 | 203 | | } |
| | | 204 | | } |