Control Flow
Branching, iteration, scoping, validation, and error recovery composites.
Control-flow blocks are composites: they embed sub-flows and are assembled by the flow builder rather than the leaf-block registry.
Composite blocks are configured with top-level YAML keys on the block (condition, then, cases, branches, items, body, ...), not under settings:. A composite that declares a slot it does not own fails at startup.
Each sub-flow slot (then, else, default, body, a fork branch, a switch case) has the same shape as a flow: an optional name and a process block list. Sub-flows must not declare source, workers, buffer, pool, or error, which are root-flow-only fields. See Flow File for how a block entry is written and for the whole slot vocabulary.
All condition and key fields are CEL expressions evaluated against the message (body, vars, eventID, correlationID, env, now). They are compiled once at startup, so a malformed expression fails deployment rather than a message.
if
Runs one of two sub-flows depending on a boolean condition. With no else, a false condition passes the message through unchanged.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
condition | expression (bool) | Yes | none | Boolean CEL expression evaluated against the message. Errors if the result is not a bool. |
then | flow | Yes | none | Sub-flow run when the condition is true. |
else | flow | No | none | Sub-flow run when the condition is false. Omitted: pass through. |
- type: if
name: any-orders
condition: "size(body.orders) > 0"
then:
process:
- type: log
settings:
message: '"processing " + string(size(body.orders)) + " orders"'
else:
process:
- type: log
settings:
message: '"no orders to process"'switch
Runs the sub-flow of the first case whose when guard is true, in order. When no case matches, the default flow runs; with no default, the message passes through unchanged.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
cases | list | Yes | none | Ordered case list. Each entry has a when (expression, bool, required) plus an inline flow (process, optional name). |
default | flow | No | none | Sub-flow run when no case matches. Omitted: pass through. |
- type: switch
name: classify-order
cases:
- when: "vars.order.amount >= vars.threshold"
process:
- type: log
settings:
message: '"HIGH order " + string(vars.order.id)'
default:
process:
- type: log
settings:
message: '"low order " + string(vars.order.id)'fork
Scatters the message across parallel branches, then joins. Each branch receives its own clone of the message and runs on the flow's shared worker pool (sized by the root flow's pool field), so branch effects on body and variables never leak back: on success the original input message passes through unchanged. Branch outputs are not aggregated.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
branches | list of flows | Yes | none | Parallel sub-flows; each gets a clone of the message. At least one required. |
Errors: the first branch error aborts the fork (cancelling the remaining branches) and propagates, labeled with the failing branch's name.
- type: fork
name: notify-all
branches:
- name: audit
process:
- type: log
settings: { message: '"audit " + correlationID' }
- name: metrics
process:
- type: publish-event
settings: { connector: bus, topic: orders }foreach
Sequentially iterates an array, running the body sub-flow once per element with the element bound to a loop variable. mode picks what the loop is for: running the body for its side effects (iterate, the default), or transforming the array into a new one (map).
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
items | expression (array) | Yes | none | CEL expression that must evaluate to an array; anything else errors. |
as | string | No | item | Variable name each element is bound to (vars.<as> inside the body). |
mode | enum: iterate | map | No | iterate | iterate threads the message through the body per element; map collects each element's resulting body into an array that replaces the message body. |
body | flow | Yes | none | Sub-flow run once per element, in order. |
iterate mode (default)
- Iteration is sequential and runs on the shared message, so body edits (body, variables) carry into the next iteration and past the loop.
- The loop variable is restored to its pre-loop state (or removed) after the loop.
- A body that drops the message stops the loop and drops it; a body error aborts the loop; a filter block inside the body that requests stop halts iteration and bubbles the stop to the root.
- type: foreach
name: each-order
items: "body.orders"
as: order
body:
process:
- type: log
settings:
message: '"order " + string(vars.order.id) + " amount=" + string(vars.order.amount)'map mode
Each element's body runs on its own scope of the message, and the body it produces becomes one element of a new array, which replaces the message body. This is how you reshape a payload without hand-appending to a variable.
- type: foreach
name: enrich-orders
items: "body.orders"
as: order
mode: map
body:
process:
- type: set-payload
settings:
value: '{"id": vars.order.id, "total": vars.order.amount * 1.2}'
# body is now the array of the objects each iteration produced- As many elements out as in. An iteration whose body drops the message contributes
nullat that position rather than shortening the array; downstream blocks decide what anullelement means. - Iterations are independent. Each runs on its own scope, so one element's variables cannot leak into the next and the loop variable never escapes, the opposite of
iteratemode. A scope copies the variables but shares the body, which keeps the loop linear in the size of the collection; see message copying semantics. - A body error aborts the loop. A filter block inside the body that requests stop halts iteration, writes back what was collected so far, and bubbles the stop to the root.
The mapped array replaces the message body. To keep the original body, wrap the foreach in an enrich block and use setVars to land the mapped array in a variable.
split
Splits one message into many. Each element continues through the rest of the flow as an invocation of its own, so elements are processed concurrently and each has its own error fate.
This is where a flow becomes asynchronous: the blocks after split run once per element, on borrowed workers, and the caller is answered by the buildResponse slot instead.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
mode | enum: expression | lines | delimiter | chunk | No | expression | How the message is carved into elements. |
items | expression (array) | Yes for expression mode | none | CEL expression that must evaluate to an array. |
delimiter | string | Yes for delimiter mode | none | Separator to cut the text body on. |
chunkSize | number | Yes for chunk mode | none | Characters per element. |
onError | enum: skip | abort | No | skip | What to do when an element cannot be dispatched. |
buildResponse | flow | No | none | Sub-flow that shapes the caller's response, run once after every element is dispatched. |
- type: split
name: lines
items: "body.lines"
buildResponse:
process:
- type: set-payload
settings:
value: '{"accepted": vars.groupSize}'
# every block below here runs once per lineexpressionmode takes the same inputforeachdoes.lines,delimiterandchunktokenize a text body (raw content or a string) one element at a time, so a large payload is never held in memory whole.- Each element is a full invocation with its own event ID, so one bad record fails alone and the flow's
error:chain recovers it independently. In aforeach, one failing element aborts the loop. - Backpressure is automatic: elements are scheduled onto the flow's pool while it has room and run on the splitting worker when it does not. Size it with the flow's
poolsetting. onErrorgoverns dispatch only. An element that fails once it is running always fails on its own.
Correlation variables
Every element carries three variables. They are the entire contract between split and aggregate, which is why the two pair with no configuration.
| Variable | Type | Description |
|---|---|---|
groupId | string | Identifies the group: the event ID of the message being split. |
groupIndex | number | The element's 0-based position. |
groupSize | number | How many messages the group holds, or -1 while not yet known. |
In expression mode the total is known up front, so every element carries it. A tokenizing mode cannot know it while emitting, so groupSize is -1 on every element except the last, which carries the real count. The buildResponse sub-flow always sees the real total, because it runs after every element has been dispatched.
split and aggregate must be top-level blocks in a flow's process or error chain: inside another block's slot there is no "rest of the flow" to continue into, so nesting one fails at startup. Both also require at least one block after them.
aggregate
Combines many messages into one. Messages are grouped by a CEL expression and held until the group completes (by size, by timeout, or by a predicate), then the completed group continues through the rest of the flow as a single message.
It is the only block that holds state between messages, so it can combine messages that arrive over time (batch 100 webhook events; every event in a five-second window), not just elements of one body.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
correlation | expression | No | vars.groupId | Which group a message belongs to. |
strategy | enum: append | expression | No | append | append collects bodies into an array; expression folds through expression. |
expression | expression | Yes for expression strategy | none | Fold whose result becomes the new accumulator. |
completionSize | expression (number) | No | vars.groupSize | How many messages the group expects; it completes when the count reaches this. |
completionTimeout | duration | No | none | How long a group may stay open (e.g. 5s). |
completionExpression | expression (bool) | No | none | Predicate checked after each message. |
timeoutFrom | enum: first | last | No | first | Fixed window from the group opening, or a sliding idle window. |
storeKey | string | No | the block's address | Identifies what the group state and leader election are namespaced by. |
maxGroups | number | No | 1000 | How many groups may be open at once. |
onOverflow | enum: fail | complete | drop | No | fail | At the cap: fail errors the message into the flow's error chain, complete releases the group that has been open the longest to make room, drop discards the message. |
buildResponse | flow | No | none | Sub-flow that shapes the caller's response for the message just absorbed. |
# Re-join a split: no configuration needed, the defaults read the split's variables.
- type: aggregate
name: rejoin
completionTimeout: 10s
# Or batch events that never came from a split.
- type: aggregate
correlation: "body.customerId"
completionSize: "100"
completionTimeout: "5s"All three completion conditions may be combined; the first to fire wins.
- Completion is a count, never a last-marker: elements finish out of order, so there is deliberately no
groupLastvariable. appendre-joins in the split's order. Elements are placed bygroupIndex, so the same input yields the same array every run. Messages with nogroupIndex(events off a source) are appended in arrival order, as is one whose position is beyond the group's size or beyond 100,000. A group that completes with an element missing leavesnullin the gap.- Only the timeout fires without a message, so a group whose other conditions may never be met needs a
completionTimeout.vars.aggregateReasontells the rest of the flow whether it received a whole group or a partial one.
The completed message
| Variable | Description |
|---|---|
aggregateKey | The correlation key the group folded under. |
aggregateCount | How many messages it holds. |
aggregateReason | size, timeout, expression, or overflow. |
The message inherits the correlationID of the group's first message. The three group* variables are cleared, so a second aggregation stage does not put every message in a batch of one.
vars.group: the group so far
Every expression the block evaluates sees the accumulated state under vars.group, so a predicate can ask about the group rather than only the message in hand.
| Field | Description |
|---|---|
key | The correlation key. |
count | Messages folded so far. |
size | The learned total, or -1. |
acc | The accumulator; under append, the array of bodies. |
firstAt / lastAt | RFC3339 timestamps. |
ageMs / idleMs | Milliseconds since the group opened / since the last message. |
completionExpression: "vars.group.count >= 5 && vars.group.idleMs > 2000"The fold expression sees vars.group before its message is folded (acc is the previous accumulator, null on a new group); completionSize and completionExpression see it after.
State and clustering
Group state lives in the runtime KV store under optimistic concurrency, so two replicas can never each hold half of a group. The timeout sweep is gated on leader election, so a group is reaped once per cluster. Both work the same way standalone, where the store is in-process and leadership is permanent.
storeKey defaults to the block's own address. Set it explicitly to keep in-flight groups alive across an edit that renames or moves the block; otherwise those groups are stranded until their deadline. Two blocks may not share one.
strategy: append rewrites the whole group on every message, so the stored bytes grow quadratically with group size. That is fine for batching a hundred events and wrong for re-joining a fifty-thousand-element split: use strategy: expression with a running accumulator for large groups.
enrich
Runs a body sub-flow on an isolated scope of the message, then merges back only what you name: setBody computes the new body and each setVars entry sets one variable. Everything else stays isolated.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
body | flow | Yes | none | Sub-flow run once on an isolated scope of the message. |
setBody | expression | No | none | Evaluated against the scope's result; its value becomes the message body. Omitted or empty: the incoming body is unchanged. |
setVars | map of string → expression | No | none | Each expression is evaluated against the scope's result and set as a variable on the message. |
Errors: a body error aborts the block; a body that drops the message drops it here too.
- type: enrich
name: derive-total
# setBody omitted: the incoming order body is preserved.
setVars:
total: body.total # pull just the total out of the scope's result
body:
process:
- type: set-payload
settings:
value: '{"total": body.qty * body.price, "note": "scope-only body"}'validate
A filter block: asserts boolean CEL rules against the message. If all hold, the message passes through unchanged. If any fail, the block rejects: it shapes a terminal response and requests the flow stop, so the rest of the chain never runs.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
rules | list | Yes | none | Assertions; each entry is {expr: <expression (bool)>, message: <text>}. All must hold. |
rejectStatus | int | No | 422 | HTTP status of the built-in rejection response (written to vars.httpStatus). Not used when onReject is set. |
onReject | flow | No | none | Sub-flow run on rejection, before the flow stops. It shapes the response itself (e.g. set-payload + set httpStatus) and can read vars.validationErrors. Omitted or empty: the built-in response is used. |
On rejection:
- The failing rules'
messagetexts are collected intovars.validationErrors; all rules are evaluated, so every failure is reported. - Without
onReject, the built-in response body is{"error": "validation_failed", "messages": [...]}andvars.httpStatusis set torejectStatus. Over an HTTP source that becomes the response status. - With
onReject, the sub-flow's output is the terminal response; one that drops the message drops it here too.
- type: validate
name: check-order
rejectStatus: 422
rules:
- expr: 'has(body.id)'
message: "order id is required"
- expr: 'body.amount > 0'
message: "amount must be positive"
# Only reached when every rule held.
- type: set-payload
settings:
value: '{"orderId": body.id, "status": "accepted"}'handle-errors
Inline error recovery: runs the process chain and, on error, exposes the failure as vars.error and runs the error chain. Both are bare block lists. When the error chain succeeds, its output is the block's result and the flow continues (recovery).
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
process | block list | Yes | none | The happy-path chain. |
error | block list | Yes | none | Runs when the process chain errors; reads vars.error. |
vars.error is a map: {message: <error text>, flow: <handle-errors block name>, block: <failing block label>}. If the error chain itself errors, the block fails with that error. A root flow's own error: chain provides the same recovery at flow level; see Error handling.
- type: handle-errors
name: charge
process:
- type: rest
name: call-charge
settings:
connector: payments
method: POST
path: /charges
body: '{"amount": body.amount}'
error:
- type: set-payload
settings:
value: '{"status": "degraded", "reason": vars.error.message}'cache-scope
Memoizes the body its wrapped flow produces in the runtime store, keyed by an evaluated expression. On a fresh hit it restores the cached body and skips the body flow; on a miss or an expired entry it runs the flow and stores its body.
| Key | Type | Required | Default | Description |
|---|---|---|---|---|
key | expression (string) | Yes | none | Cache-key expression, evaluated per message. An invalidate-cache block with the same key expression evicts this entry. |
ttl | duration string | No | 60s | How long an entry stays fresh (e.g. 5m). "0" never expires. Negative values are rejected at startup. |
body | flow | Yes | none | Sub-flow whose result body is cached and replayed on a hit. |
Notes:
- Only the message body is cached; variables the wrapped flow sets are not restored on a hit.
- Storing is best-effort: a version conflict (another worker cached first) or an absent store leaves the result correct, just uncached. A body error aborts the block; a body that drops the message drops it, and nothing is cached.
- Entries live in the same per-deployment store as object storage; see Runtime services.
- type: cache-scope
name: cached-report
key: '"report"'
ttl: 1m
body:
process:
- type: set-payload
settings:
value: '{"report": "generated", "event": eventID}'