StageFlow
A framework for describing and running JSON-defined pipelines: a graph of
nodes, user-defined stages, an immutable data frame, CEL expressions, retry
and block-scoped try/except, parallel branches, nested pipelines.
Requires Python 3.11+ (the common-expression-language CEL binding does).
Installation
python -m venv .venv
source .venv/bin/activate
pip install -e .
Quick start
Register a stage:
from stageflow import BaseStage, register_stage
@register_stage("HelloStage")
class HelloStage(BaseStage):
"""
description: "Custom stage example"
icon: "👋"
arguments:
name: string
outputs:
greeting: string
"""
async def run(self):
name = self.get_arguments().get("name", "world")
self.set_outputs({"greeting": f"Hello, {name}!"})
Describe a pipeline:
pipeline_dict = {
"nodes": [
{
"id": "start",
"type": "entry",
"variables": {"user_name": "Alice"},
"next": "hello",
},
{
"id": "hello",
"type": "stage",
"stage": "HelloStage",
"arguments": {"vars": {"name": "user_name"}},
"outputs": {"greeting": "greeting"},
"next": "finish",
},
{
"id": "finish",
"type": "terminal",
"result": {"status": "ok"},
"artifacts": ["greeting"],
},
],
}
Run it:
import asyncio
from stageflow import Pipeline, Session
async def main():
session = Session(id="demo", pipeline=Pipeline.from_dict(pipeline_dict))
result = await session.run()
print(result.result) # {'status': 'ok'}
print(result.artifacts) # {'greeting': 'Hello, Alice!'}
asyncio.run(main())
Values declared by the entry node are defaults: anything supplied from the
outside overrides them, so the same pipeline can be parameterised without
editing its JSON.
from stageflow import Context
session = Session(id="demo", pipeline=pipeline, context=Context({"user_name": "Bob"}))
Node types
type |
Purpose | Key fields |
|---|---|---|
entry |
graph start and the variables it begins with | variables, next |
stage |
run a registered stage | stage, arguments, outputs, consume, next |
condition |
binary branch on a CEL expression | condition, then, else |
switch |
n-way branch, first matching case wins | cases: [{when, next}], default |
parallel |
concurrent branches with independent frames | branches, cancel_on_error, next |
try |
block-scoped error handling over a graph region | body, except, next |
subpipeline |
nested pipeline with a fresh frame | subpipeline_id, inputs, artifact_outputs, result_output, next |
terminal |
end of execution | result, artifacts |
Every node additionally accepts retry, consume (drop names from the frame
after the step) and expose (copy or rename a variable without a stage).
Data model
Pipeline data is a single frame, vars, carried along the execution path.
The frame is immutable: every write produces a new one, so parallel branches
diverge independently and the same object can safely be handed to several
concurrent consumers. The only visibility boundary is subpipeline: a child
starts with a fresh frame and receives data through inputs.
A node's arguments are split into buckets: vars holds references to frame
variables, const holds literals. Buckets can be mixed in one node; on a name
collision the variable reference wins over the literal. A stage node has no
separate settings field — a literal setting is a const argument.
{
"id": "auth",
"type": "stage",
"stage": "AuthStage",
"arguments": { "vars": { "creds": "creds" }, "const": { "timeout_s": 30 } },
"outputs": { "token": "token" },
"consume": ["creds"],
"next": "fetch"
}
In outputs the key is the name of a field in the stage result and the value
is the variable to store it in. A field the stage does not return is rejected
by Pipeline.validate() against the stage specification. If a stage declares
no outputs at all, the contract is considered undeclared and the check is
skipped.
A node's outputs are applied as a simultaneous assignment: values are computed against the frame as it was on entry and only then written, so key order does not affect the result.
Expressions
The expression language is CEL. An
expression is allowed in condition of a condition node, in when of a
switch node, and in any value whose key carries the .$ suffix in
arguments, outputs and variables.
"outputs": {
"value": "n",
"attempts.$": "0",
"greeting.$": "'hello ' + string(vars.user_name)"
}
Frame variables are addressed through the vars namespace (vars.n). A name
that is not an ASCII identifier is addressed by index: vars['итог'].
A key with the .$ suffix in outputs names a variable rather than a stage
field, so a single node can introduce any number of pipeline variables.
Backend: common-expression-language (native), falling back to cel-python.
The entry node
The graph starts at a node, not at a field in the JSON header, and that node declares the pipeline's initial variables.
{
"id": "start",
"type": "entry",
"variables": {
"n": 5,
"items": [1, 2, 3],
"total.$": "vars.n * 2"
},
"next": "check"
}
- A name already present in the frame (session seed, a parent's
inputs) is neither overwritten nor evaluated. - Variables are bound in dependency order: an expression may reference a
sibling variable of the same node, key order in the JSON is irrelevant, and
a cycle is a
validate()error. - There is exactly one
entrynode per graph and jumping back into it is forbidden. The pipeline-levelentryfield is optional when such a node exists, and must point at it when given. - Value types are checked statically, during validation.
Parallel branches
{
"id": "fan_out",
"type": "parallel",
"branches": [{ "id": "hash", "entry": "hash_step" },
{ "id": "thumb", "entry": "thumb_step" }],
"next": "merge"
}
Only names that were not in the frame before parallel leave a branch; a
write to a name that existed on entry stays branch-local. Such names are
listed in the parallel_completed event:
{"merged": ["fresh"], "dropped": ["left.n"]}. Two branches writing the same
name raise BranchError naming both.
cancel_on_error (default true) decides the fate of sibling branches when
one fails: true cancels them immediately (the parallel_cancelled event),
false lets them finish. Either way the node fails with the error of the
first branch that failed.
Errors: retry per node, try/except per region
Retrying is a property of an operation, so retry is a node field:
{ "retry": [{ "error_equals": ["TimeoutError"], "max_attempts": 3, "backoff_rate": 2.0 }] }
max_attempts counts every run of the node including the first: 3 means one
run and two retries. Each policy in the list keeps its own counter.
Error handling is block-scoped: a try node covers a region of the graph, and
an error raised by any node inside it goes to the matching except.
{
"id": "safe_fetch",
"type": "try",
"body": "fetch",
"except": [
{ "error_equals": ["TimeoutError"], "next": "on_timeout", "result_var": "error" },
{ "error_equals": ["*"], "next": "on_any" }
],
"next": "after"
}
- The region is every node reachable from
bodybut not reachable fromnext; it is derived from the graph rather than listed by hand. - The failing node's
retrypolicies are exhausted first, then the error propagates to the nearest enclosingtry; an error no handler matches keeps propagating outwards. - Nested
trynodes work as expected: the inner one simply lies inside the outer one's region. - A handler sees the frame as the last successfully completed node of the body left it.
result_varputs the error object into the frame with the fieldstype,full_type,messageandnode.error_equalsaccepts a bare exception class name, a fully qualified path, or*.
Variable typing
Typing is gradual: an undeclared variable is not checked, and the type sections are optional.
{
"types": {
"UserId": "int",
"User": {
"id": "UserId",
"name": "string",
"email?": "string",
"settings": { "theme": "string" }
},
"Point": { "fields": { "x": "int", "y": "int" }, "strict": true },
"Tree": { "value": "int", "children": "list<Tree>" }
},
"variables": { "user": "User", "attempts": "int", "tags": "list<string>" }
}
The type language: primitives string, int, float, number (int|float),
bool, any, null; containers list<T> and map<T> (string keys); unions
T|U; the shorthand T? for T|null; names from the types section.
Structures support optional fields (a ? suffix on the name), nested
anonymous structures, recursion, and a strict mode (strict forbids extra
fields).
Checks come in two layers:
- statically, during graph validation: a variable's declared type is matched
against the type hints in the stage specification, and
exposerequires the source and destination to be compatible; a mismatch is aPipeline.validate()error raised before the run starts; - dynamically, during execution: every write to a declared variable (
entry,outputs,expose,except.result_var, subpipeline artifacts) and the session's input context are checked against the full structure of the value; a mismatch raisesTypeCheckErrorcarrying the node and the path to it.
A subpipeline inherits its parent's named types and may declare its own; variable types are its own.
Session control
session.stop(); session.pause(); session.resume()
await session.input("command", {"name": "skip"})
User input: a stage declares allowed_inputs and awaits
await self.wait_input("user_input", timeout=...); input is delivered from
outside with await session.input("user_input", {...}). The payload is
validated against the payload_schema from the declaration.
Snapshots: session.snapshot() returns a dict, Session.from_snapshot(snap)
restores the session, and run() resumes from the saved node.
Step debugging
Session accepts a debugger that is given control before and after every
node. The in-core implementation is StepDebugger.
from stageflow import Pipeline, Session, StepDebugger
debugger = StepDebugger(mode="step", delay=0.0, on_event=print)
session = Session("s1", Pipeline.from_dict(data), debugger=debugger)
task = asyncio.create_task(session.run()) # stops before the first node
debugger.step() # let one node run, then stop again
debugger.set_vars({"n": 42}) # applied before the next node
debugger.set_delay(0.5) # run on its own, pausing between nodes
debugger.resume() # continue without stopping
result = await task
Available from outside: the stop point (debugger.node), the frame at that
point (debugger.vars), and the on_event stream: node_enter, node_exit,
paused, var_set, var_rejected. Frame edits are checked against declared
types — a mismatch is rejected with an event rather than crashing the session.
The debugger also applies inside a try body, inside parallel branches and
inside a subpipeline: every node goes through Session.execute_node, and a
child session inherits the debugger. Commands are thread-safe.
Built-in stages
| Category | Stages |
|---|---|
| vars | SetValueStage, CopyValueStage, IncrementStage, MergeDictStage |
| lists | AppendListStage, ExtendListStage, FilterListStage, UniqueListStage, PopListStage |
| dicts | PickKeysStage, DropKeysStage |
| strings | ConcatStage, TemplateStage |
| logic | AssertStage, FailStage, LogStage, SleepStage |
All of them return new values and never mutate their input.
Stage specification
A stage is specified by YAML in its docstring: description, arguments,
outputs, plus visual hints for the editor.
description: "Increment numeric value by delta"
icon: "+" # glyph, SVG link, data URI or inline <svg> markup
icon_mono: false # recolour the SVG to the node colour (monochrome sets)
color: "#ff8800" # card accent (defaults to the category colour)
icon accepts four forms:
| Value | What gets drawn |
|---|---|
"+", "👋" |
the glyph or emoji itself |
"/icons/globe.svg", "https://…/x.svg" |
an SVG by link |
"data:image/svg+xml;utf8,…" |
a data URI |
"<svg …>…</svg>" |
markup straight from the docstring |
icon_mono: true draws the SVG as a mask in the node colour, which suits
monochrome sets (lucide, feather, tabler) that paint via currentColor.
Without icon the editor draws a monogram of the stage name
(IncrementStage → IS); without color it picks a deterministic colour for
the category.
A stage may also declare allowed_events and allowed_inputs (EventSpec /
InputSpec with a payload_schema), a category and a timeout. All of it
ends up in get_specs().
Schema and stage specifications
The pipeline JSON Schema and the specifications of registered stages:
from stageflow.docs import generate_pipeline_schema, generate_stages_json, load_pipeline_schema
from stageflow import get_stages
schema = generate_pipeline_schema(get_stages()) # schema with the stage-name enum
stages = generate_stages_json(get_stages()) # stage specs for an editor
load_pipeline_schema() returns the schema without the injected enum; it is
the one Pipeline.validate() uses. These two functions supply everything an
external tool needs: an editor, a CI validator, a documentation generator.
Tests
python -m unittest discover -s tests
For declarative pipeline testing there is stageflow.testing:
from stageflow.testing import PipelineTestSpec, run_pipeline_test
Releasing
Pushing a tag builds the package and publishes it to PyPI
(.github/workflows/publish.yml):
# bump project.version in pyproject.toml first — the workflow refuses
# a tag that does not match it
git tag -a 1.0.0 -m "StageFlow 1.0.0"
git push origin 1.0.0
The workflow runs the test suite, builds an sdist and a wheel, validates both
with twine check, and uploads them through PyPI Trusted Publishing, so no
API token is stored in the repository. The upload step runs in the pypi
environment, which can be configured to require a manual approval.
License
MIT — see LICENSE.
Package layout
stageflow/
core/ the engine: pipeline, session, nodes/, context, cel, stage, typesys, inputs, debug
builtins/ built-in stages
docs/ pipeline JSON Schema and stage specifications
exceptions.py exception hierarchy
testing.py pipeline testing helper
tests/ unit tests
Metadata
Release files for stageflow-framework 0.9.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| stageflow_framework-0.9.0.tar.gz | 65.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| stageflow_framework-0.9.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 120.1 kB
Release files / stageflow_framework-0.9.0.tar.gz
| Download URL | stageflow_framework-0.9.0.tar.gz |
|---|---|
| Size | 65.7 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
b509c7f6e7bd8bbd16ad1b5bd77ad6309c5b2886d54d898531f200d36eeec0b2
|
|
BLAKE2b-256 checksum How to use checksums |
5583c91cf816545f2dd0cbea9c10728958b39f0a729fb21ad0e66ec90f8c758b
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 19, 2026.
Transparency logRelease files / stageflow_framework-0.9.0-py3-none-any.whl
| Download URL | stageflow_framework-0.9.0-py3-none-any.whl |
|---|---|
| Size | 54.4 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
c31c2693ed83e7a7c0d1d5e4853e52d01bff011fe8310b5d8f940a751ce8fd06
|
|
BLAKE2b-256 checksum How to use checksums |
61ae21666cc40b2be1391506c62e9a394e62678b10ab4e1df72d049aae3f7777
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 19, 2026.
Transparency log