Skip to main content

Statements Reference

Statements are the building blocks of Moco workflows. Every body in a workflowspec is a statement. Statements fall into two categories: primitives (leaf nodes that do work) and composites (containers that orchestrate other statements).

For full workflowspec context, see the Workflowspec Reference. For the activities an activity statement can invoke, see the Activity Catalog.

Common Parameters

All statements support these optional parameters:

ParameterTypeDescription
namestringUnique identifier for the statement
descriptionstringHuman-readable description
conditionexpression or listSkip this statement if the expression is falsy
output_namestringVariable to store the statement result
output_datalistData transformations to apply after execution

Primitive Statements

transform

Evaluates expressions and assigns variables. The primary way to compute or reshape data.

- transform:
input_data:
- temp: "{{ price * 1.08 }}"
output_data:
- total: "{{ temp }}"
- message: "Total is {{ total }}"
ParameterDescription
input_dataVariable assignments evaluated before output_data
output_dataMain transformation assignments

abort

Terminates or breaks execution with different behaviors.

- abort:
condition: "{{ price < 0 }}"
type: raise
message: "Invalid price: {{ price }}"
TypeBehavior
abortAbort the entire workflow
terminateGracefully terminate the workflow
breakBreak out of the current sequence or parallel block
break_iterationBreak out of the current iteration loop
raiseRaise an error and fail the workflow

activity

Executes a registered activity (HTTP call, database query, custom function, etc.).

- activity:
type: http.request
input_data:
method: POST
url: https://api.example.com/orders
body:
order_id: "{{ order_id }}"
output_name: api_response
retry_policy:
timeout_sec: 30
max_attempts: 3
ParameterDescription
typeActivity type identifier (required)
versionActivity version (default: 1.0.0)
config_dataStatic configuration, evaluated once at workflow start
input_dataDynamic input, evaluated each time the activity runs
output_nameVariable to store the activity result
retry_policyNested timeout/retry config (Temporal only): timeout_sec (per-attempt execution timeout), schedule_to_close_timeout_sec, heartbeat_timeout_sec, heartbeat_interval_sec (heartbeat cadence; heartbeating is enabled only when both heartbeat_timeout_sec and heartbeat_interval_sec are set), max_attempts (total attempts = initial + retries), initial_interval_sec, backoff_coefficient, maximum_interval_sec, non_retryable_error_types
execute_locallyForce local execution, bypassing Temporal. Overrides the activity's own default; omit it to keep that default
enable_cacheEnable result caching
cache_policyCache configuration (TTL, key)

Every available activity type, with its input and output contract, is in the Activity Catalog.

Some activities already default to local execution, so execute_locally is rarely needed. Setting it to false on a selenium.* or playwright.* activity breaks browser sessions — see Activities that are already local by default.


workflow

Executes a child workflow by reference or inline definition.

- workflow:
wfspec:
name: process-order
version: 1.0.0
child_mode: sync
input_data:
order_id: "{{ order_id }}"
output_name: order_result
Child ModeBehavior
inlineRuns in the parent's context, shares variables (default)
syncRuns independently; parent waits for the result
asyncRuns independently; parent waits only for start, gets workflow_id
detachedRuns completely independently; parent doesn't wait

To define a workflow inline instead of by name:

- workflow:
wfspec:
content:
wfspec_name: inline-helper
wfspec_version: 1.0.0
input_data:
x:
output_name: result
body:
transform:
output_data:
- result: "{{ x * 2 }}"
child_mode: inline
input_data:
x: 21
output_name: doubled

call

Invokes a function defined in the enclosing wfspec's functions list. A call is a compact shorthand for a workflow statement running in inline mode: the function executes in a fresh context and its return value is mapped back via output_name / output_data.

- call:
function: add # matches a `function` in the wfspec's `functions`
input_data:
a: 1
b: 2
output_name: result
ParameterDescription
functionName of the function to call (required)
input_dataArguments passed to the function; supports expressions
output_nameContext variable to store the function's return value
output_dataTransform the return value (_raw_output holds the raw return)
conditionPre-condition; the call is skipped when it evaluates to false

Functions may call sibling functions defined in the same wfspec. Recursion (a function calling itself, directly or indirectly) is not supported — it is rejected by the cyclic call-stack guard.


wait_for

Waits for an event matching filter criteria, or until a timeout.

- wait_for:
event:
topic: order_events
match_expression: >
{{ event.data.get('order_id') == order_id and
event.data.get('status') == 'completed' }}
timeout_sec: 60
output_name: completion_event
ParameterDescription
event.topicEvent topic to subscribe to
event.event_typeOnly events with this exact type match
event.match_expressionPython expression to filter incoming events (event variable is the event object)
timeout_secMaximum wait time in seconds (required)
output_nameVariable to store the received event

To collect the result of an activity started with async_mode: true, pass the token that activity returned as event.event_type — see Async activity results.


emit_event

Emits an event to the event bus.

- emit_event:
input_data:
topic: notification_events
data:
type: order_created
order_id: "{{ order_id }}"
target_workflow_id: "{{ parent_id }}"
metadata:
priority: high
ParameterDescription
topicEvent topic (required)
dataEvent payload (required)
target_workflow_idRoute the event to a specific workflow (optional)
metadataAdditional event metadata (optional)

continue_as_new_checkpoint

Checks whether the runtime suggests restarting the workflow (e.g., Temporal event history nearing its size limit). If suggested, serializes workflow state and restarts execution from the beginning with the preserved state.

This is a no-op in the in-memory runtime. In Temporal, it triggers a continue-as-new when the SDK signals that history is getting large.

- continue_as_new_checkpoint:
name: checkpoint-after-processing
serialize_data_context: true
ParameterTypeDefaultDescription
namestringnullOptional identifier for logging
serialize_data_contextbooleantrueWhether to include data context variables in the serialized state
conditionexpressionnullSkip this statement if the expression is falsy
enforcebooleanfalseRestart regardless of runtime suggestion. Testing only

After the restart the workflow body runs from the beginning, so the spec must skip work it has already done. If the body is a single state_machine, prefer a state checkpoint_policy instead — the runtime resumes in the right state by itself.


checkpoint_policy

Declared on a state, not as a statement. The runtime checkpoints while the machine sits in that state and resumes it there afterwards, so the spec does not have to describe recovery.

states:
- name: idle
checkpoint_policy:
timeout_sec: 300
event_count: 100
ParameterTypeDefaultDescription
timeout_secnumber | expressionnullCheckpoint after this long in the state
event_countintegernullCheckpoint after this many events handled in the state
serialize_data_contextbooleantrueWhether to include data context variables
auto_resumebooleantrueResume in this state, skipping on_enter. False restarts the machine at initial_state
enforcebooleanfalseCheckpoint regardless of runtime suggestion. Testing only

At least one of timeout_sec / event_count is required; if both are set, whichever trips first wins.

An event counts toward event_count when a transition matched it for the current state and the state is unchanged afterwards — internal transitions, self-transitions, and events whose transition condition was false. Unmatched events do not count, and moving to a different state resets the count.

Requires the wfspec body to be a single state_machine statement; a policy anywhere else is rejected at parse time. On resume the state's on_enter is skipped (it already ran before the checkpoint) and its timers are re-armed from their full duration.

Opting out of auto-resume

auto_resume: false keeps the checkpoint but skips the jump: the machine restarts at initial_state with on_enter running normally.

Use it when an earlier state has a side effect the new execution needs again — typically an init state that starts a long-running async_mode activity feeding the machine's event_source_topic. Those activity handles do not survive continue-as-new, so resuming straight into the state that consumed their events leaves the machine with no producer.

The record is still written, so the state the checkpoint fired in stays readable as __sys_info__.state_machine.checkpointed_from_state. An init state can use it to re-run its side effect and then route straight back, skipping one-time setup:

states:
- name: init
on_enter:
activity: # the subscriber, restarted every execution
name: subscriber
type: rabbit.receive
async_mode: true
async_event_topic: monitor_events
output_name: subscriber_token
input_data: { queue: prices }

- name: monitoring
checkpoint_policy:
event_count: 100
auto_resume: false

transitions:
# Listed first: the first satisfied automatic transition wins.
- from_state: init
to_state: monitoring
trigger:
condition: '{{ __sys_info__.get("state_machine", {}).get("checkpointed_from_state") == "monitoring" }}'

- from_state: init
to_state: warmup
trigger: null

resumed_from_state is set only when the machine actually jumped, so under auto_resume: false it is null while checkpointed_from_state names the state. Both are null on a first run.

Two more keys sit alongside them, available in any state on_enter/on_exit and any transition action/condition while the machine runs:

KeyDescription
nameThe machine's name, or unnamed_state_machine if the spec omits it
current_stateThe state the machine is in right now

current_state is updated before a state's on_enter runs, so a state always sees itself rather than the one it just left. During a transition's action it is still the source state. The whole state_machine key is removed once the machine completes, and a nested machine shadows it and restores the outer values on exit.


Composite Statements

sequence

Executes statements one after another.

sequence:
elements:
- transform:
output_data:
- status: "validating"
- activity:
type: http.request
input_data:
url: https://api.example.com/validate
output_name: validation
- abort:
condition: "{{ not validation.valid }}"
type: raise
message: "Validation failed"

parallel

Executes statements concurrently with configurable join semantics.

- parallel:
join_type: and
elements:
- activity:
name: fetch-user
type: http.request
input_data:
url: https://api.example.com/users/{{ user_id }}
output_name: user_data
- activity:
name: fetch-orders
type: http.request
input_data:
url: https://api.example.com/orders?user={{ user_id }}
output_name: order_data
Join TypeBehavior
andAll branches must succeed
orAt least one branch must succeed

iteration

Loops over a collection, either sequentially or in parallel.

- iteration:
iter_type: sequence
input_data: "{{ items }}"
body:
activity:
type: process-item
input_data:
item: "{{ iter_item }}"
ParameterDescription
iter_typesequence (sequential) or parallel (concurrent)
input_dataCollection to iterate over (list, dict.items(), range(), etc.)
bodyStatement executed for each item
join_typeFor parallel iteration: and or or

Special variables available inside the body:

  • iter_item — the current item
  • iter_items — all items in the collection

Use abort with type: break_iteration to exit the loop early.


state_machine

Event-driven finite state machine. See State Machines Reference for full documentation.

- state_machine:
name: order-fsm
initial_state: pending
timeout_sec: 300
states:
- name: pending
- name: processing
- name: completed
is_terminal: true
transitions:
- from_state: pending
to_state: processing
trigger:
event_type: start

trigger.event_type accepts an expression, resolved against the workflow context on every incoming event — see expression triggers.


Next Steps