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_iterationsper 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 forwhile/until;countandfor_eachdefault 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). Theiterateoutcome is an internal "start" port inside the frame and the pairedloop_endrenders as a compact docked "end" port (auto-created with the loop,loop_refpreset, 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, neverscope_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_iterationsiterations 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_endguarantees 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.varsis 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, tokenloop_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:
| Foundation | Where | Why it matters |
|---|---|---|
| Jinja everywhere, resolved JIT worker-side | apps/worker/template_resolver.py, packages/core/templates/ | Conditions need no new expression language |
Bare-Jinja-expression fields with secret()/env()/file() banned | FlowOutputDefinition in apps/api/schemas.py | Exact precedent for condition fields |
Edges keyed by (source, outcome), per-outcome handles | flow_graph_runtime.py, editor types.ts | A branch is "just" a node with more outcomes |
Cycles legal (validation warns only), guarded by max_total_node_executions | apps/api/routers/flows/lifecycle.py | Back-edges — the skeleton of a loop — already run |
Per-node trigger_limit (1 = first-arrival-wins, 0 = unlimited) | flow_graph_runtime.py | Node re-execution semantics already exist |
Per-visit execution identity exec_key = {node_id}_{visit} | graph_engine/state.py | Iterations are naturally distinguishable |
| Step outputs fetched latest-wins per node | apps/api/routers/internal.py (get_step_outputs_internal) | Inside a loop, {{ steps.X }} reads the current iteration |
Subflows (flow.run handler) | worker handler registry | Future 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
- Logic lives in nodes, not edges. No hidden edge conditions. One node = one visible decision; run mode colors the fired outcome handle (existing
usedOutcomemechanic). - 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()andfile()are rejected in condition expressions, as they already are in flow outputs. - 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
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:
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
| Aspect | Behavior |
|---|---|
| Outcomes | case:<id> per case, plus else, plus optional error |
| Token flow | one token in → exactly one outcome out (like a step) |
else | fires when no case matches; must be wired (validation error otherwise) |
error | fires when an expression raises / references undefined values; if unwired, the run fails like a step failure |
| Evidence | StepRun with output = {chosen: "c1", evaluated: {c1: true, c2: null…}} (later cases after the winner are not evaluated — short-circuit, mirroring elif) |
| Convergence | no 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 analysis | identical 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
Elsehandle 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 existingVariablePicker; a rules/value mode toggle swaps the row shape (whenexpression vs.equalsvalue). - 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 construct | Graph 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
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 | continueSemantics
| Aspect | Behavior |
|---|---|
| Head outcomes | iterate (enter body), done (condition false / items exhausted / break), max_iterations (optional; unwired = run fails) |
while vs until | while checks before each iteration (0..N runs); until checks after (1..N runs, do-while) |
loop_end | no 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 context | token 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 loops | innermost frame wins for {{ loop.* }}; outer frames reachable as {{ loops.<loop_node_id>.index }} |
| Body re-arming | trigger/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 body | latest-wins semantics naturally read the current/previous iteration's output |
collect | the head's own StepRun output becomes {results: [ ... ]} — downstream nodes read {{ steps.LOOP_ID.output.results }} |
| Failure in body | on_body_failure: fail is the only accepted value today (fail the run); break / continue are deferred and rejected at commit |
| Evidence | one run event per iteration (graph.loop_iteration) + the head StepRun records iteration count, exit reason, collected results |
| Parallel iterations | max_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 failure | parallel_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, optionalcollectexpression.max_parallel_iterations(count/for_each only; hidden for while/until) and, once above 1, theparallel_failure_policyselect. 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_executionsis 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 nestedpolicyblock 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-newis 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 > 1the body may not contain fork/join nodes, nested loops,set_varsnodes, or steps whose handler declaresnode_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, concurrentrun.varswrites 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 withcollectinstead.
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.
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 theVariablePickerunder 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+ oneerroroutcome; 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 newrun.varsscope). - Inline lint via a new
POST /flows/validate-expressionendpoint (parse-only: Jinja syntax + banned-helper check — the same rulesFlowOutputDefinitionenforces today). - "Test against a run…": pick a past run of this flow and evaluate the expression against its recorded context (the
variable_contextspreview 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;
elseedge required; expressions parse. - Loop: head and
loop_endpaired 1:1; body = nodes on head→loop_endpaths; 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, UItypes.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
donetoward 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 > 1rejects 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 aredone/max_iterationsand body → outside break edges (rejected when the body contains a fork, since a sibling branch could still be running).
Phasing
- 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: 0on body steps) before real loops land. - Loop pair (while / until / count / for-each sequential) +
loop.*context + per-iteration re-arming + iteration run-visualization. - Set Variables +
run.vars(independent of 2; order can swap). - Parallel iterations:
max_parallel_iterations+parallel_failure_policyon count/for_each heads, with iteration-scoped body bookkeeping. - 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
| Decision | Outcome |
|---|---|
| Loop visual | Resizable container frame with docked start/end ports |
| Run-vars namespace | {{ run.vars.X }} — unambiguous next to global {{ vars.X }} |
| Parallel for-each | max_parallel_iterations on the count/for_each head (in-engine), not a subflow map |
Branch else | Explicit wiring required (unwired = validation error) |
Affected components (for later implementation planning)
| Area | Files |
|---|---|
| Enums / schemas | packages/core/enums.py, apps/api/schemas.py |
| Validation | apps/api/routers/flows/lifecycle.py |
| Graph runtime | apps/worker/flow_graph_runtime.py (node consts, outcomes, Token.loop_frames, join-analysis boundaries) |
| Engine | apps/worker/graph_engine/engine.py (_dispatch_node + handlers, iteration-scoped trigger keys) |
| Activities | apps/worker/flow_activities.py (evaluate_branch, evaluate_loop, evaluate_set_vars, persist_run_vars) |
| Template context | packages/core/templates/resolver.py, apps/worker/template_resolver.py (loop.*, loops.*, run.vars.*) |
| API internals | apps/api/routers/internal.py (run-vars store, context payload) |
| UI nodes | apps/ui/src/pages/FlowGraphEditor/ (BranchNode, LoopNode, LoopEndNode, SetVarsNode, editor panels, types.ts, palette) |
| UI shared | apps/ui/src/components/flow/ExpressionInput.tsx |
| Sync / YAML | apps/api/services/sync/serialization/ |
| Docs | Control-flow user guide (see TODO.md) |