< Summary

Information
Class: Elsa.Bpmn.Hosting.BpmnScopeHost
Assembly: Elsa.Bpmn
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Bpmn/Hosting/BpmnScopeHost.cs
Line coverage
95%
Covered lines: 136
Uncovered lines: 6
Coverable lines: 142
Total lines: 387
Line coverage: 95.7%
Branch coverage
89%
Covered branches: 52
Total branches: 58
Branch coverage: 89.6%
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%
For(...)100%11100%
DispatcherOf(...)100%11100%
get_Graph()100%11100%
StartAsync()100%11100%
OnWorkCompletedAsync(...)75%4494.11%
OnScopeSignalledAsync(...)66.66%6687.5%
OnWorkFaultedAsync()100%66100%
EvaluateAsync(...)100%1195.23%
ApplyAsync()83.33%6685.71%
ProjectDiagnostics(...)93.75%171686.36%
SeedCursorFrom(...)100%88100%
TryGetDiagnosticSequence(...)100%22100%
ResolveFaultedWork(...)100%66100%
Snapshot(...)100%11100%
get_InvocationCorrelation()100%22100%
BuildGraph()50%22100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Bpmn/Hosting/BpmnScopeHost.cs

#LineLine coverage
 1using System.Globalization;
 2using Bpmn.Model;
 3using Bpmn.Model.State;
 4using Bpmn.Semantics;
 5using Elsa.Bpmn.Activities;
 6using Elsa.Bpmn.Exceptions;
 7using Elsa.Bpmn.Signals;
 8using Elsa.Extensions;
 9using Elsa.Workflows;
 10using Elsa.Workflows.Activities.Flowchart.Models;
 11using Elsa.Workflows.Signals;
 12using Microsoft.Extensions.Logging;
 13
 14namespace Elsa.Bpmn.Hosting;
 15
 16/// <summary>
 17/// The host side of the <c>Bpmn.Semantics</c> port for one BPMN scope: it feeds the interpreter's four entry points
 18/// and applies what comes back onto the scope's <see cref="ActivityExecutionContext"/>.
 19/// </summary>
 20/// <remarks>
 21/// <para>
 22/// Every entry point is synchronous and returns a value; the interpreter never calls back. The host's job is to say
 23/// what its world looks like — a snapshot — and then to do what it is told, in the order it is told.
 24/// </para>
 25/// <para>
 26/// A host instance is a view over one scope's context and is created per call. Everything durable lives in
 27/// <see cref="BpmnScopeMemory"/>, everything derived lives in the context's transient properties, and everything
 28/// ordering-related lives in the instance-wide <see cref="BpmnScopeDispatcher"/>.
 29/// </para>
 30/// </remarks>
 31internal sealed class BpmnScopeHost
 32{
 33    /// <summary>
 34    /// What this host promises it can do.
 35    /// </summary>
 36    /// <remarks>
 37    /// Each capability is a claim, honoured elsewhere in this module: subtree cancellation by
 38    /// <see cref="BpmnWorkTeardown"/>, scope signalling by <see cref="BpmnScopeSignal"/>, iteration scopes by the
 39    /// applier's per-instance variables, and <see cref="BpmnHostCapabilities.ScopeVariables"/> by
 40    /// <see cref="BpmnScopeVariables"/>. Defined in terms of <see cref="BpmnRuntimeCapabilities.Declared"/> — see its
 41    /// remarks for why the set is spelled out there rather than written as <c>BpmnHostCapabilities.Full</c>, and for
 42    /// why that type exists at all. The same set goes to <see cref="BpmnGraph.Build"/> and to every snapshot, which
 43    /// the port requires.
 44    /// </remarks>
 45    public const BpmnHostCapabilities Capabilities = BpmnRuntimeCapabilities.Declared;
 46
 47    /// <summary>
 48    /// The property key under which a nested scope's invocation correlation is carried on its own context.
 49    /// </summary>
 50    public const string InvocationCorrelationPropertyKey = "Bpmn:InvocationCorrelation";
 51
 52    private const string GraphTransientPropertyKey = "Bpmn:Graph";
 53
 54    // The interpreter is a pure function of its request: it holds no per-instance state, so one instance serves the
 55    // whole process. Creating one per evaluation would only re-register the built-in element behaviors.
 256    private static readonly BpmnInterpreter Interpreter = BpmnInterpreter.CreateDefault();
 57
 258    private static readonly object DispatcherKey = new();
 259    private static readonly IReadOnlyDictionary<string, string> NoCorrelation = new Dictionary<string, string>(StringCom
 60
 61    private readonly ActivityExecutionContext _context;
 62    private readonly BpmnProcess _process;
 63
 21264    private BpmnScopeHost(ActivityExecutionContext context)
 65    {
 21266        _context = context;
 21267        _process = (BpmnProcess)context.Activity;
 21268    }
 69
 70    /// <summary>Returns the host for the given BPMN scope context.</summary>
 21271    public static BpmnScopeHost For(ActivityExecutionContext context) => new(context);
 72
 73    /// <summary>The evaluation queue shared by every BPMN scope in this workflow instance.</summary>
 74    public static BpmnScopeDispatcher DispatcherOf(WorkflowExecutionContext context) =>
 28375        context.TransientProperties.GetOrAdd(DispatcherKey, () => new BpmnScopeDispatcher());
 76
 77    /// <summary>The built graph for this scope. Derived from the definition, the bound work and the capabilities, none 
 20078    public BpmnGraph Graph => _context.TransientProperties.GetOrAdd(GraphTransientPropertyKey, BuildGraph);
 79
 80    // --- The interpreter's four entry points ---------------------------------------------------------
 81
 82    /// <summary>The scope is beginning.</summary>
 6483    public ValueTask StartAsync() => EvaluateAsync(memory =>
 12884        Interpreter.Start(new BpmnStartRequest(Graph, memory.State, Snapshot(memory))));
 85
 86    /// <summary>A unit of work finished, reporting zero or more outcome names.</summary>
 12587    public ValueTask OnWorkCompletedAsync(ActivityExecutionContext childContext, object? result) => EvaluateAsync(memory
 12588    {
 12589        // Keyed on the child's own context id. A completion for work this scope no longer holds is absorbed rather
 12590        // than faulted: an interrupting boundary tears its host down while the host's work is still in flight, and a
 12591        // late completion for work that was torn down is an ordinary BPMN race.
 12592        if (memory.Work.FindByChildContextId(childContext.Id) is not { } record)
 193            return null;
 12594
 12595        // The completing work must ALREADY be gone from LiveWork when the interpreter is asked.
 12496        memory.Work.Remove(record);
 12497        memory.SaveWork();
 12598
 12499        var outcomeNames = result is Outcomes outcomes ? outcomes.Names : [];
 125100
 124101        return Interpreter.OnWorkCompleted(new BpmnWorkCompletedRequest(
 124102            Graph, memory.State, Snapshot(memory), record.BindingRef, record.Handle, outcomeNames, record.IterationId));
 125103    });
 104
 105    /// <summary>A nested scope this one invoked signalled outward.</summary>
 106    public ValueTask OnScopeSignalledAsync(BpmnScopeSignal signal, SignalContext signalContext)
 107    {
 8108        var sender = signalContext.SenderActivityExecutionContext;
 109
 110        // The channel delivers to the sender before walking its ancestors, and a scope never signals itself.
 8111        if (string.Equals(sender.Id, _context.Id, StringComparison.Ordinal))
 4112            return default;
 113
 114        // A scope signal is for the immediate enclosing scope, which is the one that started the sender's work. Any
 115        // other receiver lets it keep bubbling; that is also how an unrelated container in between composes.
 4116        if (BpmnScopeMemory.Load(_context).Work.FindByChildContextId(sender.Id) is not { } signalling)
 0117            return default;
 118
 4119        signalContext.StopPropagation();
 120
 4121        return EvaluateAsync(memory =>
 4122        {
 4123            // Unlike a completion, the signalling work stays in the ledger: an escalating activity keeps running, and
 4124            // removing it makes the interpreter believe it has already gone.
 4125            if (memory.Work.FindByHandle(signalling.Handle) is not { } record)
 0126                return null;
 4127
 4128            return Interpreter.OnWorkSignalled(new BpmnWorkSignalledRequest(
 4129                Graph, memory.State, Snapshot(memory), record.BindingRef, record.Handle, signal.Code, signal.Payload, re
 4130        });
 131    }
 132
 133    /// <summary>
 134    /// A unit of work failed. Rides the <see cref="FaultSignal"/> seam: handle the signal, ask the interpreter what
 135    /// BPMN makes of the fault, and claim it only when a catcher took it.
 136    /// </summary>
 137    /// <remarks>
 138    /// The disposition has to be decided before this handler returns, so the interpreter is asked inline; only the
 139    /// commands are applied through the dispatcher. A <c>Propagated</c> disposition is left strictly alone — no
 140    /// <c>StopPropagation</c>, nothing terminalized — so the fault reaches the enclosing scope, which is how BPMN
 141    /// error propagation crosses a scope boundary, or the incident strategy, which is how it surfaces at the root.
 142    /// </remarks>
 143    public async ValueTask OnWorkFaultedAsync(FaultSignal signal, SignalContext signalContext)
 144    {
 15145        var memory = BpmnScopeMemory.Load(_context);
 146
 15147        if (ResolveFaultedWork(memory, signal.FaultedContext) is not { } record)
 7148            return;
 149
 150        // As with a completion, the failed work must ALREADY be removed before the interpreter is asked.
 8151        memory.Work.Remove(record);
 8152        memory.SaveWork();
 153
 8154        var evaluation = Interpreter.OnWorkFaulted(new BpmnWorkFaultedRequest(
 8155            Graph, memory.State, Snapshot(memory), record.BindingRef, record.Handle, signal.Exception.Message));
 156
 8157        ProjectDiagnostics(evaluation.State, memory.State);
 158
 8159        memory.State = evaluation.State.Prune();
 8160        memory.SaveState();
 161
 8162        if (evaluation.Disposition is not BpmnErrorDisposition.Caught)
 2163            return;
 164
 6165        signalContext.StopPropagation();
 166
 167        // A handler that claims a fault owns terminalizing the failed activity, and BPMN terminalizes the whole unit of
 168        // work rather than the one activity that threw: when the failure came from inside a nested scope, this scope's
 169        // failing work is that scope. Cancelling it recursively covers the activity that actually threw. The interprete
 170        // issues no teardown for failed work — it treats it as already terminal — so this is the host's own doing and
 171        // not a command out of order. RecoverFromFault stays the middleware's alone: it decrements every ancestor's
 172        // fault count, so a second call drives them negative.
 6173        if (BpmnWorkTeardown.FindContext(_context.WorkflowExecutionContext, record.ChildContextId) is { } failedWorkCont
 6174            await BpmnWorkTeardown.CancelSubtreeAsync(failedWorkContext, $"element '{record.ElementId}' failed");
 175
 12176        await DispatcherOf(_context.WorkflowExecutionContext).PostAsync(() => ApplyAsync(memory, evaluation));
 15177    }
 178
 179    // --- Plumbing ------------------------------------------------------------------------------------
 180
 181    private ValueTask EvaluateAsync(Func<BpmnScopeMemory, BpmnEvaluation?> evaluate) =>
 193182        DispatcherOf(_context.WorkflowExecutionContext).PostAsync(async () =>
 193183        {
 193184            var memory = BpmnScopeMemory.Load(_context);
 193185            var evaluation = evaluate(memory);
 193186
 189187            if (evaluation is null)
 1188                return;
 193189
 193190            // Diagnostics are audit-only and capped: projecting from persisted state later would lose whatever
 193191            // Prune() already dropped, so this runs on the evaluation's own state, before pruning. memory.State is
 193192            // still what was loaded before this evaluation ran -- the prior state -- since it is not overwritten
 193193            // until after this call.
 188194            ProjectDiagnostics(evaluation.State, memory.State);
 193195
 193196            // Persist the state before acting on the commands: a command applied against a state that was never
 193197            // recorded is how a crash produces work with no token behind it.
 188198            memory.State = evaluation.State.Prune();
 188199            memory.SaveState();
 193200
 188201            await ApplyAsync(memory, evaluation);
 380202        });
 203
 204    private async ValueTask ApplyAsync(BpmnScopeMemory memory, BpmnEvaluation evaluation)
 205    {
 194206        await new BpmnCommandApplier(_context, _process, memory).ApplyAsync(evaluation.Commands);
 207
 193208        switch (evaluation.Continuation)
 209        {
 210            case BpmnContinuation.Complete complete:
 211                // The scope completes because the interpreter said so, never because it ran out of children.
 51212                await _context.CompleteActivityAsync(new Outcomes(complete.Outcome));
 51213                break;
 214            case BpmnContinuation.Defer:
 215                break;
 216            case BpmnContinuation.Fault fault:
 1217                throw new BpmnScopeFaultException(fault.Code, fault.Message);
 218            default:
 0219                throw new NotSupportedException($"The BPMN continuation '{evaluation.Continuation.GetType().Name}' is no
 220        }
 192221    }
 222
 223    /// <summary>
 224    /// Projects every diagnostic the interpreter has appended since the last evaluation onto this scope's own
 225    /// execution log, keyed by element id. Under Option A only bound work has an activity id, so a gateway, an
 226    /// intermediate event or a sequence flow has nothing else in the journal to say where a token went; this is
 227    /// write-only and never read back by the interpreter or this host.
 228    /// </summary>
 229    /// <remarks>
 230    /// Runs on <see cref="_context"/> — this scope's own context — and never a child's: the diagnostic describes
 231    /// this scope's decision about a child, and the child may already be torn down by the time this runs. Called
 232    /// with the evaluation's own <see cref="BpmnEvaluation.State"/>, before <c>Prune()</c> caps
 233    /// <see cref="BpmnExecutionState.Diagnostics"/> at 200 entries, because projecting from what was actually
 234    /// persisted would lose whatever pruning already dropped. The last diagnostic id it has projected is kept in
 235    /// <see cref="BpmnScopeMemory.DiagnosticsCursorPropertyKey"/> so a resumed scope does not re-emit one a
 236    /// previous evaluation already turned into a journal entry.
 237    /// </remarks>
 238    /// <param name="state">The evaluation's own state, not yet pruned.</param>
 239    /// <param name="priorState">
 240    /// The state this scope had persisted before this evaluation ran, or <c>null</c> for a scope that has never
 241    /// been evaluated before. When the cursor property is absent -- a scope persisted before diagnostics projection
 242    /// existed -- its diagnostics are already accounted for, not new: the cursor is seeded from the highest valid
 243    /// sequence among <paramref name="priorState"/>'s own diagnostics before anything is projected, so only what
 244    /// this evaluation produced gets journaled. A genuinely new scope has no prior state and still starts at zero.
 245    /// </param>
 246    private void ProjectDiagnostics(BpmnExecutionState state, BpmnExecutionState? priorState)
 247    {
 196248        var storedCursor = BpmnScopeMemory.Read<BpmnDiagnosticsCursor>(_context, BpmnScopeMemory.DiagnosticsCursorProper
 196249        var lastProjectedSequence = storedCursor ?? SeedCursorFrom(priorState);
 196250        var highWaterMark = lastProjectedSequence;
 251
 3530252        foreach (var diagnostic in state.Diagnostics)
 253        {
 1569254            if (!TryGetDiagnosticSequence(diagnostic.DiagnosticId, out var sequence))
 255            {
 0256                _context.GetRequiredService<ILogger<BpmnScopeHost>>()
 0257                    .LogWarning("BPMN diagnostic id '{DiagnosticId}' is not in the expected 'diag:N' format and was skip
 0258                continue;
 259            }
 260
 1569261            if (sequence <= lastProjectedSequence)
 262                continue;
 263
 595264            highWaterMark = Math.Max(highWaterMark, sequence);
 265
 266            // Only a diagnostic that names neither an element nor a flow is scope-level and already journaled as the
 267            // activity's own lifecycle: the terminal "Completed" summary. Everything else -- including a start
 268            // event's own token emission, and a boundary event's token emission when it has no inbound flow (an
 269            // error or cancel boundary fires without one) -- names an element or a flow and is projected.
 595270            if (string.IsNullOrEmpty(diagnostic.ElementId) && string.IsNullOrEmpty(diagnostic.FlowId))
 271                continue;
 272
 544273            var payload = new BpmnDiagnosticLogPayload(
 544274                diagnostic.DiagnosticId,
 544275                diagnostic.ElementId,
 544276                diagnostic.FlowId,
 544277                diagnostic.TokenId,
 544278                diagnostic.Kind.ToString(),
 544279                diagnostic.Details);
 280
 544281            _context.AddExecutionLogEntry(diagnostic.Kind.ToString(), diagnostic.Message, BpmnDiagnosticEventNames.Sourc
 282        }
 283
 196284        if (highWaterMark != lastProjectedSequence)
 196285            BpmnScopeMemory.Write(_context, BpmnScopeMemory.DiagnosticsCursorPropertyKey, new BpmnDiagnosticsCursor(high
 196286    }
 287
 288    /// <summary>
 289    /// The starting cursor for a scope that has no <see cref="BpmnScopeMemory.DiagnosticsCursorPropertyKey"/> yet:
 290    /// the highest valid sequence already present in <paramref name="priorState"/>'s diagnostics, or zero when there
 291    /// is no prior state at all. A scope persisted before diagnostics projection existed has diagnostics in its
 292    /// state but no cursor; treating that absence as zero would make its next evaluation journal every one of those
 293    /// already-historical diagnostics as if they were new.
 294    /// </summary>
 295    private static int SeedCursorFrom(BpmnExecutionState? priorState)
 296    {
 61297        var highest = 0;
 298
 61299        if (priorState is null)
 60300            return highest;
 301
 14302        foreach (var diagnostic in priorState.Diagnostics)
 303        {
 6304            if (TryGetDiagnosticSequence(diagnostic.DiagnosticId, out var sequence) && sequence > highest)
 6305                highest = sequence;
 306        }
 307
 1308        return highest;
 309    }
 310
 311    /// <summary>
 312    /// The numeric ordinal in a diagnostic id (<c>diag:N</c>) — a pure function of the interpreter's own
 313    /// mutation-order sequence, so it sorts the same as arrival order. Never throws: an id that does not match the
 314    /// expected format fails to parse rather than faulting the evaluation, since this runs on every evaluation.
 315    /// </summary>
 316    /// <remarks>
 317    /// Requires the exact ordinal prefix <c>diag:</c> and a non-negative integer suffix with no leading sign, digit
 318    /// grouping or surrounding whitespace: a malformed id that happened to parse as a large number would poison the
 319    /// durable cursor and silently drop every later, genuinely valid, lower-sequence diagnostic forever.
 320    /// </remarks>
 321    internal static bool TryGetDiagnosticSequence(string diagnosticId, out int sequence)
 322    {
 323        const string prefix = "diag:";
 324
 1584325        if (diagnosticId.StartsWith(prefix, StringComparison.Ordinal))
 1582326            return int.TryParse(diagnosticId.AsSpan(prefix.Length), NumberStyles.None, CultureInfo.InvariantCulture, out
 327
 2328        sequence = 0;
 2329        return false;
 330    }
 331
 332    /// <summary>
 333    /// Finds the unit of work this scope started that the failing activity belongs to, walking outward from the
 334    /// failure.
 335    /// </summary>
 336    /// <remarks>
 337    /// A fault raised deep inside a nested scope is, to this scope, its own subprocess work failing. The nested scope
 338    /// sees the signal first and claims it if it has a catcher; if it does not, the signal arrives here and this walk
 339    /// is what turns "some activity failed" into "the work I started failed", which is exactly what BPMN error
 340    /// propagation across a scope boundary means.
 341    /// </remarks>
 342    private BpmnWorkRecord? ResolveFaultedWork(BpmnScopeMemory memory, ActivityExecutionContext faultedContext)
 343    {
 38344        for (var current = faultedContext; current is not null && !string.Equals(current.Id, _context.Id, StringComparis
 345        {
 12346            if (memory.Work.FindByChildContextId(current.Id) is { } record)
 8347                return record;
 348        }
 349
 7350        return null;
 351    }
 352
 353    private BpmnHostSnapshot Snapshot(BpmnScopeMemory memory)
 354    {
 196355        var invocationCorrelation = InvocationCorrelation;
 356
 196357        return new BpmnHostSnapshot(
 196358            ScopeInstanceId: _context.Id,
 196359            // A scope has an enclosing one exactly when another scope started it, which is what the carried
 196360            // correlation records. A root process has none, so an unhandled escalation is a documented no-op.
 196361            HasEnclosingScope: invocationCorrelation.Count > 0,
 196362            LiveWork: memory.Work.ToLiveWork(),
 196363            InvocationCorrelation: invocationCorrelation,
 196364            Variables: new BpmnScopeVariables(_context),
 196365            Capabilities: Capabilities);
 366    }
 367
 368    /// <summary>
 369    /// The correlation of the work that started this scope. It belongs to the scope and is fixed for its lifetime;
 370    /// a completing unit of work's correlation is never written here, because the event-subprocess start hint is read
 371    /// from this same dictionary.
 372    /// </summary>
 373    private IReadOnlyDictionary<string, string> InvocationCorrelation =>
 196374        BpmnScopeMemory.Read<Dictionary<string, string>>(_context, InvocationCorrelationPropertyKey) ?? NoCorrelation;
 375
 376    private BpmnGraph BuildGraph()
 377    {
 112378        var definition = _process.Process
 112379                         ?? throw new InvalidOperationException($"BPMN process activity '{_process.Id}' has no process d
 380
 381        // Every binding the definition declares, with the nested definition attached where the bound activity is
 382        // itself a BPMN scope. The graph validator reads that for an event subprocess body's start trigger.
 381383        var boundWork = BpmnBoundWork.Derive(definition, bindingRef => (_process.FindWorkActivity(bindingRef) as BpmnPro
 384
 112385        return BpmnGraph.Build(definition, boundWork, Capabilities);
 386    }
 387}