Flows¶
A flow is a Python function decorated with @flow from prefect_compat (not prefect). When you invoke it, IronFlow creates a flow run in the control plane: the Rust engine records state transitions and append-only history for that run and its task runs.
Basics¶
- Import:
from prefect_compat import flow(see Prefect → IronFlow). - A flow coordinates task runs by calling
task.submit(...),task.map(...), andwait(...)on futures—see Tasks. - You typically register a control plane (for example
InMemoryControlPlane) withset_control_planebefore executing the flow; see Quick start (demo flow).
Final state (wait_all default)¶
By default IronFlow uses @flow(final_state="wait_all"): after the flow body returns, it drains in-process submits, waits for non-detached children (including deployment-backed subflows), then resolves the flow terminal state in Rust from contributing task-run rows (CANCELLED > FAILED > all COMPLETED). Unobserved failed concurrent submits fail the flow (FlowChildrenFailed).
Escape hatches:
submit(..., detach=True)— exclude one task/subflow from the wait set (true fire-and-forget / safe-to-fail).@flow(final_state="explicit")— body return/exception remains authoritative (closer to Prefect’s return-value model).
Design notes: flow-run final state plan. Compatibility: Compatibility matrix.
Subflows (nesting flows)¶
IronFlow supports two nesting mechanisms:
- Inline (blocking) — call another
@flowas a normal Python function inside the parent. Same process; parent waits for the return value; child run is linked (execution_mode=inline). - Deployment-backed (subflow as task) —
deployment_ref("deployment-name").submit(...).result()enqueues work on a deployment’s work pool. Returns aSubflowFuturethat works withwait_for/wait. Underwait_all, omitting.result()still waits unlessdetach=True.
Step-by-step examples, UI notes, and choosing between the two: How to compose flows with subflows. Supported subset and limits: Compatibility matrix.
Transition hooks (IronFlow extension)¶
Flows support transition_hooks: a sequence of TransitionHookSpec values built with on_transition(fn, from_state=..., to_state=...). Use None for from_state or to_state to match any state on that side.
Hooks run after a successful control-plane transition, in process, without holding the control-plane lock. They are not the same API names as Prefect’s on_running / on_failure hooks; map your logic to explicit edges (for example PENDING → RUNNING). For full semantics (including the batched start path and error handling), see Compatibility matrix.
Relevant exports: TransitionHookSpec, on_transition, TransitionContext from prefect_compat.
Runtime context and logging¶
Inside an active flow (or task) body:
from prefect_compat import get_run_context, get_run_logger
@flow
def pipeline(n: int) -> int:
ctx = get_run_context() # flow_run_id, flow_name, parameters, …
log = get_run_logger()
log.info("starting with n=%s run=%s", n, ctx.flow_run_id)
return n
get_run_logger() messages appear under GET /api/flow-runs/{id}/logs and the UI Logs tab. Outside a run, the logger writes to stderr and does not persist. See Tasks for task-scoped association.
Operator pause / cancel (subset)¶
- Cancel —
POST /api/flow-runs/{id}/cancel(always terminate semantics). - Pause —
POST /api/flow-runs/{id}/pausewith required JSON{"mode": "drain"}or{"mode": "terminate"}(no ambiguous default). - Resume —
POST /api/flow-runs/{id}/resumefor operator pauses only (gate waits are different).
Drain lets in-flight tasks finish then holds PAUSED; further submit in the same in-process body raises FlowRunSchedulingHeld. Terminate / cancel cancel RUNNING rows (late COMPLETED fenced) and, under ProcessPoolTaskRunner, SIGTERM→SIGKILL registered child processes. Thread-pool bodies remain cooperative-only. After terminate pause, in-process resume prepares P1 lineage for the next @flow() invoke (prior attempt is terminalized); deployment-backed resume uses retry-with-resume_from.
Step-by-step: How to cancel, pause, and resume. Design: lifecycle plan.