dagGraphType, the dependency-graph form that ships, when the
run must be stored and resumed by the executor. For a pipeline you start
and watch yourself, pipeline() from the runtime is shorter. For a shape of
your own, read Outside graph types.
A pipeline is one use of the dag form: each stage is a node, and each edge
says which stage must finish before another can start:
examples/pipeline.ts (excerpt)
globalConcurrency, the cap on running nodes, andkeyedConcurrencyfor nodes that share akey.stopOnError, whether a required failure stops new work.- Nodes with a
kind,requiredoroptional, and edges from each stage to the one it must finish before.
Store and run
persistRunDefinition stores the definition with its resolved plan, and
createGraphExecutor gets one data function per node:
examples/pipeline.ts (excerpt)
dag/review/1, and its input has the completed results of its direct
predecessors, keyed by node id. Run the file with npx tsx pipeline.ts; it
runs with no model behind it:
Output
Restart and resume
If the executor process stops after it records a stage dispatch and before that stage completes, a new executor over the same store returnswaiting
with that stage’s position. The executor doesn’t lock the run, so wait
until the old process has stopped, then call resume with that position.
Resume writes no second dispatch event; earlier stages that completed stay
complete, and the pipeline carries on from that stage. With
globalConcurrency at 1, at most one stage is in flight and
waiting.positions has one entry.
Choose node behaviour
- Required nodes block their dependants after failure and fail the pipeline after finalizers run.
- Optional nodes can fail without blocking their dependants. A failed optional node is left out of the final output because it has no completed payload.
- Expected skips use a completed result with
skipped: trueand don’t block dependants. - Finalizer nodes have no edges and run after the other nodes settle. A failed finalizer fails the graph even when every worker completed.
- Paused nodes stop new dispatches and keep the recorded pause reason.
lane object uses the engine route stored in the plan; a
node without one uses its data function and is data-only. Nodes can share a
lane id when every copy of its declaration is identical, and a conflicting
declaration is rejected. stopOnError: true stops new work after a
required failure; set it to false when independent branches and declared
retries must continue.
Limits
globalConcurrencycaps all running nodes. A node with a non-nullkeyalso uses the matching cap inkeyedConcurrency.- The final output maps every completed node’s id to its returned payload.
Full file
Full file
examples/pipeline.ts
Next steps
- Graph executor: how a dispatch runs, fallback, and resume after a process stops.
- Plan admission: the resolved plan a stored pipeline runs under.
- Running: the runtime’s
pipeline()for a run you start and watch.