/** * src/engine/dispatch.ts — engine assembly: gates -> registry -> models -> * lanes -> result, with run lifecycle + hash-only ledger. * * This is the final assembly layer. Every dispatch runs the full preflight * BEFORE any spawn; a child is NEVER launched if any gate fails (I2/I10). The * six-part contract, path policy, write-scope, tool allowlist, output contract * and explicit-model gate are all enforced up front; failures yield a * `preflight_failed` run (failureKind `preflight` or `config`) and no child. * * runSingle / runParallel / runChain dispatch one or many runs through a * shared LanePool (hot sessions). runChain substitutes `{previous}` with the * prior step's output and reuses the SAME lane/sessionPath for continuity. * * Zero @earendil-works/* imports; the only child-process touch is the * injectable spawner supplied by src/lanes (never a hard-coded pi spawn here). */ import type { AgentScope, ChildProgressEvent, ChildResult, DelegationDetails } from "../core/types.js"; import { type ResolveAgentDirs } from "../registry/index.js"; import { type EnabledModelsReader, type ModelByClassMap, type VerifiedModelCatalog } from "../models/index.js"; import { type ProviderExtensionResolver } from "../models/provider-extensions.js"; import { type SpawnFn } from "../lanes/index.js"; import { type DelegationMonitorState, type DelegationRunMode, type DelegationRunSource } from "./runs.js"; import { type SteerOutcome } from "./steer.js"; import { type ContextMode, type ContextProvider } from "./context.js"; import { type MemoryScope } from "./memory.js"; import type { BackgroundRunRegistry } from "./background.js"; /** Concurrency-limited async mapper (used by runParallel). */ export declare function mapWithConcurrency(items: readonly T[], concurrency: number, fn: (item: T, index: number) => Promise): Promise; /** Local `usageEmpty` mirror (telemetry.ts in the harness). */ export declare function usageEmpty(): ChildResult["usage"]; /** Normalize a declarative allowlist (trim/drop empties; undefined = none). */ export declare function normalizeNestedAllowlist(values: readonly string[] | undefined): string[] | undefined; /** * Byte size of a session file (0 when missing/unreadable). Used as the * continuation base offset so already-known session content is not re-fed. */ export declare function getSessionFileSize(sessionPath: string): number; export interface DispatchEngineOptions { repoRoot: string; monitor?: DelegationMonitorState; spawn?: SpawnFn; piCommand?: string; sessionDir?: string; env?: NodeJS.ProcessEnv; maxParallel?: number; resolveAgentDirs?: ResolveAgentDirs; scope?: AgentScope; verifiedCatalog?: VerifiedModelCatalog; /** * F2: injectable verified-catalog source evaluated at EVERY preflight, so a * `/.pi/model-catalog.json` created or updated mid-session is * honored without rebuilding the engine. When set it wins over the static * `verifiedCatalog` value. The reader must never throw; a throwing reader * is reported as a catalog read error (the override stays blocked). */ verifiedCatalogReader?: () => VerifiedModelCatalog; classModels?: ModelByClassMap; parentModel?: string; /** * C3 model scope: injectable source of the pi `enabledModels` allowlist. * Undefined (or an allowlist that reads back absent/empty) skips the scope * check entirely (no-op safety). The concrete reader (pi settings, project * over global) is supplied by src/extension; the engine stays fs-free. */ enabledModelsReader?: EnabledModelsReader; ledgerDir?: string; /** When true, ledger entries are persisted to ``. */ persistLedger?: boolean; /** Steer-file dir for mid-run steering (default `/.pi/steer`). */ steerDir?: string; /** Engine-level B3 default: max turns before the soft wrap-up steer (undefined = unlimited). */ maxTurns?: number; /** Engine-level B3 default: grace turns before the hard abort (default 5, min 1). */ graceTurns?: number; parentToolCallId?: string; source?: DelegationRunSource; onLedger?: (entry: Record) => void; now?: () => number; backgroundRegistry?: BackgroundRunRegistry; /** * C2: injectable main-context source. Supplied by the extension (built from * the live pi session); the engine NEVER reads the session itself. Used only * when a dispatch sets `contextMode: "main"`. */ contextProvider?: ContextProvider; /** * C5: root for `user`-scope agent memory (`/agent-memory//`). * Injectable so hosts and tests never touch the real home directory; * defaults to `/.pi/agent` inside resolveMemoryDir. */ memoryAgentDir?: string; /** * C1 nested subagents: explicit max-subagent-depth override (clamped 0..4). * Defaults to the shared config value (`maxSubagentDepth` from * `.subagents/config.json` or `SUBAGENTS_MAX_DEPTH`, default 2). 0/1 = off. */ maxSubagentDepth?: number; /** * P2 provider-extension resolver: provider id -> extension entry path. * Children spawn `--no-extensions`; when the RESOLVED child model has a * provider with a known extension, the spawn argv gains `-e ` (a * custom-provider child like `ollama-cloud/...` otherwise crashes with * 'Model not found'). Injectable for tests; the default resolver is built * ONCE from the shared config — `.subagents/config.json` * `providerExtensions` map over the built-in * `ollama-cloud` -> `~/pi-provider-ollama-cloud/src/index.ts` default — with * every path existsSync-guarded (a path that does not exist is never * injected). A throwing resolver degrades to no extension, never a crash. */ providerExtensions?: ProviderExtensionResolver; /** * Engine-level streaming child progress feed: ChildProgressEvent per * assistant turn / captured toolCall / completed text of every dispatched * run. Overridable per run (per-run option wins). In-memory only. */ onChildEvent?: (event: ChildProgressEvent) => void; } export interface SingleDispatchOptions { task?: string; cwd?: string; model?: string; modelClass?: string; thinking?: string; tools?: string[]; allowedPaths?: string[]; forbiddenPaths?: string[]; outputContract?: string; index?: number; runId?: string; lane?: string; mode?: DelegationRunMode; parentToolCallId?: string; background?: boolean; source?: DelegationRunSource; scope?: AgentScope; /** Explicit session path override (continue reuses the original run's session file). */ sessionPath?: string; /** B3 per-run override: max turns before the soft wrap-up steer (wins over engine default). */ maxTurns?: number; /** B3 per-run override: grace turns before the hard abort (wins over engine default). */ graceTurns?: number; /** * C2 context projection mode. `isolated` (default): the child gets the task * verbatim, no parent context. `main`: the task is wrapped with the condensed * main-context projection (compaction + last-20 + session-file pointer) in a * REQUEST/HISTORY authority hierarchy (anti-persona-bleed). Requires the * engine-level `contextProvider`; without one, `main` degrades gracefully to * an unwrapped (isolated-equivalent) dispatch. */ contextMode?: ContextMode; /** * C6 worktree isolation: run the child inside a disposable DETACHED git * worktree of the repo (child cwd = worktree path, monorepo subfolders * mirrored). After the run the worktree is cleaned up: changes are * auto-committed onto a pi-agent- branch (timestamp suffix on * conflict); a clean run is just removed. The result carries the outcome on * `worktree` (see getWorktreeOutcome in lanes/worktree.ts). Requires a git * repository with at least one commit (enforced at preflight). */ isolation?: "worktree"; /** * C5 persistent memory scope override for this run. Wins over the agent * card frontmatter `memory` field; `false` explicitly disables memory even * when the card requests it. When memory is active, a memory block is * appended to the agent SYSTEM PROMPT (after the card prompt) — never into * the task. Read-only vs R/W is detected from the card tools (write/edit). */ memory?: MemoryScope | false; /** * C1 nested subagents: enable the child's nested `subagent` tool by * positioning the nesting envs on the child spawn (PI_SUBAGENTS_NESTED=1, * PI_SUBAGENTS_ALLOWED_SUBAGENTS, PI_SUBAGENTS_DEPTH=1 and * PI_SUBAGENTS_MAX_DEPTH from the config maxSubagentDepth, clamped 0..4). * `true` uses the agent card's `allowed_subagents` frontmatter allowlist; * `{ allowedSubagents: [...] }` (names or `all`) overrides the card. * Nesting stays OFF without a non-empty allowlist or when maxSubagentDepth * < 2 (0/1 = off). */ nested?: NestedDispatchOptions | true; /** * F8: fresh per-run session (default TRUE for single/parallel). Each run * gets its OWN lane/session file `-.jsonl`, so successive * dispatches to the same agent never stack onto one shared session * (context-bleed fix). Continuity paths are structural and unaffected: * chain shares an explicit lane and continue_run reuses the original run's * explicit sessionPath (the pool's sessionPath override wins over the * lane). `fresh: false` restores the legacy stable agent-named lane * (`.jsonl`); an explicit `lane` always wins over this flag. */ fresh?: boolean; /** * Per-run streaming child progress feed (wins over the engine-level * onChildEvent). Fired per assistant turn / captured toolCall / completed * text with runId + agent filled by the engine. In-memory only. */ onChildEvent?: (event: ChildProgressEvent) => void; } /** Options for `continueRun`: the run id to continue plus single-run options. */ export interface ContinueRunOptions extends SingleDispatchOptions { /** Optional id for the new (continued) run; defaults to makeRunId("delegate"). */ runId?: string; } /** C1 nested-subagent dispatch option: the strict grandchild allowlist. */ export interface NestedDispatchOptions { /** Allowed grandchild agent names (or `all`); wins over the agent card. */ allowedSubagents?: string[]; } export interface ParallelTaskInput { agent: string; task: string; model?: string; modelClass?: string; cwd?: string; tools?: string[]; allowedPaths?: string[]; forbiddenPaths?: string[]; outputContract?: string; } export interface ParallelDispatchOptions { concurrency?: number; parentToolCallId?: string; source?: DelegationRunSource; /** Batch-level child progress feed, forwarded to every task's run (per-run option wins over engine). */ onChildEvent?: (event: ChildProgressEvent) => void; } export interface ChainStepInput { agent: string; task: string; model?: string; modelClass?: string; cwd?: string; tools?: string[]; allowedPaths?: string[]; forbiddenPaths?: string[]; outputContract?: string; } export interface ChainDispatchOptions { parentToolCallId?: string; source?: DelegationRunSource; lane?: string; /** Chain-level child progress feed, forwarded to every step's run (per-run option wins over engine). */ onChildEvent?: (event: ChildProgressEvent) => void; } /** * The final engine assembly. Owns a bounded monitor, a lazily-created LanePool * (hot sessions), model routing, preflight gates, and the hash-only ledger. */ export declare class DispatchEngine { readonly repoRoot: string; readonly monitor: DelegationMonitorState; private readonly options; private pool?; constructor(options: DispatchEngineOptions); private get now(); private get source(); private get parentToolCallId(); /** Steer-file dir shared by manual steers and B3 wrap-up steers. */ private get steerDir(); private getPool; /** * C1: spawn wrapper merging per-dispatch nested env overrides. The pool's * env is shared across lanes, so the engine scopes the overrides with an * AsyncLocalStorage around each pool.dispatch call — concurrent dispatches * each get their own nested env (or none at all). */ private wrappedSpawn; private nestedConfig?; private providerExtensionResolver?; /** C1 config: max depth (clamped 0..4) + absolute agents dir, memoized. */ private loadNestedConfig; /** * C1: resolve the nesting envs for a child spawn. Undefined = nesting OFF * (option absent, maxSubagentDepth < 2 (0/1 = off), or no allowlist — the * strict posture: no allowlist ever means no nested spawning). */ private resolveNestedSpawnEnv; /** * P2: resolve the provider extension (`-e`) for the child's EFFECTIVE * model. Injectable resolver wins; the lazy default is built once from the * shared config (providerExtensions map + existsSync-guarded ollama-cloud * default). Builtin/no-slash models and unknown providers get undefined; a * throwing resolver degrades to undefined — never a dispatch crash. */ private resolveProviderExtension; /** Close the owned lane pool (aborts in-flight tasks, no leaks). */ close(): void; /** * Steer a running run with a mid-run message (B2, parent side). Writes the * `.steer` file and marks the run view `steered`. Refuses unknown or * non-running runs. Returns a structured outcome (never throws). */ steer(runId: string, message: string): SteerOutcome; /** * Dispatch a single run. If `opts.background` is set and a background * registry is available, returns immediately with a running stub while the * child runs in the background; otherwise awaits the child result. */ single(agent: string, task: string, opts?: SingleDispatchOptions): Promise; private singleBackground; /** Core single-run path (foreground). Runs preflight then dispatches on a lane. */ runOne(agent: string, task: string, opts: SingleDispatchOptions): Promise; /** * Continue a TERMINAL run with its original session file. * * Finds the run (monitor or background registry), verifies it is terminal * (complete/failed/aborted), reuses the SAME sessionPath — pi resumes the * session file AS-IS (no --session-offset: the option does not exist in * pi >= 0.84 and would exit 1; the byte-offset stays ledger metadata only). * The new run links `continuedFromRunId` and increments `turnCount`. * Non-terminal runs are refused with a clear error. */ continueRun(runId: string, task: string, opts?: ContinueRunOptions): Promise; /** Dispatch a batch of tasks concurrently with a configurable cap. */ parallel(tasks: readonly ParallelTaskInput[], opts?: ParallelDispatchOptions): Promise; /** * Dispatch a chain of steps. Each step's task may contain the `{previous}` * placeholder, which is substituted with the previous step's output. All * steps share the SAME lane (hence sessionPath) for context continuity. */ chain(steps: readonly ChainStepInput[], opts?: ChainDispatchOptions): Promise; private runPreflight; /** Record non-blocking preflight warnings (C3 model scope) on the run view. */ private recordRunWarnings; /** * Live usage patch on the run view from a child `kind=turn` event * (cumulative turns + current context snapshot + model). Best effort — a * live-usage failure must never break a dispatch (I10). */ private recordLiveUsage; /** * F2: resolve the verified catalog for one preflight. The per-dispatch * reader wins over the static value; a throwing reader is captured as a * readError (blocking, honest) instead of crashing the dispatch. */ private readVerifiedCatalog; private failPreflight; private preflightFailureResult; private dispatchOnLane; private settleRun; /** * C4: write the hash-only attestation sidecar for a settled run to * `/attestations/.json`. Only when * ledger persistence is enabled (same posture as the ledger itself); * best-effort — returns undefined on ANY error, never throws, never * blocks settlement. */ private writeAttestationForRun; private appendLedger; } /** * Standalone single-run dispatch. Builds a throwaway engine over `ctx`, * dispatches one run, and closes the pool. */ export declare function runSingle(ctx: Omit & { repoRoot: string; }, agent: string, task: string, opts?: SingleDispatchOptions): Promise; /** Standalone parallel dispatch (shared pool, configurable cap). */ export declare function runParallel(ctx: Omit & { repoRoot: string; }, tasks: readonly ParallelTaskInput[], opts?: ParallelDispatchOptions): Promise; /** Standalone chain dispatch (same lane/sessionPath for continuity). */ export declare function runChain(ctx: Omit & { repoRoot: string; }, steps: readonly ChainStepInput[], opts?: ChainDispatchOptions): Promise;