> ## 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.

# Graph Executor

> Run a stored graph one recorded decision at a time, and resume it from the exact record it stopped on.

Turn a stored run into work, one recorded decision at a time: your graph
decides, the executor records each decision before any node runs, runs the
node inside the [safe attempt lifecycle](/docs/recording/node-attempts), and
records what came back. Use it directly when you write a host. For a
watchdog that starts the executor in a worker, restarts it after a crash
and checks engines for you, read [Supervised local runs](/docs/driving/runner);
this page keeps only what is the executor's own.

The smallest executor binds a stored run to its nodes and engines and runs
it:

```ts examples/preflight-executor.ts (excerpt) {8,11-13} theme={null}
  const makeExecutor = async () => createGraphExecutor({
    ...await bindRun({
      definition: loaded.record.payload.definition,
      scratchDirectory: temporaryRoot,
    }),
    runId,
    storage,
    preflightScratchDirectory: temporaryRoot,
  });

  const paused = await (await makeExecutor()).run(new AbortController().signal);
  if (paused.kind !== 'pause' || !('preflightEventId' in paused)) {
    throw new Error(`Expected a preflight pause, received ${JSON.stringify(paused)}.`);
  }
```

`run()` returns a `GraphExecutorResult`: `complete`, `pause`, `fail` or
`waiting`. An executor needs:

* **The run id and the storage** the run definition and plan were persisted
  to.
* **The compiled graph, node bindings and engine bindings**, here from a
  host module's `bindRun`. Node behaviour is keyed by node id: a data-only
  node gets a function, an engine-backed node a prompt builder that
  receives that dispatch's typed input. The node uses the lane in the
  stored plan, and the engine binding joins that target to the full engine
  identity.
* **`preflightScratchDirectory`** when the plan carries a preflight policy.
  That path and each node binding's `scratchDirectory` must be absolute,
  real and normalised.

## Run a dispatch

1. The executor loads the run definition and the frozen plan from storage
   and checks that the supplied graph matches them.
2. When the plan carries a preflight policy, it checks every engine seat: a
   static check of its identity and, where the lane requires it, one
   tool-free live call with no workspace. The live call proves the seat
   answers and does no work. A failed check pauses the run with
   `PREFLIGHT_PAUSED` or retires the engine, as below. Nothing is
   dispatched before the checks finish.
3. It folds the run's recorded events into graph state and asks the graph
   for its next command.
4. It records the whole dispatch decision, then `node-attempt-started`
   immediately before each node starts.
5. It runs each node through the safe attempt lifecycle, records
   `node-completed`, `node-paused` or `node-failed`, and asks the graph
   again.

Set the lane's `unsupportedStatic` policy to say whether an engine without
`admit` blocks the run or is allowed; the static check of such an engine
comes back `unsupported`, never a success receipt. Usage from the checks is
kept separate from usage from nodes, and missing usage counts as unknown,
not zero.

## Fall back and retire

The host chooses the lane order once. The executor tries the effective
target and at most one live fallback. Only a failure that means the lane is
dead moves to the fallback; a rate limit, a timeout, a policy pause or a
denied action doesn't run again elsewhere.

Each dead-lane event records the selected identity and, when available, the
identity the engine reported. What the failure retires follows what it
proves: bad credentials retire every model on that adapter and provider; a
missing model, exhausted credit or an exhausted quota retire that provider
and model, and the reported ones if they differ; a missing command-line
tool or an invalid configuration retire the adapter; a rate limit or a
transport error retire nothing. When a recorded auth failure names no
provider, the executor recovers one from the lane, and if that's ambiguous
it refuses before any work with `ENGINE_IDENTITY_UNRESOLVED`. Later
dispatches skip whatever was retired.

During a run, if every declared target of a lane is dead, the executor
records `ENGINE_UNAVAILABLE` as a normal node failure, so the graph can
retry, fail or recover in its usual way. Before the first step, when every
target of a lane is excluded or blocked by its checks, the run ends with
`PREFLIGHT_FAILED` instead: terminal, with `preflight:failed` on the record,
`run()` returning `{ kind: 'fail', code: 'PREFLIGHT_FAILED' }`, nothing
dispatched and no pause to resume from.

## Read engine identity records

Every engine call records a `graph:engine-attempt-recorded` event before
fallback routing or the node's result, with the node id, position,
sequence, requested identity and reported effective identity. The sequence
starts at 1 for each position and continues across fallback calls and
resume. When resume recovers an engine node that started and never
finished, it adds a record with both identities null: a mark of
uncertainty, not a count of calls. If a receipt can't be saved, `run()`
rejects and the affected attempt stops before fallback or completion, while
other attempts in the batch can continue. The records describe what the
bound engine reported; they don't see inside an adapter.

## Resume after a stop

A fresh executor returns `waiting`, with the open positions, when a dispatch
has no result. Call `resume` with one exact position, and only after the
earlier process has stopped, because the executor doesn't lock the run
between processes. If two dispatches were open, the first resume returns
`waiting` again with the other position. What resume does depends on how
far the attempt got:

| The record shows | Resume |
| - | - |
| Node code did not start | Runs the attempt at its recorded position. |
| Started, with `retrySafe: true` saved in the start record | Records `node-resumed` and runs the same attempt again. |
| Started, `retrySafe` false, result unknown | Records a pause with a `reconcile-attempt` request naming the attempt, and doesn't call node code. |
| Paused | Records `node-resumed` and offers the same attempt again. |

For a preflight pause, call `resume` with the `{ preflightEventId }` from
the pause the first run returned:

```ts examples/preflight-executor.ts (excerpt) {2-3} theme={null}
  await writeFile(controlFile, 'ready\n');
  const completed = await (await makeExecutor()).resume(
    { preflightEventId: paused.preflightEventId },
    new AbortController().signal,
  );
  if (completed.kind !== 'complete') {
    throw new Error(`Expected completion, received ${JSON.stringify(completed)}.`);
  }
```

The static checks run again and live receipts that still apply are reused.
A stale id is refused with `RESUME_EVENT_MISMATCH`, which names both ids.
The saved `retrySafe` value and the saved plan govern recovery; a change in
the live binding doesn't rewrite an earlier record. Run the file with
`npx tsx preflight-executor.ts`. It stores a one-node graph, runs the
executor against a scripted engine that isn't ready, pauses before any
dispatch, makes the engine ready and resumes with the exact pause event:

```json Output theme={null}
{
  "pause": {
    "phase": "paused",
    "code": "PREFLIGHT_PAUSED",
    "dispatches": 0
  },
  "resume": {
    "usedReturnedToken": true,
    "phase": "admitted",
    "result": "complete",
    "output": {
      "nodes": {
        "check": {
          "checked": "offline"
        }
      }
    }
  },
  "calls": [
    "static",
    "live:not-ready",
    "static",
    "live:ready",
    "ordinary"
  ],
  "temporaryDirectoryRemoved": true
}
```

No node was dispatched before the resume, and the resume used the pause
event the first run returned. `readRunPreflight` reads a stored run's
preflight record without touching an engine.

## Failure

* **`INVALID_PREFLIGHT_CONFIG`**: `preflightScratchDirectory` is optional in
  the type, so TypeScript won't catch a missing path. `await
  createGraphExecutor(...)` rejects with a `GraphExecutionError` of this
  code when the plan needs the path and it's missing or malformed, and a
  relative or symlinked binding path rejects with the same code inside
  `run()` and `resume()` when its lane is checked. The error doesn't say
  which path failed, so handle the code at both call sites.
* **A denied or aborted attempt** records `node-failed` with `DENIED` or
  `ABORTED`. An action-policy wait records `node-paused` with its reason
  and request. A result too large for the event stream records a small
  `node-failed` with `RESULT_TOO_LARGE`, so the dispatch never stays in
  flight.
* **An interrupted check.** A live check is an engine call, and the process
  can die in the middle of it, leaving an open probe on the record. A fresh
  executor's `run()` then fails with `PROTOCOL`, names the probe and tells
  you what to do: check that the old worker and its engine process are
  gone, then call `interruptRunPreflight(storage, runId)`. That call closes
  the probe as `interrupted` and writes a preflight pause in one write that
  fails if anything else wrote first, returns the pause for the next
  resume, returns `null` when no probe is open, and returns the existing
  pause when one is already recorded. It never retries and never calls an
  engine. Call it only when you own the run and have checked the worker is
  gone: if the old call is still alive, the record says interrupted while
  the call keeps spending.

## Limits

* **The executor doesn't lock the run.** Two processes over the same store
  are yours to keep apart; the [supervised runner](/docs/driving/runner) holds a
  process lock for you.
* **One live fallback per attempt.** The lane order is the host's, chosen
  once.

<Accordion title="Full file">
  ```ts examples/preflight-executor.ts theme={null}
  import assert from 'node:assert/strict';
  import { existsSync } from 'node:fs';
  import { mkdtemp, readFile, realpath, rm, writeFile } from 'node:fs/promises';
  import { tmpdir } from 'node:os';
  import { join } from 'node:path';

  import {
    compileGraph,
    createGraphExecutor,
    dagGraphType,
    loadRunDefinition,
    persistRunDefinition,
    readRunPreflight,
    resolveGraphPlan,
    type DagDefinition,
    type ExecutionTarget,
    type RunStoragePolicy,
  } from '@obversa/runtime';
  import { createLocalRunStorage } from '@obversa/runtime/storage/local';
  import { bindRun } from './preflight-host.mjs';

  const target = {
    adapter: 'scripted-local',
    provider: 'local',
    modelFamily: 'scripted',
    model: 'offline-check',
    tools: [],
  } as const satisfies ExecutionTarget;

  const definition = {
    id: 'preflight-executor',
    definitionVersion: 1,
    data: {
      globalConcurrency: 1,
      keyedConcurrency: {},
      stopOnError: true,
      retryCapPerNode: 0,
    },
    nodes: [{
      id: 'check',
      data: {
        kind: 'required',
        key: null,
        lane: { id: 'local', requested: target, knownSubstitutions: [] },
      },
    }],
    edges: [],
  } as const satisfies DagDefinition;

  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;

  const temporaryRoot = await realpath(await mkdtemp(join(tmpdir(), 'obversa-preflight-executor-')));
  const controlFile = join(temporaryRoot, 'control.txt');
  const runId = 'preflight-executor-run';
  const callsFile = join(temporaryRoot, 'engine-calls.log');
  let report: Record<string, unknown> | undefined;

  try {
    await writeFile(controlFile, 'not-ready\n');
    await writeFile(callsFile, '');
    const graph = compileGraph(dagGraphType, definition);
    const packageIdentity = {
      source: 'npm:@example/preflight-executor',
      version: '1.0.0',
      digest: 'sha256:4444444444444444444444444444444444444444444444444444444444444444',
    } as const;
    const plan = resolveGraphPlan(graph.describe(), {
      package: packageIdentity,
      admission: { package: packageIdentity, permissions: [] },
      executionLanes: [{ id: 'local', effective: target }],
      preflight: {
        timeoutMs: 2_000,
        lanes: [{ laneId: 'local', live: 'required', unsupportedStatic: 'block' }],
      },
    });
    const storage = createLocalRunStorage({
      directory: join(temporaryRoot, 'storage'),
      namespace: 'preflight-executor-example',
      policy: storagePolicy,
    });
    await persistRunDefinition(storage, {
      runId,
      eventId: 'preflight-executor-started',
      timestamp: new Date().toISOString(),
      graphDefinition: graph.definition,
      resolvedPlan: plan,
      resolvedInputs: { controlFile, callsFile },
      workspaceBinding: null,
      hostBinding: null,
    });

    const loaded = await loadRunDefinition(storage, runId);
    const makeExecutor = async () => createGraphExecutor({
      ...await bindRun({
        definition: loaded.record.payload.definition,
        scratchDirectory: temporaryRoot,
      }),
      runId,
      storage,
      preflightScratchDirectory: temporaryRoot,
    });

    const paused = await (await makeExecutor()).run(new AbortController().signal);
    if (paused.kind !== 'pause' || !('preflightEventId' in paused)) {
      throw new Error(`Expected a preflight pause, received ${JSON.stringify(paused)}.`);
    }
    assert.equal(paused.code, 'PREFLIGHT_PAUSED');
    assert.equal((await readRunPreflight(storage, runId)).phase, 'paused');

    let dispatchesBeforeResume = 0;
    for await (const event of storage.eventStore.read({
      namespace: storage.record.namespace,
      streamId: runId,
    })) {
      if (event.type === 'graph:node-dispatched') dispatchesBeforeResume += 1;
    }
    assert.equal(dispatchesBeforeResume, 0);

    await writeFile(controlFile, 'ready\n');
    const completed = await (await makeExecutor()).resume(
      { preflightEventId: paused.preflightEventId },
      new AbortController().signal,
    );
    if (completed.kind !== 'complete') {
      throw new Error(`Expected completion, received ${JSON.stringify(completed)}.`);
    }
    assert.deepEqual(completed.output, { nodes: { check: { checked: 'offline' } } });
    const finalPreflight = await readRunPreflight(storage, runId);
    assert.equal(finalPreflight.phase, 'admitted');
    assert.equal(finalPreflight.resumedPreflightEventId, paused.preflightEventId);
    const calls = (await readFile(callsFile, 'utf8')).trim().split('\n').filter(Boolean);
    assert.deepEqual(calls, ['static', 'live:not-ready', 'static', 'live:ready', 'ordinary']);

    report = {
      pause: { phase: 'paused', code: paused.code, dispatches: dispatchesBeforeResume },
      resume: {
        usedReturnedToken: finalPreflight.resumedPreflightEventId === paused.preflightEventId,
        phase: finalPreflight.phase,
        result: completed.kind,
        output: completed.output,
      },
      calls,
    };
  } finally {
    await rm(temporaryRoot, { recursive: true, force: true });
  }

  console.log(JSON.stringify({
    ...report,
    temporaryDirectoryRemoved: !existsSync(temporaryRoot),
  }, null, 2));
  ```
</Accordion>

The engine and its bindings are in `examples/preflight-host.mjs`. The types
are `GraphExecutorResult` for what `run()` and `resume()` return,
`RunPreflightPolicy` for the plan's checks, `RunPreflightState` for what
`readRunPreflight` returns, and `PreflightPauseResult` and
`PreflightFailureResult` for the two ways the checks stop a run early.

## Next steps

* [Supervised local runs](/docs/driving/runner): the watchdog that runs this
  executor in a worker and restarts it.
* [Safe node attempts](/docs/recording/node-attempts): what one dispatch records.
* [Plan admission](/docs/graphs/plan-admission): the frozen plan the executor
  checks the graph against.


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