From d5b00f5611aa03535cfd2badc14f09c7ca42ee83 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Wed, 22 Jul 2026 00:59:27 +0530 Subject: [PATCH] refactor(workflow): model operational failures with better-result --- package-lock.json | 7 + package.json | 3 +- src/cli.ts | 6 +- src/local-agent-adapters.ts | 19 +- src/local-agent-errors.ts | 133 ++++++++++++++ src/local-agent-runtime.ts | 47 +---- src/workflow-api.ts | 32 +++- src/workflow-cli.ts | 87 +++++---- src/workflow-contracts.ts | 1 + src/workflow-engine.test.ts | 46 +++++ src/workflow-engine.ts | 7 + src/workflow-errors.test.ts | 150 +++++++++++++++ src/workflow-errors.ts | 356 ++++++++++++++++++++++++++++++++++++ src/workflow-files.ts | 190 ++++++++++++++----- src/workflow-replay.ts | 8 +- src/workflow-schema.test.ts | 3 + src/workflow-schema.ts | 149 ++++++++++++--- src/workflow-store.ts | 335 ++++++++++++++++++++++++--------- src/workflow-tools.ts | 70 +++++-- src/workflow-worktrees.ts | 170 +++++++++++------ 20 files changed, 1507 insertions(+), 312 deletions(-) create mode 100644 src/local-agent-errors.ts create mode 100644 src/workflow-errors.test.ts create mode 100644 src/workflow-errors.ts diff --git a/package-lock.json b/package-lock.json index 39efe57c..11cd43d2 100644 --- a/package-lock.json +++ b/package-lock.json @@ -20,6 +20,7 @@ "@opencode-ai/sdk": "^1.17.13", "@pierre/diffs": "^1.2.5", "ajv": "^8.20.0", + "better-result": "^2.10.0", "better-sqlite3": "^12.10.0", "diff": "^8.0.3", "drizzle-orm": "^0.45.2", @@ -3446,6 +3447,12 @@ ], "license": "MIT" }, + "node_modules/better-result": { + "version": "2.10.0", + "resolved": "https://registry.npmjs.org/better-result/-/better-result-2.10.0.tgz", + "integrity": "sha512-oQhh0y1qo2/ZKdAAEvHZAqKKiHOFU5k/bW96fE2ScgQOVkJRiHwB+nOS1SgFsYqRlxMDWvefXi9Q3px7QvgNDw==", + "license": "MIT" + }, "node_modules/better-sqlite3": { "version": "12.10.0", "resolved": "https://registry.npmjs.org/better-sqlite3/-/better-sqlite3-12.10.0.tgz", diff --git a/package.json b/package.json index 07f242f6..f30a7c98 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-contracts.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", + "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-contracts.test.ts && tsx src/workflow-errors.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": [], @@ -45,6 +45,7 @@ "@opencode-ai/sdk": "^1.17.13", "@pierre/diffs": "^1.2.5", "ajv": "^8.20.0", + "better-result": "^2.10.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 4789a324..cd70843c 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -43,6 +43,10 @@ import { expandHomePath } from "./roots.js"; import { shutdownHttpServer } from "./server-shutdown.js"; import { runWorkflowCommand } from "./workflow-cli.js"; +import { + isWorkflowOperationError, + workflowCliExitCode, +} from "./workflow-errors.js"; type Command = "serve" | "init" | "doctor" | "config" | "agents" | "workflow" | "help" | "version"; const require = createRequire(import.meta.url); @@ -775,5 +779,5 @@ function checkBashShell(): string { main(process.argv.slice(2)).catch((error) => { console.error(error instanceof Error ? error.message : String(error)); - process.exitCode = 1; + process.exitCode = isWorkflowOperationError(error) ? workflowCliExitCode(error) : 1; }); diff --git a/src/local-agent-adapters.ts b/src/local-agent-adapters.ts index d452c1ba..ce91c90f 100644 --- a/src/local-agent-adapters.ts +++ b/src/local-agent-adapters.ts @@ -1,11 +1,16 @@ import { spawn, spawnSync, type ChildProcessWithoutNullStreams } from "node:child_process"; import { resolve } from "node:path"; import { Readable, Writable } from "node:stream"; +import { Result, type Result as BetterResult } from "better-result"; import type { EffortLevel, OutputFormat, } from "@anthropic-ai/claude-agent-sdk"; import type { JsonSchema } from "./json-types.js"; +import { + classifyAgentProviderError, + type AgentProviderError, +} from "./local-agent-errors.js"; import type { LocalAgentProvider } from "./local-agent-profiles.js"; import { removeDevspaceNodeModulesBinFromPath } from "./local-agent-path.js"; import { @@ -31,7 +36,19 @@ export async function runLocalAgentProvider( provider: LocalAgentProvider, input: LocalAgentRunInput, ): Promise { - return createLocalAgentAdapter(provider).run(input); + const result = await runLocalAgentProviderResult(provider, input); + if (result.isErr()) throw result.error; + return result.value; +} + +export async function runLocalAgentProviderResult( + provider: LocalAgentProvider, + input: LocalAgentRunInput, +): Promise> { + return Result.tryPromise({ + try: () => createLocalAgentAdapter(provider).run(input), + catch: (cause) => classifyAgentProviderError(provider, cause), + }); } export function createLocalAgentAdapter(provider: LocalAgentProvider): LocalAgentAdapter { diff --git a/src/local-agent-errors.ts b/src/local-agent-errors.ts new file mode 100644 index 00000000..19058d0f --- /dev/null +++ b/src/local-agent-errors.ts @@ -0,0 +1,133 @@ +import { TaggedError } from "better-result"; +import type { LocalAgentProvider } from "./local-agent-profiles.js"; + +export class ProviderUnavailableError extends TaggedError( + "ProviderUnavailableError", +)<{ + provider: LocalAgentProvider; + message: string; +}>() { + constructor(provider: LocalAgentProvider, message?: string) { + super({ + provider, + message: message ?? `Agent provider is unavailable: ${provider}`, + }); + } +} + +export class ProviderSchemaUnsupportedError extends TaggedError( + "ProviderSchemaUnsupportedError", +)<{ + provider: LocalAgentProvider; + cause: unknown; + message: string; +}>() { + constructor(provider: LocalAgentProvider, cause: unknown) { + super({ + provider, + cause, + message: `${provider} does not support the requested native output schema: ${errorMessage(cause)}`, + }); + } +} + +export class ProviderCancelledError extends TaggedError( + "ProviderCancelledError", +)<{ + provider: LocalAgentProvider; + cause: unknown; + message: string; +}>() { + constructor(provider: LocalAgentProvider, cause: unknown) { + super({ + provider, + cause, + message: `Agent provider was cancelled: ${provider}`, + }); + } +} + +export class ProviderExecutionError extends TaggedError( + "ProviderExecutionError", +)<{ + provider: LocalAgentProvider; + retryable: boolean; + cause: unknown; + message: string; +}>() { + constructor(input: { + provider: LocalAgentProvider; + cause: unknown; + retryable?: boolean; + }) { + super({ + provider: input.provider, + retryable: input.retryable ?? false, + cause: input.cause, + message: `${input.provider} agent execution failed: ${errorMessage(input.cause)}`, + }); + } +} + +export type AgentProviderError = + | ProviderUnavailableError + | ProviderSchemaUnsupportedError + | ProviderCancelledError + | ProviderExecutionError; + +export function isAgentProviderError(error: unknown): error is AgentProviderError { + return ( + ProviderUnavailableError.is(error) || + ProviderSchemaUnsupportedError.is(error) || + ProviderCancelledError.is(error) || + ProviderExecutionError.is(error) + ); +} + +export function isProviderSchemaUnsupportedError( + error: unknown, +): error is ProviderSchemaUnsupportedError { + return ProviderSchemaUnsupportedError.is(error); +} + +export function isNativeSchemaUnsupportedFailure(error: unknown): boolean { + const message = errorMessage(error).toLowerCase(); + const mentionsSchema = + /output[ _-]?schema/.test(message) || + /json[ _-]?schema/.test(message) || + /structured[ _-]?output/.test(message) || + /output[ _-]?format/.test(message); + const unsupported = + /not supported/.test(message) || + /unsupported/.test(message) || + /invalid (?:output|json )?schema/.test(message) || + /schema (?:is )?invalid/.test(message) || + /unknown (?:field|parameter|option)/.test(message) || + /not available/.test(message); + return mentionsSchema && unsupported; +} + +export function classifyAgentProviderError( + provider: LocalAgentProvider, + cause: unknown, +): AgentProviderError { + if (isAgentProviderError(cause)) return cause; + if (isCancellation(cause)) return new ProviderCancelledError(provider, cause); + if (isNativeSchemaUnsupportedFailure(cause)) { + return new ProviderSchemaUnsupportedError(provider, cause); + } + return new ProviderExecutionError({ provider, cause }); +} + +function isCancellation(error: unknown): boolean { + return Boolean( + error && + typeof error === "object" && + "name" in error && + String((error as { name?: unknown }).name) === "AbortError", + ); +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/src/local-agent-runtime.ts b/src/local-agent-runtime.ts index 0f3896ed..e43a4aff 100644 --- a/src/local-agent-runtime.ts +++ b/src/local-agent-runtime.ts @@ -9,6 +9,16 @@ import type { } from "@openai/codex-sdk"; import type { JsonSchema } from "./json-types.js"; import type { LocalAgentProvider } from "./local-agent-profiles.js"; +import { + isNativeSchemaUnsupportedFailure, + ProviderSchemaUnsupportedError, +} from "./local-agent-errors.js"; + +export { + isNativeSchemaUnsupportedFailure, + isProviderSchemaUnsupportedError, + ProviderSchemaUnsupportedError, +} from "./local-agent-errors.js"; export type LocalAgentWriteMode = "read_only" | "allowed" | "full_access"; @@ -38,39 +48,6 @@ export interface LocalAgentRuntime { run(input: LocalAgentRunInput): Promise; } -export class ProviderSchemaUnsupportedError extends Error { - constructor( - readonly provider: string, - readonly cause: unknown, - ) { - super(`${provider} does not support the requested native output schema: ${errorMessage(cause)}`); - this.name = "ProviderSchemaUnsupportedError"; - } -} - -export function isProviderSchemaUnsupportedError( - error: unknown, -): error is ProviderSchemaUnsupportedError { - return error instanceof ProviderSchemaUnsupportedError; -} - -export function isNativeSchemaUnsupportedFailure(error: unknown): boolean { - const message = errorMessage(error).toLowerCase(); - const mentionsSchema = - /output[ _-]?schema/.test(message) || - /json[ _-]?schema/.test(message) || - /structured[ _-]?output/.test(message) || - /output[ _-]?format/.test(message); - const unsupported = - /not supported/.test(message) || - /unsupported/.test(message) || - /invalid (?:output|json )?schema/.test(message) || - /schema (?:is )?invalid/.test(message) || - /unknown (?:field|parameter|option)/.test(message) || - /not available/.test(message); - return mentionsSchema && unsupported; -} - interface CodexThreadLike { readonly id: string | null; run(prompt: string, turnOptions?: TurnOptions): Promise; @@ -159,7 +136,3 @@ async function defaultCodexFactory(): Promise { const module = await import("@openai/codex-sdk"); return (options) => new module.Codex(options) as Codex; } - -function errorMessage(error: unknown): string { - return error instanceof Error ? error.message : String(error); -} diff --git a/src/workflow-api.ts b/src/workflow-api.ts index b1b09c93..78953f54 100644 --- a/src/workflow-api.ts +++ b/src/workflow-api.ts @@ -269,6 +269,7 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { await semaphore.acquire(deps.signal); let worktree: WorkflowWorktreeHandle | null = null; let worktreePath: string | undefined; + let agentCallBegun = false; try { throwIfCancelled(deps); @@ -307,6 +308,7 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { isolation, worktreePath, }); + agentCallBegun = true; deps.journal.appendEvent({ runId: deps.runId, type: "agent_call_started", @@ -422,6 +424,7 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { return returnValue; } catch (error) { const message = error instanceof Error ? error.message : String(error); + let cleanupError: string | undefined; if (worktree) { try { const finalized = await worktree.finalize("failure"); @@ -438,22 +441,33 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { outcome: "failure", }, }); - } catch { - // preserve original error + } catch (cleanupFailure) { + cleanupError = + cleanupFailure instanceof Error + ? cleanupFailure.message + : String(cleanupFailure); } } - deps.journal.failAgentCall({ - runId: deps.runId, - callIndex: index, - error: message, - worktreePath, - }); + if (agentCallBegun) { + deps.journal.failAgentCall({ + runId: deps.runId, + callIndex: index, + error: message, + worktreePath, + }); + } deps.journal.appendEvent({ runId: deps.runId, type: "agent_call_failed", phase, label: agentOpts.label, - data: { callIndex: index, error: message, isolation, worktreePath }, + data: { + callIndex: index, + error: message, + cleanupError, + isolation, + worktreePath, + }, }); throw error; } finally { diff --git a/src/workflow-cli.ts b/src/workflow-cli.ts index 68d9f917..5ed5bc2c 100644 --- a/src/workflow-cli.ts +++ b/src/workflow-cli.ts @@ -5,7 +5,7 @@ import { resolve } from "node:path"; import { fileURLToPath } from "node:url"; import type { ServerConfig } from "./config.js"; import { parseJsonText, type JsonObject, type JsonValue } from "./json-types.js"; -import { runLocalAgentProvider } from "./local-agent-adapters.js"; +import { runLocalAgentProviderResult } from "./local-agent-adapters.js"; import { getLocalAgentProviderAvailabilitySnapshot } from "./local-agent-availability.js"; import { isLocalAgentProvider, @@ -14,11 +14,11 @@ import { } from "./local-agent-profiles.js"; import { executeWorkflow, mapEngineErrorKind } from "./workflow-engine.js"; import { - parseWorkflowArgFlags, - persistWorkflowScript, - readWorkflowScriptFile, + parseWorkflowArgFlagsResult, + persistWorkflowScriptResult, + readWorkflowScriptFileResult, resolveNamedWorkflowScript, - resolveWorkflowScriptFromPathOrName, + resolveWorkflowScriptFromPathOrNameResult, } from "./workflow-files.js"; import { createWorkflowReplay } from "./workflow-replay.js"; import { parseWorkflowScript } from "./workflow-script.js"; @@ -33,6 +33,11 @@ import { type WorkflowRunSource, } from "./workflow-types.js"; import { parseWorkflowEventPayload } from "./workflow-contracts.js"; +import { + InvalidWorkflowInputError, + WorkflowNotFoundError, + WorkflowStoredDataError, +} from "./workflow-errors.js"; import { createWorkflowWorktreeFactory, resolveWorkspaceHead, @@ -92,12 +97,16 @@ async function runWorkflowRun(args: string[], config: ServerConfig): Promise | --name | --resume )", - ); + throw new InvalidWorkflowInputError({ + code: "missing_source", + message: + "Usage: devspace workflow run (--file | --name | --resume )", + }); } const store = createWorkflowStore(config); @@ -111,11 +120,15 @@ async function runWorkflowRun(args: string[], config: ServerConfig): Promise"); const store = createWorkflowStore(config); try { - const run = store.requestCancel(runId); + const requested = store.requestCancelResult(runId); + if (requested.isErr()) throw requested.error; + const run = requested.value; console.log(formatRunLine(run)); if (run.pid && (run.status === "running" || run.status === "starting")) { try { @@ -233,7 +254,8 @@ async function runWorkflowCancel(args: string[], config: ServerConfig): Promise< } const latest = store.getRun(runId); if (latest && (latest.status === "running" || latest.status === "starting")) { - store.cancelRun(runId, "cancelled (hard kill)"); + const cancelled = store.cancelRunResult(runId, "cancelled (hard kill)"); + if (cancelled.isErr()) throw cancelled.error; } } } @@ -266,11 +288,12 @@ export async function runWorkflowWorker( if (!runId) throw new Error("Usage: devspace workflow __worker "); const store = createWorkflowStore(config); - const claimed = store.claimRun(runId, process.pid); - if (!claimed) { + const claim = store.claimRunResult(runId, process.pid); + if (claim.isErr()) { store.close(); - throw new Error(`Cannot claim workflow run ${runId} (missing or not starting)`); + throw claim.error; } + const claimed = claim.value; const abort = new AbortController(); const heartbeat = setInterval(() => { @@ -295,8 +318,8 @@ export async function runWorkflowWorker( try { argsValue = parseJsonText(claimed.argsJson); if (argsValue === null) argsValue = undefined; - } catch { - argsValue = undefined; + } catch (cause) { + throw new WorkflowStoredDataError(`${claimed.id}.argsJson`, cause); } const replay = claimed.resumedFromRunId @@ -327,7 +350,7 @@ export async function runWorkflowWorker( if (abort.signal.aborted || store.isCancelRequested(runId)) { throw Object.assign(new Error("Workflow cancelled"), { name: "AbortError" }); } - const providerResult = await runLocalAgentProvider(input.provider, { + const providerRun = await runLocalAgentProviderResult(input.provider, { prompt: input.prompt, workspace: input.workspace, providerSessionId: input.providerSessionId, @@ -336,6 +359,8 @@ export async function runWorkflowWorker( writeMode: "allowed", schema: input.schema, }); + if (providerRun.isErr()) throw providerRun.error; + const providerResult = providerRun.value; return { finalResponse: providerResult.finalResponse, providerSessionId: providerResult.providerSessionId ?? undefined, diff --git a/src/workflow-contracts.ts b/src/workflow-contracts.ts index bf43d58e..a6947c63 100644 --- a/src/workflow-contracts.ts +++ b/src/workflow-contracts.ts @@ -186,6 +186,7 @@ export const workflowEventPayloadSchemas = { .object({ callIndex: z.number().int().nonnegative(), error: z.string(), + cleanupError: z.string().optional(), isolation: agentIsolationModeSchema, worktreePath: z.string().optional(), }) diff --git a/src/workflow-engine.test.ts b/src/workflow-engine.test.ts index 030fb590..5c3438b3 100644 --- a/src/workflow-engine.test.ts +++ b/src/workflow-engine.test.ts @@ -241,6 +241,52 @@ import { createStubBudget } from "./workflow-types.js"; await rm(dir, { recursive: true, force: true }); } +// --------------------------------------------------------------------------- +// worktree setup failure preserves the primary error before journal begin +// --------------------------------------------------------------------------- +{ + const dir = await mkdtemp(join(tmpdir(), "wf-iso-fail-")); + const store = new WorkflowStore(dir); + const run = store.createRun({ + name: "iso-fail", + source: "inline", + scriptPath: "inline", + scriptHash: "h", + workspaceRoot: dir, + }); + const api = createWorkflowApi({ + runId: run.id, + journal: store, + meta: { name: "iso-fail", description: "d" }, + args: undefined, + concurrency: 1, + signal: new AbortController().signal, + workspaceRoot: dir, + enabledProviders: ["codex"], + createWorktree: async () => { + throw new Error("expected worktree setup failure"); + }, + runProvider: async () => ({ finalResponse: "unreachable" }), + }); + + const runIsolated = api.agent as ( + prompt: string, + opts: { isolation: "worktree" }, + ) => Promise; + await assert.rejects( + () => runIsolated("do", { isolation: "worktree" }), + /expected worktree setup failure/, + ); + assert.equal(store.listAgentCalls(run.id).length, 0); + const failed = store + .drainEvents(run.id) + .events.find((event) => event.type === "agent_call_failed"); + assert.ok(failed); + + store.close(); + await rm(dir, { recursive: true, force: true }); +} + // --------------------------------------------------------------------------- // provider resolve order + no writeMode // --------------------------------------------------------------------------- diff --git a/src/workflow-engine.ts b/src/workflow-engine.ts index ff2a4726..1611e7ca 100644 --- a/src/workflow-engine.ts +++ b/src/workflow-engine.ts @@ -18,6 +18,10 @@ import { type WorkflowMeta, type WorkflowErrorKind, } from "./workflow-types.js"; +import { + isWorkflowOperationError, + workflowErrorKind, +} from "./workflow-errors.js"; export interface ExecuteWorkflowOptions { /** Pre-parsed script, or pass `source` instead. */ @@ -185,6 +189,9 @@ export function mapEngineErrorKind(error: unknown): WorkflowErrorKind { if (error instanceof WorkflowEngineError) { return error.kind; } + if (isWorkflowOperationError(error)) { + return workflowErrorKind(error); + } if (error && typeof error === "object" && "name" in error) { const name = String((error as { name: string }).name); if (name === "WorkflowScriptError") { diff --git a/src/workflow-errors.test.ts b/src/workflow-errors.test.ts new file mode 100644 index 00000000..bbc38590 --- /dev/null +++ b/src/workflow-errors.test.ts @@ -0,0 +1,150 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + classifyAgentProviderError, + ProviderCancelledError, + ProviderExecutionError, + ProviderSchemaUnsupportedError, +} from "./local-agent-errors.js"; +import { + parseWorkflowArgFlagsResult, + readWorkflowScriptFileResult, + resolveNamedWorkflowScriptResult, +} from "./workflow-files.js"; +import { + InvalidRunTransitionError, + InvalidWorkflowInputError, + NamedWorkflowNotFoundError, + SchemaRetriesExhaustedError, + WorkflowFileNotFoundError, + WorkflowNotFoundError, + WorktreeOperationError, + serializeWorkflowError, + workflowCliExitCode, + workflowErrorKind, +} from "./workflow-errors.js"; +import { WorkflowStore } from "./workflow-store.js"; +import { enforceAgentSchemaResult } from "./workflow-schema.js"; +import { createWorkflowWorktreeResult } from "./workflow-worktrees.js"; + +{ + const invalid = parseWorkflowArgFlagsResult(["--arg", "missing-equals"]); + assert.ok(invalid.isErr()); + if (invalid.isErr()) assert.ok(InvalidWorkflowInputError.is(invalid.error)); +} + +{ + const missing = await readWorkflowScriptFileResult("/definitely/missing/workflow.js"); + assert.ok(missing.isErr()); + if (missing.isErr()) assert.ok(WorkflowFileNotFoundError.is(missing.error)); +} + +{ + const root = await mkdtemp(join(tmpdir(), "wf-result-files-")); + try { + const missing = await resolveNamedWorkflowScriptResult({ + name: "missing", + workspaceRoot: root, + }); + assert.ok(missing.isErr()); + if (missing.isErr()) assert.ok(NamedWorkflowNotFoundError.is(missing.error)); + } finally { + await rm(root, { recursive: true, force: true }); + } +} + +{ + const cancelled = Object.assign(new Error("cancel"), { name: "AbortError" }); + assert.ok(ProviderCancelledError.is(classifyAgentProviderError("codex", cancelled))); + assert.ok( + ProviderSchemaUnsupportedError.is( + classifyAgentProviderError( + "claude", + new Error("structured output format is not supported"), + ), + ), + ); + assert.ok( + ProviderExecutionError.is( + classifyAgentProviderError("opencode", new Error("authentication failed")), + ), + ); + + const unavailable = new ProviderSchemaUnsupportedError( + "codex", + new Error("output schema unsupported"), + ); + assert.equal(workflowCliExitCode(unavailable), 5); + assert.deepEqual(serializeWorkflowError(unavailable), { + code: "ProviderSchemaUnsupportedError", + message: unavailable.message, + kind: "schema", + retryable: false, + }); +} + +{ + const root = await mkdtemp(join(tmpdir(), "wf-result-store-")); + const store = new WorkflowStore(root); + try { + const missing = store.claimRunResult("wfr_missing", process.pid); + assert.ok(missing.isErr()); + if (missing.isErr()) assert.ok(WorkflowNotFoundError.is(missing.error)); + + const run = store.createRun({ + name: "result-store", + source: "inline", + scriptPath: "inline", + scriptHash: "h", + workspaceRoot: root, + }); + assert.ok(store.claimRunResult(run.id, process.pid).isOk()); + const duplicate = store.claimRunResult(run.id, process.pid); + assert.ok(duplicate.isErr()); + if (duplicate.isErr()) assert.ok(InvalidRunTransitionError.is(duplicate.error)); + } finally { + store.close(); + await rm(root, { recursive: true, force: true }); + } +} + +{ + const exhausted = await enforceAgentSchemaResult({ + schema: { type: "object" }, + prompt: "return json", + provider: "opencode", + maxRetries: 0, + run: async () => ({ finalResponse: "not json" }), + }); + assert.ok(exhausted.isErr()); + if (exhausted.isErr()) { + assert.ok(SchemaRetriesExhaustedError.is(exhausted.error)); + assert.equal(workflowErrorKind(exhausted.error), "schema"); + } +} + +{ + const root = await mkdtemp(join(tmpdir(), "wf-result-worktree-")); + try { + const created = await createWorkflowWorktreeResult( + { worktreeRoot: join(root, "worktrees") }, + { + runId: "wfr_result", + callIndex: 0, + workspaceRoot: root, + }, + ); + assert.ok(created.isErr()); + if (created.isErr()) { + assert.ok(WorktreeOperationError.is(created.error)); + assert.equal(workflowErrorKind(created.error), "worktree"); + assert.ok(created.error.cause); + } + } finally { + await rm(root, { recursive: true, force: true }); + } +} + +console.log("workflow-errors.test.ts: ok"); diff --git a/src/workflow-errors.ts b/src/workflow-errors.ts new file mode 100644 index 00000000..4349cc5f --- /dev/null +++ b/src/workflow-errors.ts @@ -0,0 +1,356 @@ +import { TaggedError } from "better-result"; +import { + isAgentProviderError, + ProviderExecutionError, + ProviderUnavailableError, + type AgentProviderError, +} from "./local-agent-errors.js"; +import type { + WorkflowErrorKind, + WorkflowRunStatus, +} from "./workflow-types.js"; + +export class InvalidWorkflowInputError extends TaggedError( + "InvalidWorkflowInputError", +)<{ + code: "ambiguous_source" | "missing_source" | "invalid_name" | "invalid_argument"; + message: string; +}>() {} + +export class WorkflowFileNotFoundError extends TaggedError( + "WorkflowFileNotFoundError", +)<{ + path: string; + message: string; +}>() { + constructor(path: string) { + super({ path, message: `Script file not found: ${path}` }); + } +} + +export class WorkflowFileReadError extends TaggedError( + "WorkflowFileReadError", +)<{ + path: string; + cause: unknown; + message: string; +}>() { + constructor(path: string, cause: unknown) { + super({ + path, + cause, + message: `Unable to read workflow script ${path}: ${errorMessage(cause)}`, + }); + } +} + +export class WorkflowFileWriteError extends TaggedError( + "WorkflowFileWriteError", +)<{ + path: string; + cause: unknown; + message: string; +}>() { + constructor(path: string, cause: unknown) { + super({ + path, + cause, + message: `Unable to persist workflow script ${path}: ${errorMessage(cause)}`, + }); + } +} + +export class NamedWorkflowNotFoundError extends TaggedError( + "NamedWorkflowNotFoundError", +)<{ + name: string; + candidates: string[]; + message: string; +}>() { + constructor(name: string, candidates: string[]) { + super({ + name, + candidates, + message: `Named workflow not found: ${name}. Looked in ${candidates.join(", ")}`, + }); + } +} + +export class WorkflowNotFoundError extends TaggedError( + "WorkflowNotFoundError", +)<{ + runId: string; + message: string; +}>() { + constructor(runId: string) { + super({ runId, message: `Unknown workflow run: ${runId}` }); + } +} + +export class InvalidRunTransitionError extends TaggedError( + "InvalidRunTransitionError", +)<{ + runId: string; + from: WorkflowRunStatus; + operation: "claim" | "complete" | "fail" | "cancel" | "set_script_path"; + message: string; +}>() { + constructor(input: { + runId: string; + from: WorkflowRunStatus; + operation: "claim" | "complete" | "fail" | "cancel" | "set_script_path"; + }) { + super({ + ...input, + message: `Cannot ${input.operation} workflow run ${input.runId} in status ${input.from}`, + }); + } +} + +export class WorkflowStoreError extends TaggedError( + "WorkflowStoreError", +)<{ + operation: string; + cause: unknown; + message: string; +}>() { + constructor(operation: string, cause: unknown) { + super({ + operation, + cause, + message: `Workflow store ${operation} failed: ${errorMessage(cause)}`, + }); + } +} + +export class WorkflowStoredDataError extends TaggedError( + "WorkflowStoredDataError", +)<{ + record: string; + cause: unknown; + message: string; +}>() { + constructor(record: string, cause: unknown) { + super({ + record, + cause, + message: `Stored workflow data is invalid (${record}): ${errorMessage(cause)}`, + }); + } +} + +export class WorktreeOperationError extends TaggedError( + "WorktreeOperationError", +)<{ + operation: "create" | "inspect" | "finalize" | "remove"; + runId?: string; + callIndex?: number; + path?: string; + cause: unknown; + message: string; +}>() { + constructor(input: { + operation: "create" | "inspect" | "finalize" | "remove"; + runId?: string; + callIndex?: number; + path?: string; + cause: unknown; + }) { + super({ + ...input, + message: `Workflow worktree ${input.operation} failed${input.path ? ` at ${input.path}` : ""}: ${errorMessage(input.cause)}`, + }); + } +} + +export interface SchemaIssue { + path: string; + message: string; +} + +export class InvalidAgentJsonError extends TaggedError( + "InvalidAgentJsonError", +)<{ + attempt: number; + mode: "native" | "prompt"; + responseExcerpt: string; + message: string; +}>() { + constructor(input: { + attempt: number; + mode: "native" | "prompt"; + responseExcerpt: string; + }) { + super({ + ...input, + message: `Agent response was not valid JSON on attempt ${input.attempt}`, + }); + } +} + +export class AgentSchemaValidationError extends TaggedError( + "AgentSchemaValidationError", +)<{ + attempt: number; + mode: "native" | "prompt"; + issues: SchemaIssue[]; + message: string; +}>() { + constructor(input: { + attempt: number; + mode: "native" | "prompt"; + issues: SchemaIssue[]; + }) { + super({ + ...input, + message: `Agent response failed schema validation on attempt ${input.attempt}: ${input.issues.map((issue) => `${issue.path} ${issue.message}`).join("; ")}`, + }); + } +} + +export class SchemaConfigurationError extends TaggedError( + "SchemaConfigurationError", +)<{ + cause: unknown; + message: string; +}>() { + constructor(cause: unknown) { + super({ + cause, + message: `Unable to compile agent JSON Schema: ${errorMessage(cause)}`, + }); + } +} + +export type SchemaAttemptError = InvalidAgentJsonError | AgentSchemaValidationError; + +export class SchemaRetriesExhaustedError extends TaggedError( + "SchemaRetriesExhaustedError", +)<{ + attempts: number; + lastFailure: SchemaAttemptError; + message: string; +}>() { + constructor(attempts: number, lastFailure: SchemaAttemptError) { + super({ + attempts, + lastFailure, + message: `Schema validation failed after ${attempts} attempts: ${lastFailure.message}`, + }); + } +} + +export type WorkflowOperationError = + | InvalidWorkflowInputError + | WorkflowFileNotFoundError + | WorkflowFileReadError + | WorkflowFileWriteError + | NamedWorkflowNotFoundError + | WorkflowNotFoundError + | InvalidRunTransitionError + | WorkflowStoreError + | WorkflowStoredDataError + | WorktreeOperationError + | InvalidAgentJsonError + | AgentSchemaValidationError + | SchemaConfigurationError + | SchemaRetriesExhaustedError + | AgentProviderError; + +export function isWorkflowOperationError(error: unknown): error is WorkflowOperationError { + return ( + InvalidWorkflowInputError.is(error) || + WorkflowFileNotFoundError.is(error) || + WorkflowFileReadError.is(error) || + WorkflowFileWriteError.is(error) || + NamedWorkflowNotFoundError.is(error) || + WorkflowNotFoundError.is(error) || + InvalidRunTransitionError.is(error) || + WorkflowStoreError.is(error) || + WorkflowStoredDataError.is(error) || + WorktreeOperationError.is(error) || + InvalidAgentJsonError.is(error) || + AgentSchemaValidationError.is(error) || + SchemaConfigurationError.is(error) || + SchemaRetriesExhaustedError.is(error) || + isAgentProviderError(error) + ); +} + +export function workflowErrorKind(error: WorkflowOperationError): WorkflowErrorKind { + switch (error._tag) { + case "InvalidWorkflowInputError": + case "WorkflowFileNotFoundError": + case "WorkflowFileReadError": + case "WorkflowFileWriteError": + case "NamedWorkflowNotFoundError": + return "path"; + case "WorkflowNotFoundError": + case "InvalidRunTransitionError": + case "WorkflowStoreError": + case "WorkflowStoredDataError": + return "internal"; + case "WorktreeOperationError": + return "worktree"; + case "InvalidAgentJsonError": + case "AgentSchemaValidationError": + case "SchemaConfigurationError": + case "SchemaRetriesExhaustedError": + case "ProviderSchemaUnsupportedError": + return "schema"; + case "ProviderCancelledError": + return "cancelled"; + case "ProviderUnavailableError": + return "provider_unavailable"; + case "ProviderExecutionError": + return "provider"; + } +} + +export function workflowCliExitCode(error: WorkflowOperationError): number { + switch (error._tag) { + case "InvalidWorkflowInputError": + return 2; + case "WorkflowFileNotFoundError": + case "NamedWorkflowNotFoundError": + case "WorkflowNotFoundError": + return 3; + case "ProviderUnavailableError": + return 4; + case "ProviderCancelledError": + return 130; + case "InvalidAgentJsonError": + case "AgentSchemaValidationError": + case "SchemaConfigurationError": + case "SchemaRetriesExhaustedError": + case "ProviderSchemaUnsupportedError": + return 5; + case "WorkflowFileReadError": + case "WorkflowFileWriteError": + case "InvalidRunTransitionError": + case "WorkflowStoreError": + case "WorkflowStoredDataError": + case "WorktreeOperationError": + case "ProviderExecutionError": + return 1; + } +} + +export function serializeWorkflowError(error: WorkflowOperationError): { + code: WorkflowOperationError["_tag"]; + message: string; + kind: WorkflowErrorKind; + retryable: boolean; +} { + return { + code: error._tag, + message: error.message, + kind: workflowErrorKind(error), + retryable: + ProviderExecutionError.is(error) ? error.retryable : ProviderUnavailableError.is(error), + }; +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/src/workflow-files.ts b/src/workflow-files.ts index 7615980f..2a0cfadf 100644 --- a/src/workflow-files.ts +++ b/src/workflow-files.ts @@ -1,8 +1,16 @@ import { createHash, randomBytes } from "node:crypto"; -import { access, mkdir, readFile, writeFile } from "node:fs/promises"; +import { mkdir, readFile, writeFile } from "node:fs/promises"; import { basename, dirname, extname, isAbsolute, join, resolve } from "node:path"; +import { Result, type Result as BetterResult } from "better-result"; import { hashSource } from "./workflow-script.js"; import { jsonValueSchema, type JsonValue } from "./json-types.js"; +import { + InvalidWorkflowInputError, + NamedWorkflowNotFoundError, + WorkflowFileNotFoundError, + WorkflowFileReadError, + WorkflowFileWriteError, +} from "./workflow-errors.js"; export class WorkflowPathError extends Error { constructor(message: string) { @@ -19,6 +27,12 @@ export interface ResolvedWorkflowScript { origin: "file" | "named" | "inline" | "resume"; } +export type WorkflowFileResolveError = + | InvalidWorkflowInputError + | NamedWorkflowNotFoundError + | WorkflowFileNotFoundError + | WorkflowFileReadError; + /** * Persist script under stateDir for worker re-read / audit. * Returns absolute path written. @@ -29,27 +43,58 @@ export async function persistWorkflowScript(input: { source: string; preferredName?: string; }): Promise { + const result = await persistWorkflowScriptResult(input); + if (result.isErr()) throw result.error; + return result.value; +} + +export async function persistWorkflowScriptResult(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; + return Result.tryPromise({ + try: async () => { + await mkdir(dir, { recursive: true }); + await writeFile(path, input.source, { encoding: "utf8", mode: 0o600 }); + return path; + }, + catch: (cause) => new WorkflowFileWriteError(path, cause), + }); } export async function readWorkflowScriptFile(path: string): Promise { + const result = await readWorkflowScriptFileResult(path); + if (result.isErr()) throwPathCompatibilityError(result.error); + return result.value; +} + +export async function readWorkflowScriptFileResult( + 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", - }; + return Result.tryPromise({ + try: async () => { + const source = await readFile(scriptPath, "utf8"); + return { + source, + scriptPath, + scriptHash: hashSource(source), + nameHint: basename(scriptPath, extname(scriptPath)), + origin: "file" as const, + }; + }, + catch: (cause) => + isFileNotFound(cause) + ? new WorkflowFileNotFoundError(scriptPath) + : new WorkflowFileReadError(scriptPath, cause), + }); } /** @@ -64,9 +109,24 @@ export async function resolveNamedWorkflowScript(input: { workspaceRoot: string; stateDir?: string; }): Promise { + const result = await resolveNamedWorkflowScriptResult(input); + if (result.isErr()) throwPathCompatibilityError(result.error); + return result.value; +} + +export async function resolveNamedWorkflowScriptResult(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)}`); + return Result.err( + new InvalidWorkflowInputError({ + code: "invalid_name", + message: `Invalid workflow name: ${JSON.stringify(input.name)}`, + }), + ); } const candidates = [ join(input.workspaceRoot, ".devspace", "workflows", `${name}.js`), @@ -76,23 +136,14 @@ export async function resolveNamedWorkflowScript(input: { 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 + const result = await readWorkflowScriptFileResult(candidate); + if (result.isOk()) { + return Result.ok({ ...result.value, nameHint: name, origin: "named" as const }); } + if (WorkflowFileNotFoundError.is(result.error)) continue; + return result; } - throw new WorkflowPathError( - `Named workflow not found: ${name}. Looked in ${candidates.join(", ")}`, - ); + return Result.err(new NamedWorkflowNotFoundError(name, candidates)); } export async function resolveWorkflowScriptFromPathOrName(input: { @@ -101,29 +152,61 @@ export async function resolveWorkflowScriptFromPathOrName(input: { workspaceRoot: string; stateDir?: string; }): Promise { + const result = await resolveWorkflowScriptFromPathOrNameResult(input); + if (result.isErr()) throwPathCompatibilityError(result.error); + return result.value; +} + +export async function resolveWorkflowScriptFromPathOrNameResult(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"); + return Result.err( + new InvalidWorkflowInputError({ + code: "ambiguous_source", + message: "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); + return readWorkflowScriptFileResult(path); } if (input.name) { - return resolveNamedWorkflowScript({ + return resolveNamedWorkflowScriptResult({ name: input.name, workspaceRoot: input.workspaceRoot, stateDir: input.stateDir, }); } - throw new WorkflowPathError("Provide --file or --name "); + return Result.err( + new InvalidWorkflowInputError({ + code: "missing_source", + message: "Provide --file or --name ", + }), + ); } export function parseWorkflowArgFlags(tokens: string[]): { args: Record; rest: string[]; } { + const result = parseWorkflowArgFlagsResult(tokens); + if (result.isErr()) throwPathCompatibilityError(result.error); + return result.value; +} + +export function parseWorkflowArgFlagsResult( + tokens: string[], +): BetterResult< + { args: Record; rest: string[] }, + InvalidWorkflowInputError +> { const args: Record = {}; const rest: string[] = []; for (let i = 0; i < tokens.length; i += 1) { @@ -131,7 +214,12 @@ export function parseWorkflowArgFlags(tokens: string[]): { if (token === "--arg") { const pair = tokens[++i]; if (!pair || !pair.includes("=")) { - throw new WorkflowPathError("--arg requires key=value"); + return Result.err( + new InvalidWorkflowInputError({ + code: "invalid_argument", + message: "--arg requires key=value", + }), + ); } const eq = pair.indexOf("="); const key = pair.slice(0, eq); @@ -142,13 +230,20 @@ export function parseWorkflowArgFlags(tokens: string[]): { if (token.startsWith("--arg=")) { const pair = token.slice("--arg=".length); const eq = pair.indexOf("="); - if (eq < 0) throw new WorkflowPathError("--arg requires key=value"); + if (eq < 0) { + return Result.err( + new InvalidWorkflowInputError({ + code: "invalid_argument", + message: "--arg requires key=value", + }), + ); + } args[pair.slice(0, eq)] = coerceArgValue(pair.slice(eq + 1)); continue; } rest.push(token); } - return { args, rest }; + return Result.ok({ args, rest }); } function coerceArgValue(raw: string): JsonValue { @@ -159,14 +254,6 @@ function coerceArgValue(raw: string): JsonValue { } } -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, "-") @@ -185,3 +272,18 @@ export function contentHash(source: string): string { export function dirnameOf(path: string): string { return dirname(path); } + +function isFileNotFound(error: unknown): boolean { + return Boolean( + error && + typeof error === "object" && + "code" in error && + (error as { code?: unknown }).code === "ENOENT", + ); +} + +function throwPathCompatibilityError(error: Error): never { + const compatible = new WorkflowPathError(error.message); + compatible.cause = error; + throw compatible; +} diff --git a/src/workflow-replay.ts b/src/workflow-replay.ts index 3038d1b0..24f35418 100644 --- a/src/workflow-replay.ts +++ b/src/workflow-replay.ts @@ -1,6 +1,7 @@ import type { WorkflowAgentCallRecord } from "./workflow-types.js"; import type { WorkflowReplay, WorkflowReplayHit } from "./workflow-api.js"; import { parseJsonText } from "./json-types.js"; +import { WorkflowStoredDataError } from "./workflow-errors.js"; /** * Resume matcher: @@ -69,8 +70,11 @@ function toHit(call: WorkflowAgentCallRecord): WorkflowReplayHit { structuredJson: call.structuredJson, providerSessionId: call.providerSessionId, }; - } catch { - // fall through to text + } catch (cause) { + throw new WorkflowStoredDataError( + `${call.runId}.agentCalls[${call.callIndex}].structuredJson`, + cause, + ); } } return { diff --git a/src/workflow-schema.test.ts b/src/workflow-schema.test.ts index 8e92ffcd..a05b2dc6 100644 --- a/src/workflow-schema.test.ts +++ b/src/workflow-schema.test.ts @@ -37,6 +37,7 @@ assert.ok(!supportsNativeStructuredOutput("opencode")); additionalProperties: false, }, prompt: "give n", + provider: "opencode", run: async () => { attempts += 1; if (attempts === 1) return { finalResponse: '{"n":"x"}' }; @@ -55,6 +56,7 @@ assert.ok(!supportsNativeStructuredOutput("opencode")); enforceAgentSchema({ schema: { type: "object", properties: { n: { type: "number" } }, required: ["n"] }, prompt: "x", + provider: "opencode", maxRetries: 1, run: async () => ({ finalResponse: "not json" }), }), @@ -205,6 +207,7 @@ assert.ok(!supportsNativeStructuredOutput("opencode")); enforceAgentSchema({ schema: { type: "object" }, prompt: "x", + provider: "opencode", maxRetries: 0, onRetry: ({ attempt }) => retries.push(attempt), run: async () => ({ finalResponse: "not json" }), diff --git a/src/workflow-schema.ts b/src/workflow-schema.ts index b5234b45..0249c2dd 100644 --- a/src/workflow-schema.ts +++ b/src/workflow-schema.ts @@ -1,8 +1,13 @@ import { createRequire } from "node:module"; +import { Result, type Result as BetterResult } from "better-result"; import { WORKFLOW_MAX_SCHEMA_RETRIES } from "./workflow-types.js"; import { tryExtractJson, WorkflowEngineError } from "./workflow-api.js"; import type { WorkflowProviderRunResult, WorkflowRunProvider } from "./workflow-api.js"; -import { isProviderSchemaUnsupportedError } from "./local-agent-runtime.js"; +import { + classifyAgentProviderError, + isProviderSchemaUnsupportedError, + type AgentProviderError, +} from "./local-agent-errors.js"; import { supportsNativeStructuredOutput } from "./local-agent-capabilities.js"; import type { LocalAgentProvider } from "./local-agent-profiles.js"; import { @@ -10,6 +15,13 @@ import { type JsonSchema, type JsonValue, } from "./json-types.js"; +import { + AgentSchemaValidationError, + InvalidAgentJsonError, + SchemaConfigurationError, + SchemaRetriesExhaustedError, + type SchemaAttemptError, +} from "./workflow-errors.js"; const require = createRequire(import.meta.url); @@ -40,7 +52,7 @@ export interface EnforceSchemaInput { * Provider id for native-vs-prompt policy. When in NATIVE_SCHEMA_PROVIDERS, * attempt 0 uses raw prompt + native structured path; later attempts repair via prompt. */ - provider?: LocalAgentProvider; + provider: LocalAgentProvider; run: ( prompt: string, opts: { @@ -64,6 +76,11 @@ export interface EnforceSchemaResult { mode: SchemaEnforceMode; } +export type EnforceSchemaError = + | AgentProviderError + | SchemaConfigurationError + | SchemaRetriesExhaustedError; + /** * Native-first for codex/claude; otherwise prompt+extract+Ajv. Always Ajv-validate. * Retries ≤ WORKFLOW_MAX_SCHEMA_RETRIES after the first attempt. @@ -71,14 +88,35 @@ export interface EnforceSchemaResult { 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 result = await enforceAgentSchemaResult(input); + if (result.isOk()) return result.value; + if ( + SchemaConfigurationError.is(result.error) || + SchemaRetriesExhaustedError.is(result.error) + ) { + throw new WorkflowEngineError("schema", result.error.message); + } + throw result.error; +} + +export async function enforceAgentSchemaResult( + input: EnforceSchemaInput, +): Promise> { + const compiled = Result.try({ + try: () => { + const Ajv = loadAjv(); + const ajv = new Ajv({ allErrors: true, strict: false }); + return ajv.compile(input.schema); + }, + catch: (cause) => new SchemaConfigurationError(cause), + }); + if (compiled.isErr()) return compiled; + const validate = compiled.value; const maxRetries = input.maxRetries ?? WORKFLOW_MAX_SCHEMA_RETRIES; - const native = Boolean(input.provider && supportsNativeStructuredOutput(input.provider)); + const native = supportsNativeStructuredOutput(input.provider); const basePrompt = augmentPromptForSchema(input.prompt, input.schema); - let lastErrors = "unknown validation error"; + let lastFailure: SchemaAttemptError | undefined; let providerSessionId: string | undefined; for (let attempt = 0; attempt <= maxRetries; attempt += 1) { @@ -89,26 +127,39 @@ export async function enforceAgentSchema( ? input.prompt : attempt === 0 ? basePrompt - : `${basePrompt}\n\nPrevious JSON failed validation:\n${lastErrors}\nReturn only corrected JSON.`; - - let result: WorkflowProviderRunResult; - try { - result = await input.run(prompt, { mode, providerSessionId }); - } catch (error) { - if (mode === "native" && isProviderSchemaUnsupportedError(error) && attempt < maxRetries) { - lastErrors = error.message; - input.onRetry?.({ attempt: attempt + 1, errors: lastErrors, mode }); + : `${basePrompt}\n\nPrevious JSON failed validation:\n${lastFailure?.message ?? "unknown validation error"}\nReturn only corrected JSON.`; + + const runResult = await Result.tryPromise({ + try: () => input.run(prompt, { mode, providerSessionId }), + catch: (cause) => classifyAgentProviderError(input.provider, cause), + }); + if (runResult.isErr()) { + if ( + mode === "native" && + isProviderSchemaUnsupportedError(runResult.error) && + attempt < maxRetries + ) { + input.onRetry?.({ + attempt: attempt + 1, + errors: runResult.error.message, + mode, + }); continue; } - throw error; + return runResult; } + const result = runResult.value; providerSessionId = result.providerSessionId ?? providerSessionId; const candidates = structuredCandidates(result); if (candidates.length === 0) { - lastErrors = "Response was not valid JSON"; + lastFailure = new InvalidAgentJsonError({ + attempt: attempt + 1, + mode, + responseExcerpt: result.finalResponse.slice(0, 500), + }); if (attempt < maxRetries) { - input.onRetry?.({ attempt: attempt + 1, errors: lastErrors, mode }); + input.onRetry?.({ attempt: attempt + 1, errors: lastFailure.message, mode }); } continue; } @@ -116,25 +167,36 @@ export async function enforceAgentSchema( for (const candidate of candidates) { const ok = validate(candidate); if (ok) { - return { + return Result.ok({ value: candidate, finalResponse: result.finalResponse, providerSessionId, attempts: attempt + 1, mode, - }; - }; + }); + } } - lastErrors = formatAjvErrors(validate.errors); + lastFailure = new AgentSchemaValidationError({ + attempt: attempt + 1, + mode, + issues: toSchemaIssues(validate.errors), + }); if (attempt < maxRetries) { - input.onRetry?.({ attempt: attempt + 1, errors: lastErrors, mode }); + input.onRetry?.({ attempt: attempt + 1, errors: lastFailure.message, mode }); } } - throw new WorkflowEngineError( - "schema", - `Schema validation failed after ${maxRetries + 1} attempts: ${lastErrors}`, + return Result.err( + new SchemaRetriesExhaustedError( + maxRetries + 1, + lastFailure ?? + new InvalidAgentJsonError({ + attempt: maxRetries + 1, + mode: native ? "native" : "prompt", + responseExcerpt: "", + }), + ), ); } @@ -159,6 +221,18 @@ export function formatAjvErrors( .join("; "); } +function toSchemaIssues( + errors: Array<{ instancePath?: string; message?: string }> | null | undefined, +): Array<{ path: string; message: string }> { + if (!errors || errors.length === 0) { + return [{ path: "/", message: "validation failed" }]; + } + return errors.map((error) => ({ + path: error.instancePath || "/", + message: error.message ?? "invalid", + })); +} + /** Helper for wiring into agent(): wrap a one-shot provider as retrying schema runner. */ export function schemaAwareRunProvider( runProvider: WorkflowRunProvider, @@ -181,6 +255,27 @@ export function schemaAwareRunProvider( }); } +export function schemaAwareRunProviderResult( + runProvider: WorkflowRunProvider, + schema: JsonSchema, + base: Parameters[0], + onRetry?: EnforceSchemaInput["onRetry"], +): Promise> { + return enforceAgentSchemaResult({ + schema, + prompt: base.prompt, + provider: base.provider, + onRetry, + run: (prompt, options) => + runProvider({ + ...base, + prompt, + providerSessionId: options.providerSessionId, + ...(options.mode === "native" ? { schema } : {}), + }), + }); +} + function structuredCandidates(result: WorkflowProviderRunResult): JsonValue[] { const candidates: JsonValue[] = []; if (result.structured !== undefined) { diff --git a/src/workflow-store.ts b/src/workflow-store.ts index 90d866c3..085757d2 100644 --- a/src/workflow-store.ts +++ b/src/workflow-store.ts @@ -1,5 +1,6 @@ import { randomUUID } from "node:crypto"; import { resolve } from "node:path"; +import { Result, type Result as BetterResult } from "better-result"; import { openDatabase, type DatabaseHandle } from "./db/client.js"; import type { ServerConfig } from "./config.js"; import { @@ -22,6 +23,16 @@ import { workflowRunSourceSchema, workflowRunStatusSchema, } from "./workflow-contracts.js"; +import { + InvalidRunTransitionError, + WorkflowNotFoundError, + WorkflowStoreError, +} from "./workflow-errors.js"; + +export type WorkflowRunTransitionError = + | WorkflowNotFoundError + | InvalidRunTransitionError + | WorkflowStoreError; export interface CreateWorkflowRunInput { name: string; @@ -207,6 +218,15 @@ export class WorkflowStore { return row ? rowToRun(row) : undefined; } + getRunResult( + id: string, + ): BetterResult { + return Result.try({ + try: () => this.getRun(id), + catch: (cause) => new WorkflowStoreError("get_run", cause), + }); + } + listRuns(limit = 50): WorkflowRunRecord[] { const rows = this.database.sqlite .prepare("select * from workflow_runs order by updated_at desc limit ?") @@ -219,31 +239,102 @@ export class WorkflowStore { * 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); + return unwrapRunResult(this.setScriptPathResult(id, scriptPath)); + } + + setScriptPathResult( + id: string, + scriptPath: string, + ): BetterResult { + const current = this.getRunResult(id); + if (current.isErr()) return current; + const run = current.value; + if (!run) return Result.err(new WorkflowNotFoundError(id)); + const updated = Result.try({ + try: () => { + const now = isoNow(); + this.database.sqlite + .prepare( + `UPDATE workflow_runs SET script_path = ?, updated_at = ? WHERE id = ?`, + ) + .run(scriptPath, now, id); + return this.getRun(id); + }, + catch: (cause) => new WorkflowStoreError("set_script_path", cause), + }); + if (updated.isErr()) return updated; + return updated.value + ? Result.ok(updated.value) + : Result.err(new WorkflowNotFoundError(id)); } claimRun(id: string, pid: number): WorkflowRunRecord | undefined { - const now = isoNow(); - const result = this.database.sqlite - .prepare( - `update workflow_runs set - status = 'running', - pid = ?, - heartbeat_at = ?, - started_at = coalesce(started_at, ?), - updated_at = ? - where id = ? and status = 'starting'`, - ) - .run(pid, now, now, now, id); - if (result.changes === 0) return undefined; - return this.getRun(id); + const result = this.claimRunResult(id, pid); + if (result.isOk()) return result.value; + if ( + WorkflowNotFoundError.is(result.error) || + InvalidRunTransitionError.is(result.error) + ) { + return undefined; + } + throw result.error; + } + + claimRunResult( + id: string, + pid: number, + ): BetterResult { + const currentResult = this.getRunResult(id); + if (currentResult.isErr()) return currentResult; + const current = currentResult.value; + if (!current) return Result.err(new WorkflowNotFoundError(id)); + if (current.status !== "starting") { + return Result.err( + new InvalidRunTransitionError({ + runId: id, + from: current.status, + operation: "claim", + }), + ); + } + + const claimed = Result.try({ + try: () => { + const now = isoNow(); + const update = this.database.sqlite + .prepare( + `update workflow_runs set + status = 'running', + pid = ?, + heartbeat_at = ?, + started_at = coalesce(started_at, ?), + updated_at = ? + where id = ? and status = 'starting'`, + ) + .run(pid, now, now, now, id); + return update.changes; + }, + catch: (cause) => new WorkflowStoreError("claim_run", cause), + }); + if (claimed.isErr()) return claimed; + if (claimed.value === 0) { + const latestResult = this.getRunResult(id); + if (latestResult.isErr()) return latestResult; + const latest = latestResult.value; + return latest + ? Result.err( + new InvalidRunTransitionError({ + runId: id, + from: latest.status, + operation: "claim", + }), + ) + : Result.err(new WorkflowNotFoundError(id)); + } + const runResult = this.getRunResult(id); + if (runResult.isErr()) return runResult; + const run = runResult.value; + return run ? Result.ok(run) : Result.err(new WorkflowNotFoundError(id)); } setHeartbeat(id: string, at = isoNow()): void { @@ -255,16 +346,34 @@ export class WorkflowStore { } requestCancel(id: string): WorkflowRunRecord { - const run = this.requireRun(id); - if (TERMINAL_STATUSES.has(run.status)) return run; + return unwrapRunResult(this.requestCancelResult(id)); + } - const now = isoNow(); - this.database.sqlite - .prepare( - `update workflow_runs set cancel_requested = 'true', updated_at = ? where id = ?`, - ) - .run(now, id); - return this.requireRun(id); + requestCancelResult( + id: string, + ): BetterResult { + const current = this.getRunResult(id); + if (current.isErr()) return current; + const run = current.value; + if (!run) return Result.err(new WorkflowNotFoundError(id)); + if (TERMINAL_STATUSES.has(run.status)) return Result.ok(run); + + const updated = Result.try({ + try: () => { + const now = isoNow(); + this.database.sqlite + .prepare( + `update workflow_runs set cancel_requested = 'true', updated_at = ? where id = ?`, + ) + .run(now, id); + return this.getRun(id); + }, + catch: (cause) => new WorkflowStoreError("request_cancel", cause), + }); + if (updated.isErr()) return updated; + return updated.value + ? Result.ok(updated.value) + : Result.err(new WorkflowNotFoundError(id)); } isCancelRequested(id: string): boolean { @@ -272,69 +381,114 @@ export class WorkflowStore { } completeRun(id: string, input: CompleteRunInput = {}): WorkflowRunRecord { - if (input.resultJson !== undefined) assertResultSize(input.resultJson); - const now = isoNow(); - const result = this.database.sqlite - .prepare( - `update workflow_runs set - status = 'completed', - result_json = ?, - completed_at = ?, - updated_at = ?, - error = null, - error_kind = null - where id = ? and status in ('starting', 'running')`, - ) - .run(input.resultJson ?? null, now, now, id); - if (result.changes === 0) { - const run = this.requireRun(id); - if (TERMINAL_STATUSES.has(run.status)) return run; - throw new Error(`Cannot complete workflow run ${id} in status ${run.status}`); - } - return this.requireRun(id); + return unwrapRunResult(this.completeRunResult(id, input)); + } + + completeRunResult( + id: string, + input: CompleteRunInput = {}, + ): BetterResult { + return this.transitionRunResult(id, "complete", () => { + if (input.resultJson !== undefined) assertResultSize(input.resultJson); + const now = isoNow(); + return this.database.sqlite + .prepare( + `update workflow_runs set + status = 'completed', + result_json = ?, + completed_at = ?, + updated_at = ?, + error = null, + error_kind = null + where id = ? and status in ('starting', 'running')`, + ) + .run(input.resultJson ?? null, now, now, id).changes; + }); } failRun(id: string, input: FailRunInput): WorkflowRunRecord { - const now = isoNow(); - const result = this.database.sqlite - .prepare( - `update workflow_runs set - status = 'failed', - error = ?, - error_kind = ?, - completed_at = ?, - updated_at = ? - where id = ? and status in ('starting', 'running')`, - ) - .run(input.error, input.errorKind ?? "internal", now, now, id); - if (result.changes === 0) { - const run = this.requireRun(id); - if (TERMINAL_STATUSES.has(run.status)) return run; - throw new Error(`Cannot fail workflow run ${id} in status ${run.status}`); - } - return this.requireRun(id); + return unwrapRunResult(this.failRunResult(id, input)); + } + + failRunResult( + id: string, + input: FailRunInput, + ): BetterResult { + return this.transitionRunResult(id, "fail", () => { + const now = isoNow(); + return this.database.sqlite + .prepare( + `update workflow_runs set + status = 'failed', + error = ?, + error_kind = ?, + completed_at = ?, + updated_at = ? + where id = ? and status in ('starting', 'running')`, + ) + .run(input.error, input.errorKind ?? "internal", now, now, id).changes; + }); } cancelRun(id: string, error = "cancelled"): WorkflowRunRecord { - const now = isoNow(); - const result = this.database.sqlite - .prepare( - `update workflow_runs set - status = 'cancelled', - error = ?, - error_kind = 'cancelled', - cancel_requested = 'true', - completed_at = ?, - updated_at = ? - where id = ? and status in ('starting', 'running')`, - ) - .run(error, now, now, id); - if (result.changes === 0) { - const run = this.requireRun(id); - if (TERMINAL_STATUSES.has(run.status)) return run; - throw new Error(`Cannot cancel workflow run ${id} in status ${run.status}`); + return unwrapRunResult(this.cancelRunResult(id, error)); + } + + cancelRunResult( + id: string, + error = "cancelled", + ): BetterResult { + return this.transitionRunResult(id, "cancel", () => { + const now = isoNow(); + return this.database.sqlite + .prepare( + `update workflow_runs set + status = 'cancelled', + error = ?, + error_kind = 'cancelled', + cancel_requested = 'true', + completed_at = ?, + updated_at = ? + where id = ? and status in ('starting', 'running')`, + ) + .run(error, now, now, id).changes; + }); + } + + private transitionRunResult( + id: string, + operation: "complete" | "fail" | "cancel", + update: () => number, + ): BetterResult { + const currentResult = this.getRunResult(id); + if (currentResult.isErr()) return currentResult; + const current = currentResult.value; + if (!current) return Result.err(new WorkflowNotFoundError(id)); + if (TERMINAL_STATUSES.has(current.status)) return Result.ok(current); + + const updated = Result.try({ + try: update, + catch: (cause) => new WorkflowStoreError(`${operation}_run`, cause), + }); + if (updated.isErr()) return updated; + if (updated.value === 0) { + const latestResult = this.getRunResult(id); + if (latestResult.isErr()) return latestResult; + const latest = latestResult.value; + if (!latest) return Result.err(new WorkflowNotFoundError(id)); + if (TERMINAL_STATUSES.has(latest.status)) return Result.ok(latest); + return Result.err( + new InvalidRunTransitionError({ + runId: id, + from: latest.status, + operation, + }), + ); } - return this.requireRun(id); + const runResult = this.getRunResult(id); + if (runResult.isErr()) return runResult; + const run = runResult.value; + return run ? Result.ok(run) : Result.err(new WorkflowNotFoundError(id)); } appendEvent(input: AppendWorkflowEventInput): WorkflowEventRecord { @@ -664,3 +818,10 @@ function truncateJson(value: unknown, maxBytes: number): string { const slice = Buffer.from(text, "utf8").subarray(0, budget).toString("utf8"); return JSON.stringify({ truncated: true, preview: slice }); } + +function unwrapRunResult( + result: BetterResult, +): WorkflowRunRecord { + if (result.isErr()) throw result.error; + return result.value; +} diff --git a/src/workflow-tools.ts b/src/workflow-tools.ts index faefce75..7de46bf6 100644 --- a/src/workflow-tools.ts +++ b/src/workflow-tools.ts @@ -6,9 +6,9 @@ import type { ServerConfig } from "./config.js"; import { jsonValueSchema, parseJsonText, type JsonValue } from "./json-types.js"; import type { WorkspaceRegistry } from "./workspaces.js"; import { - persistWorkflowScript, - resolveNamedWorkflowScript, - readWorkflowScriptFile, + persistWorkflowScriptResult, + resolveNamedWorkflowScriptResult, + readWorkflowScriptFileResult, } from "./workflow-files.js"; import { parseWorkflowScript } from "./workflow-script.js"; import { createWorkflowStore } from "./workflow-store.js"; @@ -27,6 +27,13 @@ import { LOCAL_AGENT_PROVIDERS, type LocalAgentProvider, } from "./local-agent-profiles.js"; +import { + InvalidWorkflowInputError, + isWorkflowOperationError, + serializeWorkflowError, + WorkflowNotFoundError, + WorkflowStoredDataError, +} from "./workflow-errors.js"; const WORKFLOW_API_CHEATSHEET = ` Workflow scripts (JS only): @@ -81,7 +88,10 @@ export function registerWorkflowTools( try { const provided = [script, name, resumeFromRunId].filter((v) => v !== undefined); if (provided.length !== 1) { - throw new Error("Provide exactly one of script, name, or resumeFromRunId"); + throw new InvalidWorkflowInputError({ + code: provided.length === 0 ? "missing_source" : "ambiguous_source", + message: "Provide exactly one of script, name, or resumeFromRunId", + }); } let source: string; @@ -93,10 +103,12 @@ export function registerWorkflowTools( if (resumeFromRunId) { const prior = store.getRun(resumeFromRunId); - if (!prior) throw new Error(`Unknown run: ${resumeFromRunId}`); + if (!prior) throw new WorkflowNotFoundError(resumeFromRunId); priorRunId = prior.id; priorScriptPath = prior.scriptPath; - const resolved = await readWorkflowScriptFile(prior.scriptPath); + const resolvedResult = await readWorkflowScriptFileResult(prior.scriptPath); + if (resolvedResult.isErr()) throw resolvedResult.error; + const resolved = resolvedResult.value; source = resolved.source; scriptHash = prior.scriptHash; nameHint = prior.name; @@ -104,16 +116,18 @@ export function registerWorkflowTools( if (args === undefined && prior.argsJson && prior.argsJson !== "null") { try { args = parseJsonText(prior.argsJson); - } catch { - // keep undefined + } catch (cause) { + throw new WorkflowStoredDataError(`${prior.id}.argsJson`, cause); } } } else if (name) { - const resolved = await resolveNamedWorkflowScript({ + const resolvedResult = await resolveNamedWorkflowScriptResult({ name, workspaceRoot: workspace.root, stateDir: config.stateDir, }); + if (resolvedResult.isErr()) throw resolvedResult.error; + const resolved = resolvedResult.value; source = resolved.source; scriptHash = resolved.scriptHash; nameHint = resolved.nameHint; @@ -140,15 +154,19 @@ export function registerWorkflowTools( baseSha, }); - const persisted = - priorScriptPath ?? - (await persistWorkflowScript({ + let persisted = priorScriptPath; + if (!persisted) { + const persistedResult = await persistWorkflowScriptResult({ stateDir: config.stateDir, runId: run.id, source, preferredName: parsed.meta.name || nameHint, - })); - if (!priorScriptPath) store.setScriptPath(run.id, persisted); + }); + if (persistedResult.isErr()) throw persistedResult.error; + persisted = persistedResult.value; + const updated = store.setScriptPathResult(run.id, persisted); + if (updated.isErr()) throw updated.error; + } const cliEntry = fileURLToPath( import.meta.url.replace(/workflow-tools\.(ts|js)$/, "cli.$1"), @@ -158,6 +176,9 @@ export function registerWorkflowTools( const yieldMs = yieldTimeMs ?? 2_000; const page = await yieldEvents(store, run.id, 0, yieldMs); return toolResult(page); + } catch (error) { + if (isWorkflowOperationError(error)) return workflowToolError(error); + throw error; } finally { store.close(); } @@ -187,9 +208,12 @@ export function registerWorkflowTools( async ({ runId, sinceSeq, yieldTimeMs }) => { const store = createWorkflowStore(config); try { - if (!store.getRun(runId)) throw new Error(`Unknown workflow run: ${runId}`); + if (!store.getRun(runId)) throw new WorkflowNotFoundError(runId); const page = await yieldEvents(store, runId, sinceSeq ?? 0, yieldTimeMs ?? 0); return toolResult(page); + } catch (error) { + if (isWorkflowOperationError(error)) return workflowToolError(error); + throw error; } finally { store.close(); } @@ -211,7 +235,9 @@ export function registerWorkflowTools( async ({ runId }) => { const store = createWorkflowStore(config); try { - const run = store.requestCancel(runId); + const requested = store.requestCancelResult(runId); + if (requested.isErr()) throw requested.error; + const run = requested.value; if (run.pid && (run.status === "running" || run.status === "starting")) { try { process.kill(run.pid, "SIGTERM"); @@ -224,6 +250,9 @@ export function registerWorkflowTools( content: [{ type: "text" as const, text: JSON.stringify({ runId, status: latest.status }) }], structuredContent: { runId, status: latest.status }, }; + } catch (error) { + if (isWorkflowOperationError(error)) return workflowToolError(error); + throw error; } finally { store.close(); } @@ -289,6 +318,15 @@ function toolResult(page: { }; } +function workflowToolError(error: Parameters[0]) { + const payload = { error: serializeWorkflowError(error) }; + return { + content: [{ type: "text" as const, text: JSON.stringify(payload, null, 2) }], + structuredContent: payload, + isError: true, + }; +} + function safeJson(text: string): JsonValue { try { return parseJsonText(text); diff --git a/src/workflow-worktrees.ts b/src/workflow-worktrees.ts index 1bbb3414..1ccf6089 100644 --- a/src/workflow-worktrees.ts +++ b/src/workflow-worktrees.ts @@ -2,8 +2,9 @@ import { execFile } from "node:child_process"; import { mkdir, rm } from "node:fs/promises"; import { join } from "node:path"; import { promisify } from "node:util"; +import { Result, type Result as BetterResult } from "better-result"; import type { CreateAgentWorktree, WorkflowWorktreeHandle } from "./workflow-api.js"; -import { WorkflowEngineError } from "./workflow-api.js"; +import { WorktreeOperationError } from "./workflow-errors.js"; const execFileAsync = promisify(execFile); @@ -21,44 +22,65 @@ 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", + const result = await createWorkflowWorktreeResult(host, input); + if (result.isErr()) throw result.error; + return result.value; + }; +} + +export async function createWorkflowWorktreeResult( + host: WorkflowWorktreeHost, + input: Parameters[0], +): Promise> { + return Result.tryPromise({ + try: async () => { + 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 Error("isolation: 'worktree' requires Git on PATH", { cause: error }); + } + throw new Error( + `isolation: 'worktree' requires a Git repository (not found at ${input.workspaceRoot})`, + { cause: error }, ); } - 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 }); - }; + + const baseSha = + input.baseSha ?? + (await git(["rev-parse", "--verify", "HEAD^{commit}"], sourceRoot)).trim(); + + try { + await git(["worktree", "add", "--detach", path, baseSha], sourceRoot); + } catch (error) { + try { + await rm(path, { recursive: true, force: true }); + } catch (cleanupError) { + throw new AggregateError( + [error, cleanupError], + "Failed to create and clean up agent worktree", + ); + } + const message = error instanceof Error ? error.message : String(error); + throw new Error(`Failed to create agent worktree: ${message}`, { cause: error }); + } + + return createHandle({ path, sourceRoot }); + }, + catch: (cause) => + new WorktreeOperationError({ + operation: "create", + runId: input.runId, + callIndex: input.callIndex, + cause, + }), + }); } function createHandle(input: { @@ -68,9 +90,12 @@ function createHandle(input: { return { path: input.path, finalize: async (outcome) => { - const dirty = await isDirty(input.path); + const dirtyResult = await isDirtyResult(input.path); + if (dirtyResult.isErr()) throw dirtyResult.error; + const dirty = dirtyResult.value; if (outcome === "success" && !dirty) { - await removeWorktree(input.sourceRoot, input.path); + const removed = await removeWorktreeResult(input.sourceRoot, input.path); + if (removed.isErr()) throw removed.error; return { dirty: false, removed: true }; } // Preserve dirty or failed worktrees for diagnosis. @@ -80,29 +105,62 @@ function createHandle(input: { } 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; - } + const result = await isDirtyResult(worktreePath); + return result.isOk() ? result.value : true; +} + +export async function isDirtyResult( + worktreePath: string, +): Promise> { + return Result.tryPromise({ + try: async () => { + const status = (await git(["status", "--porcelain=v1"], worktreePath)).trim(); + return status.length > 0; + }, + catch: (cause) => + new WorktreeOperationError({ + operation: "inspect", + path: worktreePath, + cause, + }), + }); } 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 - } - } + const result = await removeWorktreeResult(sourceRoot, worktreePath); + if (result.isErr()) throw result.error; +} + +export async function removeWorktreeResult( + sourceRoot: string, + worktreePath: string, +): Promise> { + return Result.tryPromise({ + try: async () => { + try { + await git(["worktree", "remove", "--force", worktreePath], sourceRoot); + } catch (removeError) { + await rm(worktreePath, { recursive: true, force: true }); + try { + await git(["worktree", "prune"], sourceRoot); + } catch (pruneError) { + throw new AggregateError( + [removeError, pruneError], + "Worktree directory was removed but Git metadata pruning failed", + ); + } + } + }, + catch: (cause) => + new WorktreeOperationError({ + operation: "remove", + path: worktreePath, + cause, + }), + }); } export async function resolveWorkspaceHead(workspaceRoot: string): Promise {