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

# Supervised Local Runs

> If the worker dies, a new one reads the record and a finished step does not run again. Engines are checked before the first step.

If the worker dies, a new one reads the record and a finished step does not
run again. Start a graph in a worker and let the calling process be its
watchdog: it checks each agent you named before the first step, waits for
the worker, stops
its child processes, and restarts it within the time and restart limits you
set. Use it for a run that must survive a crash or an engine that isn't
ready yet. For a run you start and watch yourself, `run()` from the
[runtime](/docs/concepts/running) is enough.

The run is durable: every step is written to the run's record, so a killed
run carries on where it stopped. The smallest supervised run starts one
graph under the watchdog:

```ts examples/preflight-supervised-run.ts (excerpt) {2-4,13-15} theme={null}
  handle = await startSupervisedRun({
    directory: runnerDirectory,
    runRoot,
    module: './host.mjs',
    storage,
    workspace,
    definition: {
      runId,
      graphDefinition: graph.definition,
      resolvedPlan,
      resolvedInputs: { controlFile, callsFile },
    },
    limits: { timeoutMs: 20_000, maxDispatches: 1 },
    restart,
    teardownGraceMs: 100,
  });
```

`startSupervisedRun` returns a handle with `done`, `status()` and `stop()`.
`done` resolves to the graph's result, `complete`, `pause` or `fail`, and
infrastructure errors reject it. A supervised run needs:

* **A runner directory and a worker root**, with the host module as a
  relative specifier inside that root.
* **The storage settings and the workspace provider.**
* **The stored definition**: the run id, the graph, the resolved plan and
  the resolved inputs.
* **Limits, a restart policy and a cleanup grace.**

## Bind the host module

The runner loads host modules and supervises processes; the runtime never
loads a module. The module exports `bindRun`, which receives the stored
definition and a scratch directory and returns the compiled graph, the node
bindings and the engine bindings. `examples/preflight-host.mjs` is one; the
example copies it into the disposable repository it creates and the runner
loads that copy.

The module must stay inside the worker root after symlinks resolve; a path
outside it is refused with `HOST_MODULE`. The worker root must be the top
level of the repository the workspace provider captures; a provider for
another repository is refused with `WORKSPACE_ROOT` before any run is
stored. The runner stores the module's
path and byte digest and checks the digest before and after each import.
Changed bytes fail with `HOST_MODULE_CHANGED`. The digest covers the module
file alone, not what it imports.

The worker inherits only `PATH`, `HOME`, `TMPDIR`, `TMP`, `TEMP`,
`SystemRoot`, `USERPROFILE` and `PATHEXT` from the watchdog. To pass an
engine credential, list its variable name in `environmentVariables` on the
start or resume call. The runner copies the value when it's set and stores
neither the name nor the value in the run. `resolvedInputs` is durable
storage, so keep secrets out of it.

## Check the engines

When the stored plan carries a preflight policy, the worker checks every
declared engine seat before it dispatches anything.

* **A static check** asks the engine to admit the seat, with the node's
  real configuration and no prompt. The engine reports the identity it will
  run under (adapter, provider, model family, model, executable), and the
  runtime compares it with the saved identity. An engine that doesn't
  support the check is `unsupported`, and the lane's policy says whether
  that blocks the seat or is allowed.
* **A live check**, when the lane requires one, is one tool-free engine
  call with no workspace. It runs once per eligible seat, one at a time,
  and stops at the first seat that answers. It proves the seat can answer,
  not only that it's configured.
* **A failed check that retires nothing pauses the run** on a
  `PREFLIGHT_PAUSED` record before any dispatch, carrying the pause event
  id. Nothing has run, so there's nothing to reconcile.
* **A lane with no admissible seat ends the run.** When every target of a
  lane is excluded or blocked, the worker returns `PREFLIGHT_FAILED` and
  `done` resolves to `fail` with that code. Nothing was dispatched and
  there's no pause to resume from.
* **A check the worker died in the middle of becomes a pause.** Once the
  watchdog has verified the worker's processes are gone, it closes the
  check as interrupted with `interruptRunPreflight` from
  `@obversa/runtime`, and `done` resolves to a `PREFLIGHT_PAUSED` pause you
  resume like any other. If the watchdog can't verify cleanup, the check
  stays open.

Command-line adapters record the path and version they were admitted with,
not a hash of the executable. An API-key adapter is admitted locally with no
executable and no capabilities, and its retry layers are off for the check
only. Environment variable names are validated and their values captured
before any asynchronous work; no value is written to run storage.

## Start, inspect, stop

Call `stop()` while work is running to request a stop and await cleanup. A
clean stop without a recorded worker result returns `fail` with `STOPPED`;
after completion it returns the settled result. When the worker has
recorded a result and cleanup and lease release are verified, the watchdog
returns that result even if a stop or timeout arrives before the worker
exits. Cleanup and lease failures take precedence over a recorded result.

`readSupervisedRunStatus` reads a run's recorded progress from another
process with the same storage settings. It doesn't take ownership of the
run or stop another process's watchdog.

Completion output is stored once as an artifact under the storage policy's
limits, and the watchdog verifies the stored bytes before returning
`done.output`. Engine responses are stored the same way, as
`runner-engine-parts` artifacts. A response that fits the node's output
limit can still exceed a storage limit: a size or quota refusal raises
`StorageError` with `STORAGE_LIMIT_EXCEEDED`, and a match on a known secret
raises `KNOWN_SECRET`. Inside the worker the refused write becomes a failed
node with `EFFECT_FAILED`, and with the built-in dag form a failed required
node makes `done` resolve to `fail` with `DAG_NODE_FAILED`.

## Recover after a crash

The watchdog waits for worker exit and cleanup, then starts a replacement
that reads the record and decides what remains.

| Crash boundary | Recovery |
| - | - |
| Before a dispatch is recorded | Decide and dispatch the work. |
| After dispatch, before node code starts | Start the recorded occurrence. |
| After node code starts, before its result is recorded | Pause for reconciliation unless the binding declares `retrySafe: true`. |
| After the result is recorded | Continue without repeating the finished occurrence. |

A saved start without a saved result doesn't prove whether an outward effect
happened. Set `retrySafe: true` only when repeating the node is acceptable
after an unknown outcome. Otherwise the run pauses with a
`reconcile-attempt` request for a person to answer.

The watchdog holds the workspace lease while the worker runs and while
cleanup executes. Cleanup uses the public `@obversa/core/command` API and
its owner markers; on Linux it also finds processes carrying the worker's
inherited marker, and elsewhere it follows the observed process tree.
Incomplete cleanup or a failed lease release retains the process lock and
starts no further worker. The lock is scoped to the storage directory,
namespace and run id: a second watchdog for the same run gets
`PROCESS_LOCKED`, and if a watchdog dies its lock and lease remain, with no
forced takeover.

## Resume a paused run

`resumeSupervisedRun` reopens one recorded pause. For a preflight pause,
pass the run id and the `preflightEventId`; for a node pause, the run id
and the exact `position` from its `graph:node-paused` event:

```ts examples/preflight-supervised-run.ts (excerpt) {3,11} theme={null}
  await handle.stop();
  handle = undefined;
  handle = await resumeSupervisedRun({
    directory: runnerDirectory,
    runRoot,
    storage,
    workspace,
    restart,
    teardownGraceMs: 100,
    runId,
    preflightEventId: paused.preflightEventId,
  });
```

The definition, host module and limits come from the stored run. Resume
writes no second start event and repeats no finished occurrence. For a
preflight pause the static checks run again, the live receipts that still
apply are reused, and a stale id is refused with `RESUME_EVENT_MISMATCH`
naming both ids.

Before releasing ownership after a node pause, the watchdog saves a snapshot
of the workspace. Resume checks the workspace against it before starting a
worker and never adopts edits made while paused. A changed workspace pauses
again with `WORKSPACE_DRIFT`, a missing snapshot with
`WORKSPACE_ANCHOR_MISSING`, an unreadable one with
`WORKSPACE_ANCHOR_INVALID`. If the snapshot can't be captured at pause time,
the run fails with `WORKSPACE_ANCHOR_WRITE` and no resumable pause is
recorded. The host's action policy runs again on resume; record approval
where that policy reads it, because calling resume isn't an approval.

Run the example with `npx tsx preflight-supervised-run.ts`. It starts a
one-node graph under the watchdog with a scripted engine that isn't ready,
pauses before any work is dispatched, makes the engine ready, stops the old
watchdog, and resumes from the recorded pause:

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

Five engine calls, in order: a static check of the seat, a live check that
found it not ready, and after resume a second static check, a live check
that found it ready, and the one ordinary call the graph asked for. No node
was dispatched before the resume, and the resume used the exact pause event
the first run returned.

## Read progress

Status carries the run phase, worker liveness, inspected processes, the
restart count, backoff, elapsed and remaining time, and pause reasons.
`cleanupVerified` is `null` before a terminal result, `true` when cleanup
was verified within the platform's capability, and `false` when it
couldn't be. `leaseRetained` reports a lease still held after a terminal
failure. A dead worker or an empty process list isn't proof of cleanup:
read these two fields before treating the workspace as released.

Usage is reported per engine call. Missing usage isn't zero: a call without
a receipt counts as unknown, and a node's usage is `partial` when some of
its calls have receipts and some don't. The check calls before the run
report their own usage, separate from the nodes'.

## Failure

A check that fails tells the run what not to try again, and the scope
follows what the failure proves:

* **Bad credentials** retire every model on that adapter and provider,
  because the credential is theirs.
* **A missing model, exhausted credit or an exhausted quota** retire that
  provider and model, and also the provider and model the engine reported,
  if they differ. A quota is an allowance gone for hours or longer.
* **A missing command-line tool or an invalid configuration** retire the
  adapter.
* **A rate limit or a transport error** retire nothing, because they clear
  in seconds or minutes.
* **An old auth record with no provider** recovers one from its lane. When
  that's ambiguous, the run refuses before any work with
  `ENGINE_IDENTITY_UNRESOLVED`.

If the pause event itself can't be written, the run returns `fail` with
`RUN_STORAGE` and no resumable pause exists. If the pause was written and
only the watchdog's final record failed, the result is a pause with
`RUN_STORAGE`, and a later resume can reopen it. Neither path retries the
write on its own.

## Limits

* **Elapsed time.** `limits.timeoutMs` covers worker execution and restart
  backoff across pauses and resumes. It counts from the stored run
  timestamp, freezes at each settled pause, and resumes at the next worker
  launch. Resume and preflight checks before launch don't spend it; engine
  checks inside the worker do.
* **Dispatches.** `limits.maxDispatches` counts recorded graph dispatches
  across workers.
* **Restarts.** `restart.maxRestarts` caps replacement workers, with backoff
  from `initialBackoffMs` to `maxBackoffMs`.
* **Cleanup grace.** `teardownGraceMs` is the time processes get to stop
  before they are killed.
* **Worker output.** Each worker launch may write 1,000,000 bytes to its
  combined stdout and stderr. Beyond that, the watchdog stops the worker
  without restarting it and `done` resolves to `fail` with `OUTPUT_LIMIT`.

<Accordion title="Full file">
  ```ts examples/preflight-supervised-run.ts theme={null}
  import assert from 'node:assert/strict';
  import { execFile } from 'node:child_process';
  import { createHash } from 'node:crypto';
  import { existsSync } from 'node:fs';
  import { createRequire } from 'node:module';
  import { mkdir, mkdtemp, readFile, realpath, rm, symlink, writeFile } from 'node:fs/promises';
  import { tmpdir } from 'node:os';
  import { dirname, join, relative } from 'node:path';
  import { promisify } from 'node:util';

  import {
    compileGraph,
    createGitWorktreeProvider,
    dagGraphType,
    readRunPreflight,
    resolveGraphPlan,
    type DagDefinition,
    type RunStoragePolicy,
  } from '@obversa/runtime';
  import { createLocalRunStorage } from '@obversa/runtime/storage/local';
  import {
    resumeSupervisedRun,
    startSupervisedRun,
    type SupervisedRunHandle,
  } from '@obversa/runner';

  const git = promisify(execFile);
  const temporaryRoot = await realpath(await mkdtemp(join(tmpdir(), 'obversa-preflight-supervised-')));
  const runRoot = join(temporaryRoot, 'workspace');
  const controlFile = join(temporaryRoot, 'control.txt');
  const callsFile = join(temporaryRoot, 'engine-calls.log');
  const runnerDirectory = join(temporaryRoot, 'runner');
  const runId = 'preflight-supervised-run';
  let handle: SupervisedRunHandle | undefined;
  let startupAttempted = false;
  let report: Record<string, unknown> | undefined;

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

  const definition = {
    id: 'preflight-supervised',
    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;

  try {
    await writeFile(controlFile, 'not-ready\n');
    await writeFile(callsFile, '');
    assert.equal(relative(runRoot, controlFile).startsWith('..'), true);

    await mkdir(runRoot);
    const hostBytes = await readFile(new URL('./preflight-host.mjs', import.meta.url));
    await writeFile(join(runRoot, 'host.mjs'), hostBytes);
    await writeFile(join(runRoot, '.gitignore'), 'node_modules/\n');
    await git('git', ['init', '-q', '-b', 'main'], { cwd: runRoot });
    await git('git', ['add', 'host.mjs', '.gitignore'], { cwd: runRoot });
    await git('git', [
      '-c', 'user.name=Example',
      '-c', 'user.email=example@example.com',
      '-c', 'commit.gpgsign=false',
      'commit', '-qm', 'chore: initialize disposable preflight example',
    ], { cwd: runRoot });

    await mkdir(join(runRoot, 'node_modules/@obversa'), { recursive: true });
    await symlink(
      dirname(createRequire(join(process.cwd(), 'package.json')).resolve('@obversa/runtime/package.json')),
      join(runRoot, 'node_modules/@obversa/runtime'),
      'junction',
    );

    const graph = compileGraph(dagGraphType, definition);
    const packageIdentity = {
      source: 'file:host.mjs',
      version: '1.0.0',
      digest: `sha256:${createHash('sha256').update(hostBytes).digest('hex')}`,
    } as const;
    const resolvedPlan = 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 = {
      directory: join(temporaryRoot, 'storage'),
      namespace: 'preflight-supervised-example',
      policy: storagePolicy,
    } as const;
    const workspace = createGitWorktreeProvider({ repositoryPath: runRoot });
    const restart = {
      maxRestarts: 1,
      initialBackoffMs: 100,
      maxBackoffMs: 1_000,
    } as const;

    startupAttempted = true;
    handle = await startSupervisedRun({
      directory: runnerDirectory,
      runRoot,
      module: './host.mjs',
      storage,
      workspace,
      definition: {
        runId,
        graphDefinition: graph.definition,
        resolvedPlan,
        resolvedInputs: { controlFile, callsFile },
      },
      limits: { timeoutMs: 20_000, maxDispatches: 1 },
      restart,
      teardownGraceMs: 100,
    });
    const paused = await handle.done;
    if (paused.kind !== 'pause' || paused.code !== 'PREFLIGHT_PAUSED' || paused.preflightEventId === undefined) {
      throw new Error(`Expected a preflight pause, received ${JSON.stringify(paused)}.`);
    }
    const pausedStatus = await handle.status();
    assert.equal(pausedStatus.phase, 'paused');
    assert.equal(pausedStatus.cleanupVerified, true);
    assert.equal(pausedStatus.leaseRetained, false);

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

    await writeFile(controlFile, 'ready\n');
    await handle.stop();
    handle = undefined;
    handle = await resumeSupervisedRun({
      directory: runnerDirectory,
      runRoot,
      storage,
      workspace,
      restart,
      teardownGraceMs: 100,
      runId,
      preflightEventId: paused.preflightEventId,
    });
    const completed = await handle.done;
    if (completed.kind !== 'complete') {
      throw new Error(`Expected completion, received ${JSON.stringify(completed)}.`);
    }
    assert.deepEqual(completed.output, { nodes: { check: { checked: 'offline' } } });
    const completedStatus = await handle.status();
    assert.equal(completedStatus.phase, 'completed');
    assert.equal(completedStatus.cleanupVerified, true);
    assert.equal(completedStatus.leaseRetained, false);
    const finalPreflight = await readRunPreflight(runStorage, 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: pausedStatus.phase, code: paused.code, dispatches: dispatchesBeforeResume },
      resume: {
        usedReturnedToken: finalPreflight.resumedPreflightEventId === paused.preflightEventId,
        phase: completedStatus.phase,
        result: completed.kind,
        output: completed.output,
      },
      calls,
      controlOutsideCapturedWorkspace: relative(runRoot, controlFile).startsWith('..'),
    };
  } finally {
    // Preserve the lock and lease if cleanup cannot be verified.
    if (handle !== undefined) {
      await handle.stop();
      const status = await handle.status();
      assert.equal(status.cleanupVerified, true, 'Retain the directory when cleanup is unverified.');
      assert.equal(status.leaseRetained, false, 'Retain the directory while its workspace lease is held.');
    }
    // Startup can fail before returning a handle while retaining ownership.
    if (!startupAttempted || handle !== undefined) {
      await rm(temporaryRoot, { recursive: true, force: true });
    }
  }

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

The engine and its bindings live in `examples/preflight-host.mjs`, the host
module the example copies into its disposable repository.

## Next steps

* [Runner](/docs/packages/runner): the package's public entry points.
* [The record](/docs/concepts/record): what a resumed run skips, repeats and
  asks a person to reconcile.
* [Watch a run in the browser](/docs/driving/monitor): the page a supervised run
  serves while it works.


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