Skip to main content
Run stages in order through the stored graph executor, with each dispatch recorded before the stage runs, so a stopped pipeline carries on where it left off. Use 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)
The form works only on frozen data: it folds recorded events into state and decides what runs next. A dag definition needs:
  • globalConcurrency, the cap on running nodes, and keyedConcurrency for nodes that share a key.
  • stopOnError, whether a required failure stops new work.
  • Nodes with a kind, required or optional, 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)
The review function reads the draft’s result, and the publish function reads the review’s. Each dispatch receives a stable position such as 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
Three dispatches in order, the final output keyed by node id, and the plan declares exactly three dispatches with a maximum concurrency of one.

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 returns waiting 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: true and 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.
A node with a 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

  • globalConcurrency caps all running nodes. A node with a non-null key also uses the matching cap in keyedConcurrency.
  • The final output maps every completed node’s id to its returned payload.
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.