Skip to content

Control Flow in Flow Graphs ​

Design record for the flow-graph control-flow constructs: branch, loops (while / until / count / for-each, sequential and parallel), run-scoped variables (Set Variables), and loop break (body → outside edges).

Implementation notes (locked decisions) ​

The loop engine was designed against an adversarial correctness review before coding; the decisions that shaped the shipped implementation:

  • Execution budget (3 layers). Per-run node-execution cap is per-flow configurable (default 750, hard ceiling 2000); validation estimates the multiplicative nested-loop budget (Π of enclosing max_iterations per node, summed) against the flow's effective cap (its configured budget clamped to the ceiling, else 750) — warning above 80% of the cap and rejecting above it; a runtime history-event estimator fails the run cleanly before Temporal's hard history-size termination.
  • max_iterations (min 1) defaults to 100 for while/until; count and for_each default to their natural bound (the count / item total), so an explicit value is only ever an authored cap. All gated by the budget.
  • Loop UI = resizable container frame with internal start/end ports. The loop head renders as a group-style frame the body nodes are dropped into (children via parentId/scope_id). The iterate outcome is an internal "start" port inside the frame and the paired loop_end renders as a compact docked "end" port (auto-created with the loop, loop_ref preset, pinned non-draggable to the inner bottom-right corner — the mirror of the start port at the top-left), so the body reads start → … → end within the container. Only the exits (done/max_iterations, right) and the input (left) live outside, centered on the frame's sides, and containers are excluded from smart-edge obstacles. Visual containment is cosmetic only: body membership is defined by iterate→loop_end reachability, never scope_id.
  • Break = a body edge leaving the loop. Any edge from a body node to a node outside the body exits the loop early: when the token is dequeued at the outside node while still carrying the loop's frame, the engine finalizes the loop (StepRun closes SUCCESS with exit_reason: "break", frame popped, loop state cleared) before the node executes — innermost loop first when breaking out of nesting. Break is rejected when the body contains a fork (a sibling branch could still be running when the token leaves; the engine also fails the run if body work is in flight at break time as defense in depth). on_body_failure = break/continue modes remain deferred — a failure edge wired outside the body already gives break-on-failure explicitly.
  • Loop boundary rules (commit-gated validation): a loop body may only be entered through its head (no outside → body edges); a body node must not edge back to its own head (the loop end is the only continue path — a direct edge would repeat the iteration without advancing the index); two loop bodies must be disjoint or fully nested (overlap computes as mutual head containment and is rejected); and a loop head must not be reachable from two parallel branches of the same fork before a join consolidates them (loop state is per-head, one active pass at a time — a pass being a single iteration by default, or the max_parallel_iterations iterations a count/for_each loop may keep in flight at once).
  • In a loop body, only all-mode + wait joins are allowed (any / fail_fast / ignore_failures joins and until_join monitors are rejected), so a single consolidated token at loop_end guarantees the body has drained before the next iteration re-arms it. The engine also asserts body quiescence before resetting.
  • Determinism. Condition/items/collect and set_vars assignments evaluate in activities (results in workflow history). run.vars is authoritative in workflow memory (never re-read from the DB); the DB snapshot is best-effort for the UI. Body reset is a command-free set of in-memory pops (order-safe).
  • continue-as-new for very long loops is deferred; the budget guards keep worst-case history below the termination zone. The state model (run_vars, loop_states, token loop_frames) is kept CAN-ready.

Motivation ​

Before this design, flow graphs expressed parallelism well (fork/join) and outcome-based routing narrowly (a step's success/failure/cancelled handles). Anything richer — "branch three ways on a value", "retry until healthy", "do this for every device in a list" — was worked around with scripts in containers, which hid logic from the canvas, from run visualization, and from the audit trail.

The design introduced first-class control-flow constructs that:

  • are visible on the canvas and in run visualization,
  • reuse the existing Jinja variable system ({{ steps.* }}, {{ inputs.* }}, {{ vars.* }}, {{ device.* }}),
  • compile onto the existing token/outcome/edge machinery — no second execution model,
  • keep the proof-of-execution ethos: every decision is recorded with the evaluated values that produced it.

Baseline: what the engine already gave us ​

This is the baseline the design started from — the architecture was already closer to control flow than it looked:

FoundationWhereWhy it matters
Jinja everywhere, resolved JIT worker-sideapps/worker/template_resolver.py, packages/core/templates/Conditions need no new expression language
Bare-Jinja-expression fields with secret()/env()/file() bannedFlowOutputDefinition in apps/api/schemas.pyExact precedent for condition fields
Edges keyed by (source, outcome), per-outcome handlesflow_graph_runtime.py, editor types.tsA branch is "just" a node with more outcomes
Cycles legal (validation warns only), guarded by max_total_node_executionsapps/api/routers/flows/lifecycle.pyBack-edges — the skeleton of a loop — already run
Per-node trigger_limit (1 = first-arrival-wins, 0 = unlimited)flow_graph_runtime.pyNode re-execution semantics already exist
Per-visit execution identity exec_key = {node_id}_{visit}graph_engine/state.pyIterations are naturally distinguishable
Step outputs fetched latest-wins per nodeapps/api/routers/internal.py (get_step_outputs_internal)Inside a loop, {{ steps.X }} reads the current iteration
Subflows (flow.run handler)worker handler registryFuture parallel per-item fan-out for free

What was missing was purely: (1) a way to evaluate a condition and route on it, (2) iteration bookkeeping (index/item context, re-arming, caps), and (3) mutable run-scoped state (counters, flags).

Design principles ​

  1. Logic lives in nodes, not edges. No hidden edge conditions. One node = one visible decision; run mode colors the fired outcome handle (existing usedOutcome mechanic).
  2. Conditions are Jinja expressions over the same context steps already use. They are evaluated worker-side in a small activity, so the result is recorded in Temporal history (deterministic replay) and persisted as a StepRun whose output contains the evaluated values — decisions become auditable evidence. secret(), env() and file() are rejected in condition expressions, as they already are in flow outputs.
  3. Constructs compile onto tokens/outcomes/edges. The scheduler loop, join analysis, pause/approval handling, and run visualization all keep working unchanged.

Construct 1: Branch node — if / elif / else and switch ​

The cornerstone. A branch node holds an ordered list of cases; at runtime the first case whose condition is true wins, and exactly one outcome fires.

Node model ​

yaml
type: branch
branch_mode: rules            # rules | value
cases:                        # ordered; first true wins
  - id: c1
    label: "disk low"
    when: "steps.check.output.free_gb | int < 10"
  - id: c2
    label: "no backup"
    when: "not steps.backup.output.ok"

value mode is the switch variant — one expression, compared per case:

yaml
type: branch
branch_mode: value
value_expr: "steps.detect.output.platform"
cases:
  - {id: c1, label: "IOS-XE",  equals: "iosxe"}
  - {id: c2, label: "NX-OS",   equals: "nxos"}

Semantics ​

AspectBehavior
Outcomescase:<id> per case, plus else, plus optional error
Token flowone token in → exactly one outcome out (like a step)
elsefires when no case matches; must be wired (validation error otherwise)
errorfires when an expression raises / references undefined values; if unwired, the run fails like a step failure
EvidenceStepRun with output = {chosen: "c1", evaluated: {c1: true, c2: null…}} (later cases after the winner are not evaluated — short-circuit, mirroring elif)
Convergenceno merge node needed: wire several cases to the same continuation node; only one path executes, and default trigger_limit=1 first-arrival-wins dedups anyway
Join analysisidentical to a step (multi-outcome, single token) — _compute_join_expected_counts unchanged

Canvas ​

  • Compact amber card (or diamond, like fork's purple diamond), with stacked labeled source handles on the right, in case order, and a grey Else handle at the bottom. Step nodes already render multiple outcome handles; the same pattern extends to N labeled handles.
  • Editor panel (BranchNodeFields): reorderable case rows — drag to reorder = reorder the elif chain — each row a label + expression input with the existing VariablePicker; a rules/value mode toggle swaps the row shape (when expression vs. equals value).
  • Run mode: fired case handle colored green, others grey (existing usedOutcome); the node detail drawer shows each case's evaluated value — "why did it go this way" answered directly in the run view.

If / elif / else mapping ​

Programming constructGraph shape
if A:branch with 1 case, else wired to the "skip" continuation
if A: ... elif B: ... else: ...branch with 2 ordered cases + else
switch x: case "a": ... case "b": ...branch in value mode
ternary "pick an input for the next step"usually better served by Jinja inside the step's params — the branch node is for routing, not value selection

Construct 2: Loop — while / until / count / for-each ​

Loops are a paired construct like fork/join: a loop head node and a loop_end tail node, linked by loop_end.loop_ref = <head id>. The back-edge is implicit — the engine routes tokens from loop_end back to the head; the editor draws it as a dashed edge automatically. Users never manage back edges or trigger_limit themselves.

Node model ​

yaml
type: loop
loop_mode: while              # while | until | count | for_each
condition: "run.vars.attempts | int < 5"     # while/until modes
count: 10                                    # count mode
items_expr: "steps.scan.output.devices"      # for_each mode
max_iterations: 100                          # hard cap, always enforced
on_body_failure: fail                        # fail (break | continue deferred, rejected at commit)
collect: "steps.apply.output.result"         # optional per-iteration gather
max_parallel_iterations: 1                   # count/for_each only; 1 (default) = sequential
parallel_failure_policy: drain               # drain (default) | cancel | continue

Semantics ​

AspectBehavior
Head outcomesiterate (enter body), done (condition false / items exhausted / break), max_iterations (optional; unwired = run fails)
while vs untilwhile checks before each iteration (0..N runs); until checks after (1..N runs, do-while)
loop_endno outgoing edges; arrival = continue (token returns to head). Break is a body-to-outside edge: routing a body token past the loop boundary finalizes the loop with exit_reason: "break"
Iteration contexttoken carries a loop-frame stack [{loop_id, index, item}]; body templates get {{ loop.index }} (0-based), {{ loop.item }} (for_each), {{ loop.count }} where known
Nested loopsinnermost frame wins for {{ loop.* }}; outer frames reachable as {{ loops.<loop_node_id>.index }}
Body re-armingtrigger/dedup bookkeeping keyed by (node_id, loop frame), so each iteration re-arms body nodes; per-visit exec_key already keeps executions distinct
steps.X inside bodylatest-wins semantics naturally read the current/previous iteration's output
collectthe head's own StepRun output becomes {results: [ ... ]} — downstream nodes read {{ steps.LOOP_ID.output.results }}
Failure in bodyon_body_failure: fail is the only accepted value today (fail the run); break / continue are deferred and rejected at commit
Evidenceone run event per iteration (graph.loop_iteration) + the head StepRun records iteration count, exit reason, collected results
Parallel iterationsmax_parallel_iterations > 1 (count/for_each only) makes the head a dispatcher: it starts iterations on their own tokens up to that width and re-queues itself while the window has room; loop_end retires an iteration by index and re-arms the head. collect values are stored per index, so results stay in iteration order however the iterations interleave. Effective concurrency is still bounded by the flow's max_parallelism window (default 20)
Parallel failureparallel_failure_policy: drain (default — stop dispatching, let in-flight iterations finish), cancel (cancel in-flight body activities), continue (run every iteration). All three fail the run at the end, reporting which iterations failed

Canvas ​

Shipped visual: the container frame (an earlier node-pair stage was the interim step and is now historical): reuse the existing group-node machinery (scope_id, resizable): the frame header is the loop head (icon, mode + condition summary, and a live run-mode badge "iteration 3 / 12" via SSE), done/max_iterations handles sit on the frame's right edge, and a small "↺ / break" dock inside-right acts as the loop_end. Body nodes live inside the frame for readability, but visual containment (scope_id) is cosmetic only — execution membership is computed from iterate→loop_end reachability, which validation enforces on the graph itself. (In the diagram above, the failure edge leaving the body directly is a break: loop_end itself has no outgoing edges — the dashed return is the engine handing the token back to the head.)

Editor panel (LoopNodeFields) ​

  • Mode select (While / Do-until / Count / For each) swapping the relevant field (condition / count / items expression), each an expression input with VariablePicker.
  • max_iterations, on_body_failure, optional collect expression.
  • max_parallel_iterations (count/for_each only; hidden for while/until) and, once above 1, the parallel_failure_policy select. The inter-iteration delay field hides while the loop runs iterations concurrently — the two are mutually exclusive and validation rejects the combination.
  • Inline hint showing what {{ loop.item }} / {{ loop.index }} will contain (for for_each, a preview of the first item when testing against a past run — see Expression UX below).

Guards and Temporal considerations ​

  • Per-loop max_iterations (default 100 for open-ended while/until; count/for_each default to their natural bound) and the global execution budget: max_total_node_executions is per-flow configurable (default 750, hard ceiling 2000) so legitimate loops don't starve big flows while still capping runaway graphs.
  • The engine's concurrency window, max_parallelism, is also per-flow configurable (default 20, ceiling 100; Settings → General → Advanced Options → Engine Concurrency). It bounds how many step activities one run keeps in flight — fork branches and parallel loop iterations draw on the same window — so a loop set to run 50 iterations at once still starts them 20 at a time under the default. Both graph-level settings are top-level keys on the graph; the legacy nested policy block is ignored.
  • Condition evaluation happens in an activity, so its result is in workflow history → deterministic replay, no sandbox issues.
  • Long loops grow Temporal history (one activity round-trip per decision + body steps). Document practical bounds (hundreds, not tens of thousands, of iterations); continue-as-new is a later optimization if real usage pushes the limit. Running iterations concurrently shortens wall-clock but does not reduce history — the per-run execution budget is unchanged either way.
  • Parallel-iteration restrictions (v1, commit-gated). With max_parallel_iterations > 1 the body may not contain fork/join nodes, nested loops, set_vars nodes, or steps whose handler declares node_scoped_state (tf.plan and tf.apply), and no edge may leave the body (break). Each restriction exists because the state involved is shared across iterations: join arrivals from different iterations would consolidate into one, a nested loop keeps a single per-head state, concurrent run.vars writes race, a node-scoped step's working state (tf.plan's directory, tf.apply's plan lookup) is keyed by node id, and a break would abandon iterations still running. Trigger bookkeeping for body nodes is keyed by (node_id, loop frame) so concurrent iterations are not absorbed as duplicate arrivals at the same node. Note also that {{ steps.X }} reads are latest-wins across iterations — under concurrency, read {{ loop.item }} / {{ loop.index }} and gather results with collect instead.

Construct 3: Set Variables node + run-scoped variables ​

while needs counters and flags; iteration pipelines need accumulators. Step outputs are immutable per execution, and {{ vars.* }} are global — neither fits. Proposal: run-scoped mutable variables.

yaml
type: set_vars
assignments:                  # evaluated in order, later sees earlier
  - {name: attempts, expr: "(run.vars.attempts | default(0) | int) + 1"}
  - {name: last_error, expr: "steps.upgrade.error | default('')"}
  • New per-run store (JSONB on the run row, size-capped, last-write-wins), injected into every template context as {{ run.vars.NAME }}, listed in the VariablePicker under a new scope, and shown on the run detail page.
  • Executed by the same expression-evaluation activity; recorded as a StepRun with before → after values (secrets are impossible here — secret() et al. are banned in expressions).
  • Small rectangle node, teal, one success + one error outcome; editor panel is an ordered name/expression list.

Typical retry-with-counter shape:

Construct 4 (superseded): parallel for-each via subflow map ​

Superseded by in-engine parallel iterations. Concurrency is now a property of the existing count/for_each loop — max_parallel_iterations on the head (see Construct 2) — rather than a separate loop_mode: map that fans out one subflow per item via flow.run.

The subflow-map idea remains on the table for a different problem: history containment for very large collections. In-engine iterations keep every iteration's activities in the parent workflow's history, so a run over thousands of items still pushes the history budget. A child workflow per item would move that cost out of the parent. Nothing in the current schema blocks adding it later.


Shared expression UX ​

One ExpressionInput component used by branch cases, loop conditions/items, and set_vars assignments:

  • Built on the existing StringInputWithPicker + VariablePicker (scopes from /variable-contexts, including the new run.vars scope).
  • Inline lint via a new POST /flows/validate-expression endpoint (parse-only: Jinja syntax + banned-helper check — the same rules FlowOutputDefinition enforces today).
  • "Test against a run…": pick a past run of this flow and evaluate the expression against its recorded context (the variable_contexts preview machinery already assembles exactly this data). Shows the value or the error before you commit the flow.

Validation (commit-time, mirrored in the editor) ​

  • Branch: ≥ 1 case; case ids unique; else edge required; expressions parse.
  • Loop: head and loop_end paired 1:1; body = nodes on head→loop_end paths; a body edge leaving the body is a break (allowed unless the body contains a fork); edges entering the body from outside, body → own-head edges, overlapping non-nested bodies, and concurrent fork entry into a head are errors; max_iterations ≥ 1.
  • set_vars: names match the existing variable-name pattern; expressions parse.
  • Cycle warning suppressed for implicit loop back-edges (they are the point); still emitted for hand-made cycles.
  • Existing per-node-type outcome lists (runtime node_outcomes, API validation, UI types.ts) each gain the new types — these three copies must stay in sync (they already must today).

Interaction with fork/join ​

  • A branch inside a fork branch is transparent: single token in/out, so join expected-count analysis treats it like a step.
  • A loop inside a fork branch works: the token carries its loop frames and eventually exits via done toward the join — one arrival, as today.
  • Fork/join inside a loop body is allowed as long as both ends are inside the body (containment validation); the join consolidates per iteration. This holds for sequential loops only — a loop with max_parallel_iterations > 1 rejects fork/join in its body, because join state is per-node and arrivals from different iterations would consolidate into a single join.
  • Edges crossing a loop boundary into the middle of a body are validation errors — the only way in is the head's iterate. The ways out are done / max_iterations and body → outside break edges (rejected when the body contains a fork, since a sibling branch could still be running).

Phasing ​

  1. Branch node + expression activity + ExpressionInput + validation endpoint. Covers if/elif/else and switch immediately. Also unlocks a documented power-user "while via back-edge" pattern (a branch case wired backward + trigger_limit: 0 on body steps) before real loops land.
  2. Loop pair (while / until / count / for-each sequential) + loop.* context + per-iteration re-arming + iteration run-visualization.
  3. Set Variables + run.vars (independent of 2; order can swap).
  4. Parallel iterations: max_parallel_iterations + parallel_failure_policy on count/for_each heads, with iteration-scoped body bookkeeping.
  5. Polish & scale: continue-as-new if needed, lifting the v1 parallel-body restrictions (fork/join, nested loops, set_vars, break), and the subflow map for history containment on very large collections — pending items are tracked in TODO.md. The loop container visual and the per-flow execution budget setting shipped with 1-3.

Decision record ​

DecisionOutcome
Loop visualResizable container frame with docked start/end ports
Run-vars namespace{{ run.vars.X }} — unambiguous next to global {{ vars.X }}
Parallel for-eachmax_parallel_iterations on the count/for_each head (in-engine), not a subflow map
Branch elseExplicit wiring required (unwired = validation error)

Affected components (for later implementation planning) ​

AreaFiles
Enums / schemaspackages/core/enums.py, apps/api/schemas.py
Validationapps/api/routers/flows/lifecycle.py
Graph runtimeapps/worker/flow_graph_runtime.py (node consts, outcomes, Token.loop_frames, join-analysis boundaries)
Engineapps/worker/graph_engine/engine.py (_dispatch_node + handlers, iteration-scoped trigger keys)
Activitiesapps/worker/flow_activities.py (evaluate_branch, evaluate_loop, evaluate_set_vars, persist_run_vars)
Template contextpackages/core/templates/resolver.py, apps/worker/template_resolver.py (loop.*, loops.*, run.vars.*)
API internalsapps/api/routers/internal.py (run-vars store, context payload)
UI nodesapps/ui/src/pages/FlowGraphEditor/ (BranchNode, LoopNode, LoopEndNode, SetVarsNode, editor panels, types.ts, palette)
UI sharedapps/ui/src/components/flow/ExpressionInput.tsx
Sync / YAMLapps/api/services/sync/serialization/
DocsControl-flow user guide (see TODO.md)

Released as open source under the AGPL-3.0-or-later license. Development is sponsored by Rexonix s.r.o.. Contact — [email protected].