> ## Documentation Index
> Fetch the complete documentation index at: https://obversa.ai/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# Built-in Pipeline

> Run ordered stages through the stored graph executor with the dag form that ships.

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](/docs/graphs/contract).

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:

```ts examples/pipeline.ts (excerpt) {5,7,15-16} theme={null}
const definition: DagDefinition = defineGraphDefinition({
  id: 'release-pipeline',
  definitionVersion: 1,
  data: {
    globalConcurrency: 1,
    keyedConcurrency: {},
    stopOnError: true,
    retryCapPerNode: 0,
  },
  nodes: [
    { id: 'draft', data: { kind: 'required', key: null } },
    { id: 'review', data: { kind: 'required', key: null } },
    { id: 'publish', data: { kind: 'required', key: null } },
  ],
  edges: [
    { id: 'draft-to-review', source: 'draft', target: 'review', data: {} },
    { id: 'review-to-publish', source: 'review', target: 'publish', data: {} },
  ],
});
```

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:

```ts examples/pipeline.ts (excerpt) {1,12,14,16-19} theme={null}
  await persistRunDefinition(storage, {
    runId: 'release-pipeline-run',
    eventId: 'release-pipeline-started',
    timestamp: '2026-01-01T00:00:00.000Z',
    graphDefinition: graph.definition,
    resolvedPlan: plan,
    resolvedInputs: {},
    workspaceBinding: null,
    hostBinding: null,
  });

  const executor = await createGraphExecutor({
    runId: 'release-pipeline-run',
    graph,
    storage,
    nodes: {
      draft: nodeBinding(temporaryRoot, async () => {
        order.push('draft');
        return { article: 'ready' };
      }),
      review: nodeBinding(temporaryRoot, async ({ input }) => {
        order.push('review');
        const draft = (input as {
          readonly results: { readonly draft: { readonly article: string } };
        }).results.draft;
        return { approved: draft.article === 'ready' };
      }),
```

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:

```json Output theme={null}
{
  "executor": "complete",
  "output": {
    "nodes": {
      "draft": {
        "article": "ready"
      },
      "review": {
        "approved": true
      },
      "publish": {
        "published": true
      }
    }
  },
  "dispatches": 3,
  "order": [
    "draft",
    "review",
    "publish"
  ],
  "planDigest": "sha256:ebbf122994643028c5222a0afaf8d7a99a257716e00f0e255b15f315e4b862ed",
  "bounds": {
    "dispatches": {
      "min": {
        "kind": "known",
        "value": 3
      },
      "max": {
        "kind": "known",
        "value": 3
      }
    },
    "maxConcurrency": {
      "kind": "known",
      "value": 1
    },
    "maxFanOut": {
      "kind": "known",
      "value": 1
    }
  }
}
```

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.

<Accordion title="Full file">
  ```ts examples/pipeline.ts theme={null}
  import { mkdtemp, realpath, rm } from 'node:fs/promises';
  import { tmpdir } from 'node:os';
  import { join } from 'node:path';

  import {
    compileGraph,
    createGraphExecutor,
    dagGraphType,
    persistRunDefinition,
    resolveGraphPlan,
    type DagDefinition,
    type GraphNodeBinding,
    type PlanResolution,
    type RunStoragePolicy,
  } from '@obversa/runtime';
  import { defineGraphDefinition } from '@obversa/runtime/testing';
  import { createLocalRunStorage } from '@obversa/runtime/storage/local';

  const definition: DagDefinition = defineGraphDefinition({
    id: 'release-pipeline',
    definitionVersion: 1,
    data: {
      globalConcurrency: 1,
      keyedConcurrency: {},
      stopOnError: true,
      retryCapPerNode: 0,
    },
    nodes: [
      { id: 'draft', data: { kind: 'required', key: null } },
      { id: 'review', data: { kind: 'required', key: null } },
      { id: 'publish', data: { kind: 'required', key: null } },
    ],
    edges: [
      { id: 'draft-to-review', source: 'draft', target: 'review', data: {} },
      { id: 'review-to-publish', source: 'review', target: 'publish', data: {} },
    ],
  });

  const packageIdentity = {
    source: 'npm:@example/release-pipeline',
    version: '1.0.0',
    digest: 'sha256:3333333333333333333333333333333333333333333333333333333333333333',
  } as const;
  const resolution: PlanResolution = {
    package: packageIdentity,
    admission: { package: packageIdentity, permissions: [] },
    executionLanes: [],
  };
  const graph = compileGraph(dagGraphType, definition);
  const plan = resolveGraphPlan(graph.describe(), resolution);

  const storagePolicy = {
    schemaVersion: 1,
    maxEventPayloadBytes: 64_000,
    maxAppendBatchBytes: 128_000,
    maxArtifactBytes: 1_000_000,
    maxTotalArtifactBytesPerRun: 4_000_000,
    retention: 'until-run-delete',
    sensitiveContent: {
      marked: 'reject',
      exact: 'reject',
      freeText: 'redact-before-hash',
    },
  } as const satisfies RunStoragePolicy;

  function nodeBinding(
    root: string,
    runData: NonNullable<GraphNodeBinding['runData']>,
  ): GraphNodeBinding {
    return {
      prompt: null,
      scratchDirectory: root,
      workspace: { mode: 'none', directory: null, allowedPaths: [] },
      trustedCaller: {},
      permissions: [],
      policy: {
        inputBytes: 10_000,
        outputBytes: 10_000,
        timeoutMs: 5_000,
        teardownGraceMs: 100,
        memoryBytes: 10_000_000,
        filesChanged: 0,
        linesChanged: 0,
        callTokens: null,
      },
      resultContract: null,
      runData,
      parseResult: null,
      tokenBudget: null,
      decideAction: async () => ({ kind: 'allow' }),
    };
  }

  const temporaryRoot = await realpath(await mkdtemp(join(tmpdir(), 'obversa-pipeline-')));
  const order: string[] = [];
  try {
    const storage = createLocalRunStorage({
      directory: join(temporaryRoot, 'storage'),
      namespace: 'pipeline-example',
      policy: storagePolicy,
    });
    await persistRunDefinition(storage, {
      runId: 'release-pipeline-run',
      eventId: 'release-pipeline-started',
      timestamp: '2026-01-01T00:00:00.000Z',
      graphDefinition: graph.definition,
      resolvedPlan: plan,
      resolvedInputs: {},
      workspaceBinding: null,
      hostBinding: null,
    });

    const executor = await createGraphExecutor({
      runId: 'release-pipeline-run',
      graph,
      storage,
      nodes: {
        draft: nodeBinding(temporaryRoot, async () => {
          order.push('draft');
          return { article: 'ready' };
        }),
        review: nodeBinding(temporaryRoot, async ({ input }) => {
          order.push('review');
          const draft = (input as {
            readonly results: { readonly draft: { readonly article: string } };
          }).results.draft;
          return { approved: draft.article === 'ready' };
        }),
        publish: nodeBinding(temporaryRoot, async ({ input }) => {
          order.push('publish');
          const review = (input as {
            readonly results: { readonly review: { readonly approved: boolean } };
          }).results.review;
          return { published: review.approved };
        }),
      },
      engines: [],
    });
    const result = await executor.run(new AbortController().signal);
    const storedEvents = [];
    for await (const event of storage.eventStore.read({
      namespace: storage.record.namespace,
      streamId: 'release-pipeline-run',
    })) storedEvents.push(event);

    console.log(JSON.stringify({
      executor: result.kind,
      output: result.kind === 'complete' ? result.output : null,
      dispatches: storedEvents.filter(
        (event) => event.type === 'graph:node-dispatched',
      ).length,
      order,
      planDigest: plan.digest,
      bounds: plan.plan.bounds,
    }, null, 2));
  } finally {
    await rm(temporaryRoot, { recursive: true, force: true });
  }
  ```
</Accordion>

## Next steps

* [Graph executor](/docs/graphs/executor): how a dispatch runs, fallback, and
  resume after a process stops.
* [Plan admission](/docs/graphs/plan-admission): the resolved plan a stored
  pipeline runs under.
* [Running](/docs/concepts/running): the runtime's `pipeline()` for a run you
  start and watch.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.