diff --git a/package-lock.json b/package-lock.json index 0d7decc4..2029860a 100644 --- a/package-lock.json +++ b/package-lock.json @@ -19,6 +19,7 @@ "@openai/codex-sdk": "^0.142.5", "@opencode-ai/sdk": "^1.17.13", "@pierre/diffs": "^1.2.5", + "ajv": "^8.20.0", "better-sqlite3": "^12.10.0", "diff": "^8.0.3", "drizzle-orm": "^0.45.2", @@ -752,7 +753,7 @@ "typebox": "1.1.38" }, "bin": { - "pi-ai": "dist/cli.js" + "pi-ai": "./dist/cli.js" }, "engines": { "node": ">=22.19.0" @@ -1057,7 +1058,7 @@ } }, "node_modules/@earendil-works/pi-coding-agent/node_modules/@protobufjs/float": { - "version": "1.0.3", + "version": "1.0.2", "resolved": "https://registry.npmjs.org/@protobufjs/float/-/float-1.0.2.tgz", "integrity": "sha512-Ddb+kVXlXst9d+R9PfTIxh1EdNkgoRe5tOX6t01f1lYWOvJnSPDBlG241QLzcyPdoNTsblLUdujGSE4RzrTZGQ==", "license": "BSD-3-Clause" diff --git a/package.json b/package.json index d5c880d0..f43b5dff 100644 --- a/package.json +++ b/package.json @@ -28,7 +28,7 @@ "dev": "node scripts/dev-server.mjs", "postinstall": "node scripts/fix-node-pty-permissions.mjs", "start": "node dist/cli.js serve", - "test": "tsx src/config.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/review-checkpoints.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts && tsx src/workflow-types.test.ts && tsx src/workflow-store.test.ts && tsx src/workflow-script.test.ts && tsx src/workflow-sandbox.test.ts && tsx src/workflow-engine.test.ts", + "test": "tsx src/config.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/review-checkpoints.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts && tsx src/workflow-types.test.ts && tsx src/workflow-store.test.ts && tsx src/workflow-script.test.ts && tsx src/workflow-sandbox.test.ts && tsx src/workflow-engine.test.ts && tsx src/workflow-files.test.ts && tsx src/workflow-replay.test.ts && tsx src/workflow-schema.test.ts", "typecheck": "tsc -p tsconfig.json --noEmit" }, "keywords": [], @@ -44,6 +44,7 @@ "@openai/codex-sdk": "^0.142.5", "@opencode-ai/sdk": "^1.17.13", "@pierre/diffs": "^1.2.5", + "ajv": "^8.20.0", "better-sqlite3": "^12.10.0", "diff": "^8.0.3", "drizzle-orm": "^0.45.2", diff --git a/src/cli.ts b/src/cli.ts index 88dffa90..9cc41360 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -40,7 +40,9 @@ import { import { expandHomePath } from "./roots.js"; import { shutdownHttpServer } from "./server-shutdown.js"; -type Command = "serve" | "init" | "doctor" | "config" | "agents" | "help" | "version"; +import { runWorkflowCommand } from "./workflow-cli.js"; + +type Command = "serve" | "init" | "doctor" | "config" | "agents" | "workflow" | "help" | "version"; const require = createRequire(import.meta.url); const SUPPORTED_NODE_RANGE = ">=20.12 <27"; @@ -67,6 +69,9 @@ async function main(argv: string[]): Promise { case "agents": await runAgentsCommand(args); return; + case "workflow": + await runWorkflowCommand(args, loadConfig()); + return; case "help": printHelp(); return; @@ -78,7 +83,15 @@ async function main(argv: string[]): Promise { function normalizeCommand(command: string | undefined): Command { if (!command || command === "serve" || command === "start") return "serve"; - if (command === "init" || command === "doctor" || command === "config" || command === "agents") return command; + if ( + command === "init" || + command === "doctor" || + command === "config" || + command === "agents" || + command === "workflow" + ) { + return command; + } if (command === "help" || command === "--help" || command === "-h") return "help"; if (command === "version" || command === "--version" || command === "-v") return "version"; throw new Error(`Unknown command: ${command}`); @@ -312,6 +325,7 @@ function printHelp(): void { " devspace agents ls List subagent sessions", " devspace agents run [--model ] ", " devspace agents show ", + " devspace workflow run|status|cancel|ls", " devspace -v, --version Print the installed version", "", "For temporary tunnels:", diff --git a/src/workflow-api.ts b/src/workflow-api.ts index 05249df1..ea093a0f 100644 --- a/src/workflow-api.ts +++ b/src/workflow-api.ts @@ -325,7 +325,7 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { }); const cwd = worktreePath ?? deps.workspaceRoot; - const result = await deps.runProvider({ + const providerBase = { provider, prompt, model: agentOpts.model, @@ -334,25 +334,42 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { signal: deps.signal, label: agentOpts.label, phase, - }); - - throwIfCancelled(deps); + }; - let returnValue: unknown = result.finalResponse; + let returnValue: unknown; let structuredJson: string | undefined; + let result: WorkflowProviderRunResult; + if (agentOpts.schema) { - // Full Ajv enforcement lands in M6; for now extract JSON object when schema set. - const extracted = tryExtractJson(result.finalResponse); - if (extracted === undefined) { - throw new WorkflowEngineError( - "schema", - "agent() schema set but response was not valid JSON", - ); - } - returnValue = extracted; - structuredJson = JSON.stringify(extracted); + // Lazy import keeps non-schema paths free of ajv load cost. + const { enforceAgentSchema } = await import("./workflow-schema.js"); + const enforced = await enforceAgentSchema({ + schema: agentOpts.schema, + prompt, + run: (p) => deps.runProvider({ ...providerBase, prompt: p }), + onRetry: ({ attempt, errors }) => { + deps.journal.appendEvent({ + runId: deps.runId, + type: "schema_retry", + phase, + label: agentOpts.label, + data: { callIndex: index, attempt, errors }, + }); + }, + }); + returnValue = enforced.value; + structuredJson = JSON.stringify(enforced.value); + result = { + finalResponse: enforced.finalResponse, + providerSessionId: enforced.providerSessionId, + }; + } else { + result = await deps.runProvider(providerBase); + returnValue = result.finalResponse; } + throwIfCancelled(deps); + let dirty: boolean | undefined; if (worktree) { const finalized = await worktree.finalize("success"); diff --git a/src/workflow-cli.ts b/src/workflow-cli.ts new file mode 100644 index 00000000..620cd50c --- /dev/null +++ b/src/workflow-cli.ts @@ -0,0 +1,536 @@ +import { spawn } from "node:child_process"; +import { readFile } from "node:fs/promises"; +import { availableParallelism } from "node:os"; +import { resolve } from "node:path"; +import { fileURLToPath } from "node:url"; +import type { ServerConfig } from "./config.js"; +import { runLocalAgentProvider } from "./local-agent-adapters.js"; +import { getLocalAgentProviderAvailabilitySnapshot } from "./local-agent-availability.js"; +import { + isLocalAgentProvider, + LOCAL_AGENT_PROVIDERS, +} from "./local-agent-profiles.js"; +import { executeWorkflow, mapEngineErrorKind } from "./workflow-engine.js"; +import { + parseWorkflowArgFlags, + persistWorkflowScript, + readWorkflowScriptFile, + resolveNamedWorkflowScript, + resolveWorkflowScriptFromPathOrName, +} from "./workflow-files.js"; +import { createWorkflowReplay } from "./workflow-replay.js"; +import { parseWorkflowScript } from "./workflow-script.js"; +import { createWorkflowStore, type WorkflowStore } from "./workflow-store.js"; +import { + WORKFLOW_CANCEL_HARD_MS, + WORKFLOW_HEARTBEAT_MS, + WORKFLOW_LIMITS, + resolveWorkflowConcurrency, + type WorkflowEventRecord, + type WorkflowRunRecord, + type WorkflowRunSource, +} from "./workflow-types.js"; +import { + createWorkflowWorktreeFactory, + resolveWorkspaceHead, +} from "./workflow-worktrees.js"; + +export async function runWorkflowCommand( + args: string[], + config: ServerConfig, +): Promise { + const [subcommand, ...rest] = args; + switch (subcommand) { + case "run": + await runWorkflowRun(rest, config); + return; + case "status": + await runWorkflowStatus(rest, config); + return; + case "cancel": + await runWorkflowCancel(rest, config); + return; + case "ls": + case "list": + await runWorkflowList(config); + return; + case "__worker": + await runWorkflowWorker(rest, config); + return; + case undefined: + case "help": + case "--help": + case "-h": + printWorkflowHelp(); + return; + default: + throw new Error(`Unknown workflow command: ${subcommand}`); + } +} + +export function printWorkflowHelp(): void { + console.log( + [ + "DevSpace workflows", + "", + "Usage:", + " devspace workflow run (--file | --name | --resume )", + " [--arg key=value]... [--follow]", + " devspace workflow status [--follow]", + " devspace workflow cancel ", + " devspace workflow ls", + ].join("\n"), + ); +} + +async function runWorkflowRun(args: string[], config: ServerConfig): Promise { + const { flags } = splitFlags(args); + const follow = flags.has("follow"); + const file = flagValue(flags, "file"); + const name = flagValue(flags, "name"); + const resumeFrom = flagValue(flags, "resume"); + const { args: workflowArgs } = parseWorkflowArgFlags(collectArgTokens(args)); + + if (!file && !name && !resumeFrom) { + throw new Error( + "Usage: devspace workflow run (--file | --name | --resume )", + ); + } + + const store = createWorkflowStore(config); + try { + const workspaceRoot = resolve(process.env.DEVSPACE_WORKSPACE_ROOT || process.cwd()); + let source: string; + let scriptHash: string; + let nameHint: string; + let runSource: WorkflowRunSource = "inline"; + let priorRunId: string | undefined; + let priorScriptPath: string | undefined; + + if (resumeFrom) { + const prior = store.getRun(resumeFrom); + if (!prior) throw new Error(`Unknown workflow run to resume: ${resumeFrom}`); + priorRunId = prior.id; + priorScriptPath = prior.scriptPath; + const resolved = await readWorkflowScriptFile(prior.scriptPath); + source = resolved.source; + scriptHash = prior.scriptHash; + nameHint = prior.name; + runSource = "resume"; + if (!Object.keys(workflowArgs).length && prior.argsJson && prior.argsJson !== "null") { + try { + Object.assign(workflowArgs, JSON.parse(prior.argsJson) as object); + } catch { + // keep empty + } + } + } else { + const resolved = await resolveWorkflowScriptFromPathOrName({ + file, + name, + workspaceRoot, + stateDir: config.stateDir, + }); + source = resolved.source; + scriptHash = resolved.scriptHash; + nameHint = resolved.nameHint; + runSource = resolved.origin === "named" ? "named" : "inline"; + } + + const parsed = parseWorkflowScript(source, { + filename: priorScriptPath ?? file ?? name ?? "workflow:inline", + }); + const baseSha = await resolveWorkspaceHead(workspaceRoot); + + const run = store.createRun({ + name: parsed.meta.name || nameHint, + source: runSource, + scriptPath: priorScriptPath ?? "pending", + scriptHash, + workspaceRoot, + workspaceId: process.env.DEVSPACE_WORKSPACE_ID, + argsJson: JSON.stringify(Object.keys(workflowArgs).length ? workflowArgs : null), + resumedFromRunId: priorRunId, + baseSha, + }); + + const persisted = + priorScriptPath ?? + (await persistWorkflowScript({ + stateDir: config.stateDir, + runId: run.id, + source, + preferredName: parsed.meta.name || nameHint, + })); + if (!priorScriptPath) { + store.setScriptPath(run.id, persisted); + } + + spawnWorkflowWorkerFromCli( + run.id, + fileURLToPath(import.meta.url.replace(/workflow-cli\.(ts|js)$/, "cli.$1")), + ); + + console.log(formatRunLine(store.getRun(run.id) ?? { ...run, scriptPath: persisted })); + + if (follow) { + await followRun(store, run.id); + } + } finally { + store.close(); + } +} + +async function runWorkflowStatus(args: string[], config: ServerConfig): Promise { + const follow = args.includes("--follow"); + const runId = args.find((a) => !a.startsWith("-")); + if (!runId) throw new Error("Usage: devspace workflow status [--follow]"); + + const store = createWorkflowStore(config); + try { + const run = store.getRun(runId); + if (!run) throw new Error(`Unknown workflow run: ${runId}`); + console.log(formatRunLine(run)); + if (follow) { + await followRun(store, runId); + return; + } + if (run.resultJson) console.log(run.resultJson); + else if (run.error) console.log(run.error); + } finally { + store.close(); + } +} + +async function runWorkflowCancel(args: string[], config: ServerConfig): Promise { + const runId = args[0]; + if (!runId) throw new Error("Usage: devspace workflow cancel "); + const store = createWorkflowStore(config); + try { + const run = store.requestCancel(runId); + console.log(formatRunLine(run)); + if (run.pid && (run.status === "running" || run.status === "starting")) { + try { + process.kill(run.pid, "SIGTERM"); + } catch { + // already dead + } + await sleep(WORKFLOW_CANCEL_HARD_MS); + const again = store.getRun(runId); + if (again && (again.status === "running" || again.status === "starting") && again.pid) { + try { + process.kill(-again.pid, "SIGKILL"); + } catch { + try { + process.kill(again.pid, "SIGKILL"); + } catch { + // gone + } + } + const latest = store.getRun(runId); + if (latest && (latest.status === "running" || latest.status === "starting")) { + store.cancelRun(runId, "cancelled (hard kill)"); + } + } + } + console.log(formatRunLine(store.getRun(runId)!)); + } finally { + store.close(); + } +} + +async function runWorkflowList(config: ServerConfig): Promise { + const store = createWorkflowStore(config); + try { + const runs = store.listRuns(50); + if (runs.length === 0) { + console.log("No workflow runs."); + return; + } + for (const run of runs) console.log(formatRunLine(run)); + } finally { + store.close(); + } +} + +/** Detached worker entry: claim run, heartbeat, execute, complete/fail. */ +export async function runWorkflowWorker( + args: string[], + config: ServerConfig, +): Promise { + const runId = args[0]; + if (!runId) throw new Error("Usage: devspace workflow __worker "); + + const store = createWorkflowStore(config); + const claimed = store.claimRun(runId, process.pid); + if (!claimed) { + store.close(); + throw new Error(`Cannot claim workflow run ${runId} (missing or not starting)`); + } + + const abort = new AbortController(); + const heartbeat = setInterval(() => { + try { + store.setHeartbeat(runId); + if (store.isCancelRequested(runId)) abort.abort(); + } catch { + // store closed + } + }, WORKFLOW_HEARTBEAT_MS); + + try { + const source = await readFile(claimed.scriptPath, "utf8"); + const parsed = parseWorkflowScript(source, { filename: claimed.scriptPath }); + const enabledProviders = resolveEnabledProviders(); + const concurrency = resolveWorkflowConcurrency( + parsed.meta.concurrency, + availableParallelism(), + ); + + let argsValue: unknown; + try { + argsValue = JSON.parse(claimed.argsJson); + if (argsValue === null) argsValue = undefined; + } catch { + argsValue = undefined; + } + + const replay = claimed.resumedFromRunId + ? createWorkflowReplay(store.listAgentCalls(claimed.resumedFromRunId)) + : undefined; + + const createWorktree = createWorkflowWorktreeFactory({ + worktreeRoot: config.worktreeRoot, + allowedRoots: config.allowedRoots, + }); + + const { result, callCount } = await executeWorkflow({ + parsed, + runId, + journal: store, + args: argsValue, + concurrency, + signal: abort.signal, + workspaceRoot: claimed.workspaceRoot, + baseSha: claimed.baseSha, + enabledProviders, + createWorktree, + replay, + runProvider: async (input) => { + if (!isLocalAgentProvider(input.provider)) { + throw new Error(`Unknown provider: ${input.provider}`); + } + if (abort.signal.aborted || store.isCancelRequested(runId)) { + throw Object.assign(new Error("Workflow cancelled"), { name: "AbortError" }); + } + const providerResult = await runLocalAgentProvider(input.provider, { + prompt: input.prompt, + workspace: input.workspace, + model: input.model, + effort: input.effort, + writeMode: "allowed", + }); + return { + finalResponse: providerResult.finalResponse, + providerSessionId: providerResult.providerSessionId ?? undefined, + }; + }, + resolveNestedSource: async (ref) => { + if (typeof ref === "string") { + const named = await resolveNamedWorkflowScript({ + name: ref, + workspaceRoot: claimed.workspaceRoot, + stateDir: config.stateDir, + }); + return named.source; + } + return readFile(ref.scriptPath, "utf8"); + }, + }); + + if (abort.signal.aborted || store.isCancelRequested(runId)) { + store.cancelRun(runId); + return; + } + + let resultJson: string | undefined; + if (result !== undefined) { + resultJson = JSON.stringify(result); + if (Buffer.byteLength(resultJson, "utf8") > WORKFLOW_LIMITS.resultJsonBytes) { + store.failRun(runId, { + error: `result exceeds ${WORKFLOW_LIMITS.resultJsonBytes} bytes`, + errorKind: "result_too_large", + }); + return; + } + } + + store.completeRun(runId, { resultJson }); + store.appendEvent({ + runId, + type: "run_completed", + data: { callCount }, + }); + } catch (error) { + if (store.isCancelRequested(runId) || abort.signal.aborted) { + try { + store.cancelRun(runId); + } catch { + // already terminal + } + return; + } + const message = error instanceof Error ? error.message : String(error); + const errorKind = mapEngineErrorKind(error); + try { + store.failRun(runId, { error: message, errorKind }); + store.appendEvent({ + runId, + type: "run_failed", + data: { error: message, errorKind }, + }); + } catch { + // terminal race + } + } finally { + clearInterval(heartbeat); + store.close(); + } +} + +export function spawnWorkflowWorkerFromCli(runId: string, cliEntry: string): void { + const child = spawn( + process.execPath, + [...process.execArgv, cliEntry, "workflow", "__worker", runId], + { + detached: true, + stdio: "ignore", + env: process.env, + }, + ); + child.unref(); +} + +async function followRun(store: WorkflowStore, runId: string): Promise { + let sinceSeq = 0; + for (;;) { + const page = store.drainEvents(runId, sinceSeq, WORKFLOW_LIMITS.eventDrainDefault); + for (const event of page.events) printEvent(event); + sinceSeq = page.nextSeq; + if (page.terminal) { + const run = page.run; + if (run.resultJson) console.log(run.resultJson); + else if (run.error) console.log(run.error); + return; + } + await sleep(300); + } +} + +function printEvent(event: WorkflowEventRecord): void { + const prefix = event.phase ? `[${event.phase}] ` : ""; + switch (event.type) { + case "log": { + let message = event.dataJson; + try { + message = String( + (JSON.parse(event.dataJson) as { message?: string }).message ?? event.dataJson, + ); + } catch { + // raw + } + console.log(`${prefix}${message}`); + break; + } + case "phase_started": + console.log(`== phase ${event.phase ?? ""} ==`); + break; + case "agent_call_started": + console.log(`${prefix}agent start ${event.label ?? ""}`.trim()); + break; + case "agent_call_completed": + console.log(`${prefix}agent done ${event.label ?? ""}`.trim()); + break; + case "agent_call_cached": + console.log(`${prefix}agent cache ${event.label ?? ""}`.trim()); + break; + case "agent_call_failed": + console.log(`${prefix}agent fail ${event.label ?? ""} ${event.dataJson}`.trim()); + break; + case "run_completed": + case "run_failed": + case "run_cancelled": + console.log(event.type); + break; + default: + break; + } +} + +function formatRunLine( + run: Pick, +): string { + const err = run.error ? ` error=${JSON.stringify(run.error)}` : ""; + return `${run.id} ${run.status} ${run.name}${err}`; +} + +function resolveEnabledProviders(): string[] { + const snapshot = getLocalAgentProviderAvailabilitySnapshot(); + const live = new Set(snapshot.filter((row) => row.available).map((row) => row.name)); + return LOCAL_AGENT_PROVIDERS.filter((id) => live.has(id)); +} + +function splitFlags(args: string[]): { + flags: Map; + positionals: string[]; +} { + const flags = new Map(); + const positionals: string[] = []; + for (let i = 0; i < args.length; i += 1) { + const token = args[i]!; + if (token === "--") { + positionals.push(...args.slice(i + 1)); + break; + } + if (token.startsWith("--")) { + const eq = token.indexOf("="); + if (eq >= 0) { + flags.set(token.slice(2, eq), token.slice(eq + 1)); + continue; + } + const key = token.slice(2); + const next = args[i + 1]; + if (next && !next.startsWith("-") && key !== "follow") { + flags.set(key, next); + i += 1; + } else { + flags.set(key, true); + } + continue; + } + positionals.push(token); + } + return { flags, positionals }; +} + +function flagValue(flags: Map, key: string): string | undefined { + const value = flags.get(key); + return typeof value === "string" ? value : undefined; +} + +function collectArgTokens(args: string[]): string[] { + const out: string[] = []; + for (let i = 0; i < args.length; i += 1) { + const token = args[i]!; + if (token === "--arg") { + out.push(token, args[++i] ?? ""); + continue; + } + if (token.startsWith("--arg=")) out.push(token); + } + return out; +} + +function sleep(ms: number): Promise { + return new Promise((resolveSleep) => setTimeout(resolveSleep, ms)); +} diff --git a/src/workflow-files.test.ts b/src/workflow-files.test.ts new file mode 100644 index 00000000..40fd4e6b --- /dev/null +++ b/src/workflow-files.test.ts @@ -0,0 +1,64 @@ +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + parseWorkflowArgFlags, + persistWorkflowScript, + resolveNamedWorkflowScript, + resolveWorkflowScriptFromPathOrName, + WorkflowPathError, +} from "./workflow-files.js"; +import { hashSource } from "./workflow-script.js"; + +{ + const { args, rest } = parseWorkflowArgFlags([ + "--arg", + "n=1", + "--arg", + 'files=["a.ts"]', + "--follow", + "extra", + ]); + assert.deepEqual(args, { n: 1, files: ["a.ts"] }); + assert.deepEqual(rest, ["--follow", "extra"]); +} + +{ + const dir = await mkdtemp(join(tmpdir(), "wf-files-")); + const path = await persistWorkflowScript({ + stateDir: dir, + runId: "wfr_test", + source: "export const meta = { name: 'x', description: 'd' }\nreturn 1\n", + preferredName: "demo", + }); + assert.match(path, /workflow-scripts\/wfr_test\/demo\.js$/); + + const file = await resolveWorkflowScriptFromPathOrName({ + file: path, + workspaceRoot: dir, + }); + assert.equal(file.origin, "file"); + assert.equal(file.scriptHash, hashSource(file.source)); + + await mkdir(join(dir, ".devspace", "workflows"), { recursive: true }); + await writeFile( + join(dir, ".devspace", "workflows", "named.js"), + "export const meta = { name: 'named', description: 'd' }\nreturn 2\n", + ); + const named = await resolveNamedWorkflowScript({ + name: "named", + workspaceRoot: dir, + }); + assert.equal(named.origin, "named"); + assert.match(named.source, /named/); + + await assert.rejects( + () => resolveNamedWorkflowScript({ name: "missing", workspaceRoot: dir }), + WorkflowPathError, + ); + + await rm(dir, { recursive: true, force: true }); +} + +console.log("workflow-files.test.ts: ok"); diff --git a/src/workflow-files.ts b/src/workflow-files.ts new file mode 100644 index 00000000..12b12228 --- /dev/null +++ b/src/workflow-files.ts @@ -0,0 +1,186 @@ +import { createHash, randomBytes } from "node:crypto"; +import { access, mkdir, readFile, writeFile } from "node:fs/promises"; +import { basename, dirname, extname, isAbsolute, join, resolve } from "node:path"; +import { hashSource } from "./workflow-script.js"; + +export class WorkflowPathError extends Error { + constructor(message: string) { + super(message); + this.name = "WorkflowPathError"; + } +} + +export interface ResolvedWorkflowScript { + source: string; + scriptPath: string; + scriptHash: string; + nameHint: string; + origin: "file" | "named" | "inline" | "resume"; +} + +/** + * Persist script under stateDir for worker re-read / audit. + * Returns absolute path written. + */ +export async function persistWorkflowScript(input: { + stateDir: string; + runId: string; + source: string; + preferredName?: string; +}): Promise { + const dir = join(input.stateDir, "workflow-scripts", input.runId); + await mkdir(dir, { recursive: true }); + const base = + sanitizeSegment(input.preferredName ?? "script") || + `script-${randomBytes(3).toString("hex")}`; + const path = join(dir, `${base}.js`); + await writeFile(path, input.source, { encoding: "utf8", mode: 0o600 }); + return path; +} + +export async function readWorkflowScriptFile(path: string): Promise { + const scriptPath = resolve(path); + await assertReadableFile(scriptPath); + const source = await readFile(scriptPath, "utf8"); + return { + source, + scriptPath, + scriptHash: hashSource(source), + nameHint: basename(scriptPath, extname(scriptPath)), + origin: "file", + }; +} + +/** + * Resolve named workflow script. + * Search order: + * 1. `/.devspace/workflows/.js` + * 2. `/workflows/.js` + * 3. `/workflows/.js` (if stateDir provided) + */ +export async function resolveNamedWorkflowScript(input: { + name: string; + workspaceRoot: string; + stateDir?: string; +}): Promise { + const name = input.name.trim(); + if (!name || name.includes("/") || name.includes("\\") || name.includes("..")) { + throw new WorkflowPathError(`Invalid workflow name: ${JSON.stringify(input.name)}`); + } + const candidates = [ + join(input.workspaceRoot, ".devspace", "workflows", `${name}.js`), + join(input.workspaceRoot, "workflows", `${name}.js`), + ]; + if (input.stateDir) { + candidates.push(join(input.stateDir, "workflows", `${name}.js`)); + } + for (const candidate of candidates) { + try { + await assertReadableFile(candidate); + const source = await readFile(candidate, "utf8"); + return { + source, + scriptPath: candidate, + scriptHash: hashSource(source), + nameHint: name, + origin: "named", + }; + } catch { + // try next + } + } + throw new WorkflowPathError( + `Named workflow not found: ${name}. Looked in ${candidates.join(", ")}`, + ); +} + +export async function resolveWorkflowScriptFromPathOrName(input: { + file?: string; + name?: string; + workspaceRoot: string; + stateDir?: string; +}): Promise { + if (input.file && input.name) { + throw new WorkflowPathError("Pass only one of --file or --name"); + } + if (input.file) { + const path = isAbsolute(input.file) + ? input.file + : resolve(input.workspaceRoot, input.file); + return readWorkflowScriptFile(path); + } + if (input.name) { + return resolveNamedWorkflowScript({ + name: input.name, + workspaceRoot: input.workspaceRoot, + stateDir: input.stateDir, + }); + } + throw new WorkflowPathError("Provide --file or --name "); +} + +export function parseWorkflowArgFlags(tokens: string[]): { + args: Record; + rest: string[]; +} { + const args: Record = {}; + const rest: string[] = []; + for (let i = 0; i < tokens.length; i += 1) { + const token = tokens[i]!; + if (token === "--arg") { + const pair = tokens[++i]; + if (!pair || !pair.includes("=")) { + throw new WorkflowPathError("--arg requires key=value"); + } + const eq = pair.indexOf("="); + const key = pair.slice(0, eq); + const raw = pair.slice(eq + 1); + args[key] = coerceArgValue(raw); + continue; + } + if (token.startsWith("--arg=")) { + const pair = token.slice("--arg=".length); + const eq = pair.indexOf("="); + if (eq < 0) throw new WorkflowPathError("--arg requires key=value"); + args[pair.slice(0, eq)] = coerceArgValue(pair.slice(eq + 1)); + continue; + } + rest.push(token); + } + return { args, rest }; +} + +function coerceArgValue(raw: string): unknown { + try { + return JSON.parse(raw); + } catch { + return raw; + } +} + +async function assertReadableFile(path: string): Promise { + try { + await access(path); + } catch { + throw new WorkflowPathError(`Script file not found: ${path}`); + } +} + +function sanitizeSegment(value: string): string { + return value + .replace(/[^a-zA-Z0-9._-]+/g, "-") + .replace(/^-+|-+$/g, "") + .slice(0, 80); +} + +export function workflowScriptDirForRun(stateDir: string, runId: string): string { + return join(stateDir, "workflow-scripts", runId); +} + +export function contentHash(source: string): string { + return createHash("sha256").update(source).digest("hex"); +} + +export function dirnameOf(path: string): string { + return dirname(path); +} diff --git a/src/workflow-replay.test.ts b/src/workflow-replay.test.ts new file mode 100644 index 00000000..5258809e --- /dev/null +++ b/src/workflow-replay.test.ts @@ -0,0 +1,55 @@ +import assert from "node:assert/strict"; +import { createWorkflowReplay } from "./workflow-replay.js"; +import type { WorkflowAgentCallRecord } from "./workflow-types.js"; + +function call( + partial: Partial & + Pick, +): WorkflowAgentCallRecord { + return { + runId: "wfr_prior", + provider: "codex", + status: "completed", + fromCache: false, + isolation: "shared", + createdAt: "t", + updatedAt: "t", + ...partial, + }; +} + +{ + const replay = createWorkflowReplay([ + call({ callIndex: 0, cacheKey: "k0", responseText: "a" }), + call({ callIndex: 1, cacheKey: "k1", responseText: "b" }), + ]); + assert.equal(replay.match(0, "k0")?.value, "a"); + assert.equal(replay.match(1, "k1")?.value, "b"); + assert.equal(replay.match(2, "k0"), null); +} + +{ + // fan-out reorder: callIndex mismatch, consume-once by key + const replay = createWorkflowReplay([ + call({ callIndex: 0, cacheKey: "ka", responseText: "A" }), + call({ callIndex: 1, cacheKey: "kb", responseText: "B" }), + ]); + // new run asks index0 for kb first + assert.equal(replay.match(0, "kb")?.value, "B"); + assert.equal(replay.match(1, "ka")?.value, "A"); + assert.equal(replay.match(2, "ka"), null); +} + +{ + const replay = createWorkflowReplay([ + call({ + callIndex: 0, + cacheKey: "ks", + responseText: '{"ok":true}', + structuredJson: '{"ok":true}', + }), + ]); + assert.deepEqual(replay.match(0, "ks")?.value, { ok: true }); +} + +console.log("workflow-replay.test.ts: ok"); diff --git a/src/workflow-replay.ts b/src/workflow-replay.ts new file mode 100644 index 00000000..bff7b935 --- /dev/null +++ b/src/workflow-replay.ts @@ -0,0 +1,81 @@ +import type { WorkflowAgentCallRecord } from "./workflow-types.js"; +import type { WorkflowReplay, WorkflowReplayHit } from "./workflow-api.js"; + +/** + * Resume matcher: + * 1. Prefer same callIndex + cacheKey + * 2. On first miss for an index, fall back to consume-once by cacheKey + * (handles fan-out reordering vs prior run). + */ +export function createWorkflowReplay( + priorCalls: WorkflowAgentCallRecord[], +): WorkflowReplay { + const byIndex = new Map(); + const byKeyQueue = new Map(); + + for (const call of priorCalls) { + if (call.status !== "completed" && call.status !== "from_cache") continue; + byIndex.set(call.callIndex, call); + const queue = byKeyQueue.get(call.cacheKey) ?? []; + queue.push(call); + byKeyQueue.set(call.cacheKey, queue); + } + + const consumed = new Set(); // `${callIndex}` of prior rows consumed + + return { + match(callIndex: number, cacheKey: string): WorkflowReplayHit | null { + const exact = byIndex.get(callIndex); + if (exact && exact.cacheKey === cacheKey && !consumed.has(indexKey(exact))) { + consumed.add(indexKey(exact)); + removeFromKeyQueue(byKeyQueue, exact); + return toHit(exact); + } + + const queue = byKeyQueue.get(cacheKey); + if (!queue || queue.length === 0) return null; + const next = queue.shift()!; + consumed.add(indexKey(next)); + if (queue.length === 0) byKeyQueue.delete(cacheKey); + return toHit(next); + }, + }; +} + +function indexKey(call: WorkflowAgentCallRecord): string { + return `${call.runId}:${call.callIndex}`; +} + +function removeFromKeyQueue( + map: Map, + call: WorkflowAgentCallRecord, +): void { + const queue = map.get(call.cacheKey); + if (!queue) return; + const idx = queue.findIndex( + (row) => row.runId === call.runId && row.callIndex === call.callIndex, + ); + if (idx >= 0) queue.splice(idx, 1); + if (queue.length === 0) map.delete(call.cacheKey); +} + +function toHit(call: WorkflowAgentCallRecord): WorkflowReplayHit { + if (call.structuredJson) { + try { + return { + value: JSON.parse(call.structuredJson), + responseText: call.responseText, + structuredJson: call.structuredJson, + providerSessionId: call.providerSessionId, + }; + } catch { + // fall through to text + } + } + return { + value: call.responseText ?? "", + responseText: call.responseText, + structuredJson: call.structuredJson, + providerSessionId: call.providerSessionId, + }; +} diff --git a/src/workflow-schema.test.ts b/src/workflow-schema.test.ts new file mode 100644 index 00000000..67ad58f4 --- /dev/null +++ b/src/workflow-schema.test.ts @@ -0,0 +1,59 @@ +import assert from "node:assert/strict"; +import { + augmentPromptForSchema, + enforceAgentSchema, + formatAjvErrors, +} from "./workflow-schema.js"; +import { WorkflowEngineError } from "./workflow-api.js"; + +{ + const prompt = augmentPromptForSchema("find bugs", { + type: "object", + properties: { n: { type: "number" } }, + required: ["n"], + }); + assert.match(prompt, /ONLY a JSON/); + assert.match(prompt, /"n"/); +} + +assert.equal( + formatAjvErrors([{ instancePath: "/n", message: "must be number" }]), + "/n must be number", +); + +{ + let attempts = 0; + const result = await enforceAgentSchema({ + schema: { + type: "object", + properties: { n: { type: "number" } }, + required: ["n"], + additionalProperties: false, + }, + prompt: "give n", + run: async () => { + attempts += 1; + if (attempts === 1) return { finalResponse: '{"n":"x"}' }; + return { finalResponse: '{"n":2}', providerSessionId: "sess" }; + }, + }); + assert.deepEqual(result.value, { n: 2 }); + assert.equal(result.attempts, 2); + assert.equal(result.providerSessionId, "sess"); +} + +{ + await assert.rejects( + () => + enforceAgentSchema({ + schema: { type: "object", properties: { n: { type: "number" } }, required: ["n"] }, + prompt: "x", + maxRetries: 1, + run: async () => ({ finalResponse: "not json" }), + }), + (error: unknown) => + error instanceof WorkflowEngineError && error.kind === "schema", + ); +} + +console.log("workflow-schema.test.ts: ok"); diff --git a/src/workflow-schema.ts b/src/workflow-schema.ts new file mode 100644 index 00000000..476e97b6 --- /dev/null +++ b/src/workflow-schema.ts @@ -0,0 +1,132 @@ +import { createRequire } from "node:module"; +import { WORKFLOW_MAX_SCHEMA_RETRIES } from "./workflow-types.js"; +import { tryExtractJson, WorkflowEngineError } from "./workflow-api.js"; +import type { WorkflowProviderRunResult, WorkflowRunProvider } from "./workflow-api.js"; + +const require = createRequire(import.meta.url); + +type AjvLike = new (opts?: object) => { + compile: (schema: object) => ((data: unknown) => boolean) & { + errors?: Array<{ instancePath?: string; message?: string }> | null; + }; +}; + +function loadAjv(): AjvLike { + // Prefer direct package; fall back to transitive install under zod or package-lock. + try { + return require("ajv").default ?? require("ajv"); + } catch { + throw new WorkflowEngineError( + "schema", + "ajv is required for opts.schema (add dependency ajv)", + ); + } +} + +export interface EnforceSchemaInput { + schema: object; + prompt: string; + run: (prompt: string) => Promise; + onRetry?: (info: { attempt: number; errors: string }) => void; + maxRetries?: number; +} + +export interface EnforceSchemaResult { + value: unknown; + finalResponse: string; + providerSessionId?: string; + attempts: number; +} + +/** + * Augment prompt → run → extract JSON → Ajv validate → retry ≤2. + */ +export async function enforceAgentSchema( + input: EnforceSchemaInput, +): Promise { + const Ajv = loadAjv(); + const ajv = new Ajv({ allErrors: true, strict: false }); + const validate = ajv.compile(input.schema); + const maxRetries = input.maxRetries ?? WORKFLOW_MAX_SCHEMA_RETRIES; + const basePrompt = augmentPromptForSchema(input.prompt, input.schema); + + let lastResponse = ""; + let lastSession: string | undefined; + let lastErrors = "unknown validation error"; + + for (let attempt = 0; attempt <= maxRetries; attempt += 1) { + const prompt = + attempt === 0 + ? basePrompt + : `${basePrompt}\n\nPrevious JSON failed validation:\n${lastErrors}\nReturn only corrected JSON.`; + + const result = await input.run(prompt); + lastResponse = result.finalResponse; + lastSession = result.providerSessionId ?? lastSession; + + const extracted = tryExtractJson(result.finalResponse); + if (extracted === undefined) { + lastErrors = "Response was not valid JSON"; + input.onRetry?.({ attempt: attempt + 1, errors: lastErrors }); + continue; + } + + const ok = validate(extracted); + if (ok) { + return { + value: extracted, + finalResponse: result.finalResponse, + providerSessionId: result.providerSessionId, + attempts: attempt + 1, + }; + } + + lastErrors = formatAjvErrors(validate.errors); + input.onRetry?.({ attempt: attempt + 1, errors: lastErrors }); + } + + throw new WorkflowEngineError( + "schema", + `Schema validation failed after ${maxRetries + 1} attempts: ${lastErrors}`, + ); +} + +export function augmentPromptForSchema(prompt: string, schema: object): string { + return [ + prompt, + "", + "Respond with ONLY a JSON value that validates against this JSON Schema (no markdown, no prose):", + JSON.stringify(schema), + ].join("\n"); +} + +export function formatAjvErrors( + errors: Array<{ instancePath?: string; message?: string }> | null | undefined, +): string { + if (!errors || errors.length === 0) return "validation failed"; + return errors + .map((error) => { + const path = error.instancePath || "/"; + return `${path} ${error.message ?? "invalid"}`.trim(); + }) + .join("; "); +} + +/** Helper for wiring into agent(): wrap a one-shot provider as retrying schema runner. */ +export function schemaAwareRunProvider( + runProvider: WorkflowRunProvider, + schema: object, + base: Parameters[0], + onRetry?: EnforceSchemaInput["onRetry"], +): Promise { + return enforceAgentSchema({ + schema, + prompt: base.prompt, + onRetry, + run: (prompt) => + runProvider({ + ...base, + prompt, + }), + }); +} diff --git a/src/workflow-store.ts b/src/workflow-store.ts index da29c711..4e57e184 100644 --- a/src/workflow-store.ts +++ b/src/workflow-store.ts @@ -218,6 +218,17 @@ export class WorkflowStore { * Atomically claim a starting run for the worker. * Returns undefined if the run is missing or not claimable. */ + setScriptPath(id: string, scriptPath: string): WorkflowRunRecord { + this.requireRun(id); + const now = isoNow(); + this.database.sqlite + .prepare( + `UPDATE workflow_runs SET script_path = ?, updated_at = ? WHERE id = ?`, + ) + .run(scriptPath, now, id); + return this.requireRun(id); + } + claimRun(id: string, pid: number): WorkflowRunRecord | undefined { const now = isoNow(); const result = this.database.sqlite diff --git a/src/workflow-worktrees.ts b/src/workflow-worktrees.ts new file mode 100644 index 00000000..1bbb3414 --- /dev/null +++ b/src/workflow-worktrees.ts @@ -0,0 +1,146 @@ +import { execFile } from "node:child_process"; +import { mkdir, rm } from "node:fs/promises"; +import { join } from "node:path"; +import { promisify } from "node:util"; +import type { CreateAgentWorktree, WorkflowWorktreeHandle } from "./workflow-api.js"; +import { WorkflowEngineError } from "./workflow-api.js"; + +const execFileAsync = promisify(execFile); + +export interface WorkflowWorktreeHost { + worktreeRoot: string; + /** When set, assert worktree paths stay under this root. */ + allowedRoots?: string[]; +} + +/** + * Create a CreateAgentWorktree bound to host config. + * Layout: `/wf//c/` + */ +export function createWorkflowWorktreeFactory( + host: WorkflowWorktreeHost, +): CreateAgentWorktree { + return async (input) => { + const path = join(host.worktreeRoot, "wf", input.runId, `c${input.callIndex}`); + await mkdir(join(host.worktreeRoot, "wf", input.runId), { recursive: true }); + + let sourceRoot: string; + try { + sourceRoot = ( + await git(["rev-parse", "--show-toplevel"], input.workspaceRoot) + ).trim(); + } catch (error) { + if (isGitUnavailable(error)) { + throw new WorkflowEngineError( + "worktree", + "isolation: 'worktree' requires Git on PATH", + ); + } + throw new WorkflowEngineError( + "worktree", + `isolation: 'worktree' requires a Git repository (not found at ${input.workspaceRoot})`, + ); + } + + const baseSha = + input.baseSha ?? + (await git(["rev-parse", "--verify", "HEAD^{commit}"], sourceRoot)).trim(); + + try { + await git(["worktree", "add", "--detach", path, baseSha], sourceRoot); + } catch (error) { + await rm(path, { recursive: true, force: true }).catch(() => undefined); + const message = error instanceof Error ? error.message : String(error); + throw new WorkflowEngineError( + "worktree", + `Failed to create agent worktree: ${message}`, + ); + } + + return createHandle({ path, sourceRoot }); + }; +} + +function createHandle(input: { + path: string; + sourceRoot: string; +}): WorkflowWorktreeHandle { + return { + path: input.path, + finalize: async (outcome) => { + const dirty = await isDirty(input.path); + if (outcome === "success" && !dirty) { + await removeWorktree(input.sourceRoot, input.path); + return { dirty: false, removed: true }; + } + // Preserve dirty or failed worktrees for diagnosis. + return { dirty, removed: false }; + }, + }; +} + +export async function isDirty(worktreePath: string): Promise { + try { + const status = (await git(["status", "--porcelain=v1"], worktreePath)).trim(); + return status.length > 0; + } catch { + // If status fails, treat as dirty so we don't delete. + return true; + } +} + +export async function removeWorktree( + sourceRoot: string, + worktreePath: string, +): Promise { + try { + await git(["worktree", "remove", "--force", worktreePath], sourceRoot); + } catch { + await rm(worktreePath, { recursive: true, force: true }); + try { + await git(["worktree", "prune"], sourceRoot); + } catch { + // ignore + } + } +} + +export async function resolveWorkspaceHead(workspaceRoot: string): Promise { + try { + return (await git(["rev-parse", "--verify", "HEAD^{commit}"], workspaceRoot)).trim(); + } catch { + return undefined; + } +} + +async function git(args: string[], cwd: string): Promise { + try { + const { stdout } = await execFileAsync("git", args, { + cwd, + maxBuffer: 10 * 1024 * 1024, + }); + return stdout; + } catch (error) { + if (isGitUnavailable(error)) throw error; + const stderr = + typeof error === "object" && error && "stderr" in error + ? String((error as { stderr?: unknown }).stderr ?? "").trim() + : ""; + const stdout = + typeof error === "object" && error && "stdout" in error + ? String((error as { stdout?: unknown }).stdout ?? "").trim() + : ""; + const details = + stderr || stdout || (error instanceof Error ? error.message : String(error)); + throw new Error(details); + } +} + +function isGitUnavailable(error: unknown): boolean { + return Boolean( + typeof error === "object" && + error && + "code" in error && + (error as { code?: unknown }).code === "ENOENT", + ); +}