import { randomUUID } from "node:crypto"; import { EventEmitter } from "node:events"; import { basename, resolve } from "node:path"; import type { ExecutorRegistry, WorkflowExecutor } from "./executor.js"; import type { PortableAgentRunner } from "./portable-agent-runner.js"; import type { PersistedAgentState, PersistedRunState, RunStatus } from "./run-persistence.js"; import { STANDALONE_PROTOCOL_VERSION } from "./standalone-contract.js"; import { type SchedulableWorkflowManager, UsageLimitScheduler, type UsageLimitSchedulerOptions, } from "./usage-limit-scheduler.js"; import { type ExecOptions, type ManagedRun, WorkflowManager, type WorkflowManagerOptions } from "./workflow-manager.js"; import { createWorkflowStorage, type WorkflowStorage } from "./workflow-saved.js"; const FORWARDED_MANAGER_EVENTS = [ "log", "phase", "runtimeEvent", "agentStart", "agentEnd", "agentHistory", "tokenUsage", "complete", "error", "paused", "resumed", "stopped", ] as const; const SETTLED_STATUSES: ReadonlySet = new Set(["paused", "completed", "failed", "aborted"]); const CHECKPOINT_CLOSING_EVENTS = new Set(["complete", "error", "paused", "stopped"]); export interface StandaloneRunOptions { defaultExecutor?: WorkflowExecutor; maxAgents?: number; concurrency?: number; agentRetries?: number; agentTimeoutMs?: number | null; tokenBudget?: number | null; autoResume?: boolean; } export interface StartStandaloneRunRequest { script: string; args?: unknown; options?: StandaloneRunOptions; } export interface ResumeStandaloneRunRequest { script?: string; args?: unknown; } /** * Host-neutral agent surface accepted by the standalone runtime. * * It intentionally uses the portable runner's structural contract instead of * the optional Pi SDK types. A Pi-backed runner still satisfies this shape. */ export type StandaloneAgentRunner = Pick; /** * Public runtime options are deliberately declared independently from * WorkflowManagerOptions. This keeps `@fish59fish/dynamic-workflows/standalone` * type-checkable when the optional Pi peers are not installed. */ export interface StandaloneRuntimeOptions { cwd?: string; concurrency?: number; loadSavedWorkflow?: (name: string) => string | undefined; agent?: StandaloneAgentRunner; defaultExecutor?: WorkflowExecutor; executorRegistry?: ExecutorRegistry; mainModel?: string; /** Optional host model registry. It is interpreted only by a compatible runner. */ modelRegistry?: unknown; sessionId?: string; defaultAgentTimeoutMs?: number | null; defaultAgentRetries?: number; defaultTokenBudget?: number | null; /** Optional host toolsets. Values stay opaque until a compatible runner uses them. */ toolsets?: Record readonly unknown[]>; excludeSubagentTools?: string[]; persistAgentSessions?: boolean; maxTerminalRunsInMemory?: number; /** Wait for aborted executor cleanup before resume; default true. */ abortCleanupFence?: boolean; /** * Existing storage can be injected by embedders and tests. When omitted, * saved workflows use the same project/user locations as the Pi extension. */ workflowStorage?: WorkflowStorage; /** Set false to disable provider-usage-limit auto-resume. Default enabled. */ usageLimitScheduler?: false | UsageLimitSchedulerOptions; /** * Interactive mode publishes checkpoint requests for a UI/API client. * Headless mode applies each checkpoint's declared default/abort policy. * Default: interactive. */ checkpointMode?: "interactive" | "headless"; } export interface StandaloneCheckpointOptions { default?: unknown; headless?: "default" | "abort"; kind?: "confirm" | "input" | "select"; choices?: string[]; timeoutMs?: number; } export interface StandaloneWorkflowMetaPhase { title: string; detail?: string; model?: string; } export interface StandaloneWorkflowMeta { name: string; description: string; phases?: StandaloneWorkflowMetaPhase[]; model?: string; } export interface StandaloneWorkflowRunResult { meta: StandaloneWorkflowMeta; result: T; logs: string[]; phases: string[]; agentCount: number; durationMs: number; executor?: WorkflowExecutor; runId?: string; tokenUsage?: { input: number; output: number; total: number; cost: number; cacheRead?: number; cacheWrite?: number; }; } export interface StandaloneRuntimeEvent { sequence: number; timestamp: string; type: string; runId?: string; payload: unknown; } export interface PendingCheckpoint { id: string; runId: string; prompt: string; options: StandaloneCheckpointOptions; createdAt: string; } export type StandaloneRunState = PersistedRunState; export interface StandaloneRuntimeState { project: { cwd: string; name: string; }; runtime: { startedAt: string; version: string; protocolVersion: typeof STANDALONE_PROTOCOL_VERSION; pid: number; autoResumeEnabled: boolean; }; runs: StandaloneRunState[]; checkpoints: PendingCheckpoint[]; } export interface StandaloneRunSummary { runId: string; workflowName: string; description?: string; status: RunStatus; pauseReason?: string; resetHint?: string; phases: string[]; currentPhase?: string; startedAt: string; updatedAt: string; completedAt?: string; durationMs?: number; tokenUsage?: PersistedRunState["tokenUsage"]; defaultExecutor?: WorkflowExecutor; error?: PersistedRunState["error"]; agentCount: number; completedAgentCount: number; errorAgentCount: number; } export interface StandaloneRuntimeOverview { project: StandaloneRuntimeState["project"]; runtime: StandaloneRuntimeState["runtime"]; runs: StandaloneRunSummary[]; checkpoints: PendingCheckpoint[]; } interface PendingCheckpointRecord { checkpoint: PendingCheckpoint; resolve(value: unknown): void; reject(error: Error): void; timer?: ReturnType; } /** * Host-neutral facade over WorkflowManager. * * It owns no Pi UI or conversation state: callers receive a durable run model, * a normalized event stream, and explicit control methods suitable for a CLI, * daemon, desktop app, or another agent runtime. */ export class StandaloneWorkflowRuntime extends EventEmitter { readonly cwd: string; readonly startedAt = new Date().toISOString(); private readonly manager: WorkflowManager; private readonly storage: WorkflowStorage; private readonly usageLimitScheduler?: UsageLimitScheduler; private readonly checkpointMode: "interactive" | "headless"; private readonly events: StandaloneRuntimeEvent[] = []; private readonly pendingCheckpoints = new Map(); private sequence = 0; private closed = false; constructor(options: StandaloneRuntimeOptions = {}) { super(); const { workflowStorage, usageLimitScheduler, checkpointMode = "interactive", ...managerOptions } = options; this.cwd = resolve(managerOptions.cwd ?? process.cwd()); this.storage = workflowStorage ?? createWorkflowStorage(this.cwd); this.checkpointMode = checkpointMode; const loadSavedWorkflow = managerOptions.loadSavedWorkflow ?? ((name: string) => this.storage.load(name)?.script); this.manager = new WorkflowManager({ ...(managerOptions as WorkflowManagerOptions), cwd: this.cwd, loadSavedWorkflow, abortCleanupFence: managerOptions.abortCleanupFence ?? true, }); for (const eventName of FORWARDED_MANAGER_EVENTS) { this.manager.on(eventName, (payload: unknown) => { const runId = extractRunId(payload); if (runId && CHECKPOINT_CLOSING_EVENTS.has(eventName)) { this.rejectCheckpointsForRun(runId, `Workflow ${eventName}.`); } this.publish(eventName, payload); }); } if (usageLimitScheduler !== false) { this.usageLimitScheduler = new UsageLimitScheduler(this.schedulerManager(), usageLimitScheduler); } } start(request: StartStandaloneRunRequest): { runId: string; promise: Promise; } { this.assertOpen(); if (!request || typeof request.script !== "string" || !request.script.trim()) { throw new TypeError("A non-empty workflow script is required."); } let runId = ""; const confirm = this.checkpointMode === "interactive" ? this.checkpointResolver(() => runId) : undefined; const exec: ExecOptions = { ...request.options, ...(confirm ? { confirm } : {}) }; const started = this.manager.startInBackground(request.script, request.args, exec); runId = started.runId; this.publish("runStarted", { runId, options: request.options ?? {} }); return started as { runId: string; promise: Promise }; } async resume(runId: string, request: ResumeStandaloneRunRequest = {}): Promise { this.assertOpen(); const confirm = this.checkpointMode === "interactive" ? this.checkpointResolver(() => runId) : undefined; const resumed = await this.manager.resume(runId, { ...request, ...(confirm ? { confirm } : {}), }); if (resumed) this.publish("runResumeRequested", { runId }); return resumed; } pause(runId: string): boolean { this.assertOpen(); const paused = this.manager.pause(runId); if (paused) this.rejectCheckpointsForRun(runId, "Workflow paused."); return paused; } stop(runId: string): boolean { this.assertOpen(); const stopped = this.manager.stop(runId); if (stopped) this.rejectCheckpointsForRun(runId, "Workflow stopped."); return stopped; } delete(runId: string): boolean { this.assertOpen(); this.rejectCheckpointsForRun(runId, "Workflow run deleted."); const deleted = this.manager.deleteRun(runId); if (deleted) this.publish("runDeleted", { runId }); return deleted; } listRuns(): StandaloneRunState[] { return this.manager.listAllRuns().map((run) => this.mergeLiveRun(run)); } getRun(runId: string): StandaloneRunState | null { const persisted = this.manager.listAllRuns().find((run) => run.runId === runId); if (persisted) return this.mergeLiveRun(persisted); const live = this.manager.getRun(runId); return live ? this.liveRunWithoutPersistence(live) : null; } state(): StandaloneRuntimeState { return { project: { cwd: this.cwd, name: basename(this.cwd) || this.cwd }, runtime: { startedAt: this.startedAt, version: "1", protocolVersion: STANDALONE_PROTOCOL_VERSION, pid: process.pid, autoResumeEnabled: Boolean(this.usageLimitScheduler), }, runs: this.listRuns(), checkpoints: this.listCheckpoints(), }; } /** * Lightweight dashboard/list projection. Scripts, args, journals, logs, and * agent payloads stay behind the per-run endpoint instead of being * re-serialized for every SSE-driven refresh. */ overview(): StandaloneRuntimeOverview { const state = this.stateMetadata(); return { ...state, runs: this.manager.listAllRuns().map((persisted) => { const managed = this.manager.getRun(persisted.runId); const run = managed ? this.mergeLiveRun(persisted) : persisted; return { runId: run.runId, workflowName: run.workflowName, description: run.description, status: run.status, pauseReason: run.pauseReason, resetHint: run.resetHint, phases: [...run.phases], currentPhase: run.currentPhase, startedAt: run.startedAt, updatedAt: run.updatedAt, completedAt: run.completedAt, durationMs: run.durationMs, tokenUsage: run.tokenUsage ? { ...run.tokenUsage } : undefined, defaultExecutor: run.defaultExecutor, error: run.error ? { ...run.error } : undefined, agentCount: run.agents.length, completedAgentCount: run.agents.filter((agent) => agent.status === "done" || agent.status === "skipped") .length, errorAgentCount: run.agents.filter((agent) => agent.status === "error").length, }; }), }; } listCheckpoints(): PendingCheckpoint[] { return [...this.pendingCheckpoints.values()] .map(({ checkpoint }) => ({ ...checkpoint, options: { ...checkpoint.options } })) .sort((a, b) => a.createdAt.localeCompare(b.createdAt)); } respondToCheckpoint(id: string, value: unknown): boolean { const record = this.pendingCheckpoints.get(id); if (!record) return false; this.pendingCheckpoints.delete(id); if (record.timer) clearTimeout(record.timer); record.resolve(value); this.publish("checkpointResolved", { runId: record.checkpoint.runId, checkpointId: id }); return true; } recentEvents(sinceSequence = 0): StandaloneRuntimeEvent[] { return this.events.filter((event) => event.sequence > sinceSequence).map((event) => ({ ...event })); } async waitForSettled( runId: string, options: { signal?: AbortSignal; pollIntervalMs?: number } = {}, ): Promise { const first = this.getRun(runId); if (!first) throw new Error(`Unknown workflow run: ${runId}`); if (SETTLED_STATUSES.has(first.status)) return first; const pollIntervalMs = Math.max(25, options.pollIntervalMs ?? 200); return new Promise((resolveWait, rejectWait) => { let timer: ReturnType | undefined; const cleanup = () => { if (timer) clearTimeout(timer); options.signal?.removeEventListener("abort", onAbort); }; const onAbort = () => { cleanup(); rejectWait(options.signal?.reason ?? new Error("Waiting for workflow was aborted.")); }; const poll = () => { if (options.signal?.aborted) { onAbort(); return; } const run = this.getRun(runId); if (!run) { cleanup(); rejectWait(new Error(`Workflow run disappeared: ${runId}`)); return; } if (SETTLED_STATUSES.has(run.status)) { cleanup(); resolveWait(run); return; } // This poll may be the caller's only remaining completion source. timer = setTimeout(poll, pollIntervalMs); }; options.signal?.addEventListener("abort", onAbort, { once: true }); poll(); }); } /** * Gracefully detach the host. Live runs are checkpointed as paused so a * process exit never strands durable state in "running". */ close(): void { if (this.closed) return; this.closed = true; this.usageLimitScheduler?.dispose(); for (const run of this.manager.listAllRuns()) { if (run.status === "running") this.manager.pause(run.runId); } for (const record of this.pendingCheckpoints.values()) { if (record.timer) clearTimeout(record.timer); record.reject(new Error("Standalone workflow runtime closed.")); } this.pendingCheckpoints.clear(); this.publish("runtimeClosed", { cwd: this.cwd }); } private openCheckpoint( runId: string, prompt: string, options: StandaloneCheckpointOptions, resolveCheckpoint: (value: unknown) => void, rejectCheckpoint: (error: Error) => void, ): void { const id = randomUUID(); const checkpoint: PendingCheckpoint = { id, runId, prompt, options: { ...options }, createdAt: new Date().toISOString(), }; const record: PendingCheckpointRecord = { checkpoint, resolve: resolveCheckpoint, reject: rejectCheckpoint, }; if (typeof options.timeoutMs === "number" && Number.isFinite(options.timeoutMs) && options.timeoutMs > 0) { record.timer = setTimeout(() => { if (!this.pendingCheckpoints.delete(id)) return; if (options.headless === "abort") { rejectCheckpoint(new Error(`Checkpoint timed out: ${prompt}`)); } else { resolveCheckpoint(options.default ?? true); } this.publish("checkpointTimedOut", { runId, checkpointId: id }); }, options.timeoutMs); // A checkpoint-only one-shot run may have no other referenced handle. } this.pendingCheckpoints.set(id, record); this.publish("checkpointOpened", checkpoint); } private stateMetadata(): Omit { return { project: { cwd: this.cwd, name: basename(this.cwd) || this.cwd }, runtime: { startedAt: this.startedAt, version: "1", protocolVersion: STANDALONE_PROTOCOL_VERSION, pid: process.pid, autoResumeEnabled: Boolean(this.usageLimitScheduler), }, checkpoints: this.listCheckpoints(), }; } private checkpointResolver(resolveRunId: () => string): NonNullable { return (prompt: string, rawOptions: unknown) => new Promise((resolveCheckpoint, rejectCheckpoint) => { const options = (rawOptions ?? {}) as StandaloneCheckpointOptions; // startInBackground begins executing synchronously until the script's // first await. Queueing registration lets its generated runId land // first; doing the same on resume keeps both paths deterministic. queueMicrotask(() => { if (this.closed) { rejectCheckpoint(new Error("Standalone workflow runtime closed.")); return; } const runId = resolveRunId(); if (!runId) { rejectCheckpoint(new Error("Workflow checkpoint opened before the run was registered.")); return; } if (this.manager.getRun(runId)?.status !== "running") { rejectCheckpoint(new Error("Workflow stopped before its checkpoint could be opened.")); return; } this.openCheckpoint(runId, prompt, options, resolveCheckpoint, rejectCheckpoint); }); }); } /** * Usage-limit recovery must re-enter through this facade so a fresh execution * receives the current host's checkpoint resolver. Calling manager.resume() * directly would silently turn post-resume checkpoints into headless defaults. */ private schedulerManager(): SchedulableWorkflowManager { return { on: (event, listener) => this.manager.on(event, listener), off: (event, listener) => this.manager.off(event, listener), listAllRuns: () => this.manager.listAllRuns(), resume: (runId) => this.resume(runId), getPersistence: () => this.manager.getPersistence(), }; } private rejectCheckpointsForRun(runId: string, message: string): void { for (const [id, record] of this.pendingCheckpoints) { if (record.checkpoint.runId !== runId) continue; this.pendingCheckpoints.delete(id); if (record.timer) clearTimeout(record.timer); record.reject(new Error(message)); } } private mergeLiveRun(persisted: PersistedRunState): StandaloneRunState { const managed = this.manager.getRun(persisted.runId); if (!managed) return hydrateJournalResults(persisted); const tokenUsage = managed.snapshot.tokenUsage ?? persisted.tokenUsage; return { ...persisted, status: managed.status, description: managed.snapshot.description, phases: [...managed.snapshot.phases], currentPhase: managed.snapshot.currentPhase, agents: this.liveAgents(managed), logs: [...managed.snapshot.logs], result: managed.result?.result ?? persisted.result, durationMs: managed.result?.durationMs ?? Date.now() - managed.startedAt.getTime(), tokenUsage: tokenUsage ? { ...tokenUsage } : undefined, defaultExecutor: managed.defaultExecutor, error: managed.error ? { message: managed.error.message, code: managed.error.code, recoverable: managed.error.recoverable, } : undefined, }; } private liveRunWithoutPersistence(managed: ManagedRun): StandaloneRunState { return { runId: managed.runId, workflowName: managed.snapshot.name, description: managed.snapshot.description, script: managed.script, args: managed.args, status: managed.status, phases: [...managed.snapshot.phases], currentPhase: managed.snapshot.currentPhase, agents: this.liveAgents(managed), logs: [...managed.snapshot.logs], result: managed.result?.result, startedAt: managed.startedAt.toISOString(), updatedAt: new Date().toISOString(), durationMs: managed.result?.durationMs ?? Date.now() - managed.startedAt.getTime(), tokenUsage: managed.snapshot.tokenUsage ? { ...managed.snapshot.tokenUsage } : undefined, defaultExecutor: managed.defaultExecutor, error: managed.error ? { message: managed.error.message, code: managed.error.code, recoverable: managed.error.recoverable, } : undefined, }; } private liveAgents(managed: ManagedRun): PersistedAgentState[] { return managed.snapshot.agents.map((agent) => { const timestamps = managed.agentTimestamps.get(agent.id); return { ...agent, startedAt: timestamps?.startedAt, endedAt: timestamps?.endedAt, }; }); } private publish(type: string, payload: unknown): void { const runId = extractRunId(payload); const event: StandaloneRuntimeEvent = { sequence: ++this.sequence, timestamp: new Date().toISOString(), type, runId, payload, }; this.events.push(event); if (this.events.length > 1_000) this.events.splice(0, this.events.length - 1_000); this.emit("event", event); } private assertOpen(): void { if (this.closed) throw new Error("Standalone workflow runtime is closed."); } } function extractRunId(payload: unknown): string | undefined { if (!payload || typeof payload !== "object") return undefined; const runId = (payload as { runId?: unknown }).runId; return typeof runId === "string" ? runId : undefined; } function hydrateJournalResults(persisted: PersistedRunState): StandaloneRunState { if (!persisted.journal?.length) return { ...persisted }; const journal = new Map( persisted.journal.map((entry) => [`${entry.runId ?? persisted.runId}:${entry.index}`, entry.result] as const), ); return { ...persisted, agents: persisted.agents.map((agent) => { if (agent.result !== undefined || !agent.callId) return { ...agent }; const result = journal.get(agent.callId); return result === undefined ? { ...agent } : { ...agent, result }; }), }; }