From cc09a097dc28f2bc81348aee8b5c76e598e93b8b Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:05:09 +0530 Subject: [PATCH 1/7] refactor(workflow): house WorkflowEngineError in workflow-errors Move the engine domain error class next to other workflow failures and add result_too_large for upcoming oversized-return policy. Re-export from workflow-api for existing imports. --- src/workflow-api.ts | 22 +++------------------- src/workflow-errors.ts | 21 +++++++++++++++++++++ 2 files changed, 24 insertions(+), 19 deletions(-) diff --git a/src/workflow-api.ts b/src/workflow-api.ts index cd0d60da..3a2394c0 100644 --- a/src/workflow-api.ts +++ b/src/workflow-api.ts @@ -24,6 +24,9 @@ import { type WorkflowMeta, } from "./workflow-types.js"; import { agentOptsSchema } from "./workflow-contracts.js"; +import { WorkflowEngineError } from "./workflow-errors.js"; + +export { WorkflowEngineError } from "./workflow-errors.js"; // --------------------------------------------------------------------------- // Host deps (injected by engine; fakes OK in tests) @@ -182,25 +185,6 @@ export interface WorkflowApi extends WorkflowSandboxApi { getNestDepth(): number; } -export class WorkflowEngineError extends Error { - constructor( - readonly kind: - | "cancelled" - | "provider_unavailable" - | "no_provider" - | "profile" - | "nest_depth" - | "worktree" - | "schema" - | "path" - | "internal", - message: string, - ) { - super(message); - this.name = "WorkflowEngineError"; - } -} - // --------------------------------------------------------------------------- // Semaphore // --------------------------------------------------------------------------- diff --git a/src/workflow-errors.ts b/src/workflow-errors.ts index 19e18b2c..26eca02e 100644 --- a/src/workflow-errors.ts +++ b/src/workflow-errors.ts @@ -10,6 +10,27 @@ import type { WorkflowRunStatus, } from "./workflow-types.js"; +/** Domain failures inside agent()/sandbox orchestration (throw, not Result). */ +export class WorkflowEngineError extends Error { + constructor( + readonly kind: + | "cancelled" + | "provider_unavailable" + | "no_provider" + | "profile" + | "nest_depth" + | "worktree" + | "schema" + | "path" + | "result_too_large" + | "internal", + message: string, + ) { + super(message); + this.name = "WorkflowEngineError"; + } +} + export class InvalidWorkflowInputError extends TaggedError( "InvalidWorkflowInputError", )<{ From 070fbd66dcfa6e8b96b49bf00ca353b60435c616 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:06:27 +0530 Subject: [PATCH 2/7] fix(workflow): fail agent calls when return value exceeds replay budget Completed calls must remain resume-safe. Oversized return values and structured JSON now throw result_too_large instead of silently dropping persisted replay data. --- src/workflow-api.ts | 40 ++++++++++++++++++++++++------------- src/workflow-engine.test.ts | 35 +++++++++++++++++++------------- src/workflow-engine.ts | 2 +- 3 files changed, 48 insertions(+), 29 deletions(-) diff --git a/src/workflow-api.ts b/src/workflow-api.ts index 3a2394c0..86474f7f 100644 --- a/src/workflow-api.ts +++ b/src/workflow-api.ts @@ -450,6 +450,11 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { throwIfCancelled(deps); + const returnValueJson = serializeReplayValueOrThrow(returnValue); + if (structuredJson !== undefined) { + assertStructuredJsonBudget(structuredJson); + } + let dirty: boolean | undefined; if (worktree) { const finalized = await worktree.finalize("success"); @@ -473,8 +478,8 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { runId: deps.runId, callIndex: index, responseText: truncate(result.finalResponse, WORKFLOW_LIMITS.responseTextBytes), - structuredJson: boundedStructuredJson(structuredJson), - returnValueJson: serializeReplayValue(returnValue), + structuredJson, + returnValueJson, providerSessionId: result.providerSessionId, dirty, worktreePath, @@ -776,23 +781,30 @@ function truncate(text: string, maxBytes: number): string { return `${text.slice(0, end)}${marker}`; } -function boundedStructuredJson(value: string | undefined): string | undefined { - if (value === undefined) return undefined; - return Buffer.byteLength(value, "utf8") <= WORKFLOW_LIMITS.structuredJsonBytes - ? value - : undefined; +function assertStructuredJsonBudget(value: string): void { + if (Buffer.byteLength(value, "utf8") <= WORKFLOW_LIMITS.structuredJsonBytes) return; + throw new WorkflowEngineError( + "result_too_large", + `agent() structured result exceeds ${WORKFLOW_LIMITS.structuredJsonBytes} bytes; return a smaller object or write large artifacts to disk and return paths`, + ); } -function serializeReplayValue(value: unknown): string | undefined { +function serializeReplayValueOrThrow(value: unknown): string | undefined { + let json: string | undefined; try { - const json = JSON.stringify(value); - if (json === undefined) return undefined; - return Buffer.byteLength(json, "utf8") <= WORKFLOW_LIMITS.replayValueJsonBytes - ? json - : undefined; + json = JSON.stringify(value); } catch { - return undefined; + throw new WorkflowEngineError( + "result_too_large", + "agent() return value is not JSON-serializable and cannot be replayed", + ); } + if (json === undefined) return undefined; + if (Buffer.byteLength(json, "utf8") <= WORKFLOW_LIMITS.replayValueJsonBytes) return json; + throw new WorkflowEngineError( + "result_too_large", + `agent() return value exceeds ${WORKFLOW_LIMITS.replayValueJsonBytes} bytes replay budget; return a smaller summary or write large artifacts to disk and return paths`, + ); } /** Minimal JSON extract for schema path until Ajv module lands. */ diff --git a/src/workflow-engine.test.ts b/src/workflow-engine.test.ts index fc4bf57d..7378b046 100644 --- a/src/workflow-engine.test.ts +++ b/src/workflow-engine.test.ts @@ -473,7 +473,7 @@ import type { LocalAgentProfile } from "./local-agent-profiles.js"; } // --------------------------------------------------------------------------- -// oversized exact replay values do not fail the live call +// oversized exact replay values fail the live call (completed ⇒ replayable) // --------------------------------------------------------------------------- { const dir = await mkdtemp(join(tmpdir(), "wf-replay-size-")); @@ -498,14 +498,19 @@ import type { LocalAgentProfile } from "./local-agent-profiles.js"; runProvider: async () => ({ finalResponse: response }), }); - assert.equal(await api.agent("large"), response); - assert.equal(store.getAgentCall(run.id, 0)?.returnValueJson, undefined); + await assert.rejects( + () => api.agent("large"), + (error: unknown) => + error instanceof WorkflowEngineError && error.kind === "result_too_large", + ); + assert.equal(store.getAgentCall(run.id, 0)?.status, "failed"); + assert.equal(store.getAgentCall(run.id, 0)?.errorKind, "result_too_large"); store.close(); await rm(dir, { recursive: true, force: true }); } // --------------------------------------------------------------------------- -// oversized structured results remain intact without persisting invalid JSON +// oversized structured results fail the live call // --------------------------------------------------------------------------- { const dir = await mkdtemp(join(tmpdir(), "wf-structured-size-")); @@ -533,17 +538,19 @@ import type { LocalAgentProfile } from "./local-agent-profiles.js"; }), }); - assert.deepEqual( - await api.agent("large structured", { - schema: { - type: "object", - properties: { big: { type: "string" } }, - required: ["big"], - }, - }), - { big }, + await assert.rejects( + () => + api.agent("large structured", { + schema: { + type: "object", + properties: { big: { type: "string" } }, + required: ["big"], + }, + }), + (error: unknown) => + error instanceof WorkflowEngineError && error.kind === "result_too_large", ); - assert.equal(store.getAgentCall(run.id, 0)?.structuredJson, undefined); + assert.equal(store.getAgentCall(run.id, 0)?.status, "failed"); store.close(); await rm(dir, { recursive: true, force: true }); } diff --git a/src/workflow-engine.ts b/src/workflow-engine.ts index 72081b53..5fc08dac 100644 --- a/src/workflow-engine.ts +++ b/src/workflow-engine.ts @@ -151,7 +151,7 @@ export async function executeWorkflow( export function mapEngineErrorKind(error: unknown): WorkflowErrorKind { if (error instanceof WorkflowEngineError) { - return error.kind; + return error.kind as WorkflowErrorKind; } if (isWorkflowOperationError(error)) { return workflowErrorKind(error); From 0c8aceff4865726335166f615dfe9bc696d9f7cc Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:07:39 +0530 Subject: [PATCH 3/7] refactor(workflow): drop path throw shims; Result-only file APIs Remove WorkflowPathError wrappers and dual throw helpers so CLI/worker callers handle typed TaggedErrors directly. Delete unused contentHash and dirname helpers. --- src/workflow-cli.ts | 21 ++++--- src/workflow-files.test.ts | 121 ++++++++++++++++++++++--------------- src/workflow-files.ts | 85 +------------------------- 3 files changed, 85 insertions(+), 142 deletions(-) diff --git a/src/workflow-cli.ts b/src/workflow-cli.ts index bfd00367..be4f667f 100644 --- a/src/workflow-cli.ts +++ b/src/workflow-cli.ts @@ -17,9 +17,9 @@ import { executeWorkflow, mapEngineErrorKind } from "./workflow-engine.js"; import { parseWorkflowArgFlagsResult, persistWorkflowScriptResult, - readProjectWorkflowScriptFile, + readProjectWorkflowScriptFileResult, readWorkflowScriptFileResult, - resolveNamedWorkflowScript, + resolveNamedWorkflowScriptResult, resolveWorkflowScriptFromPathOrNameResult, } from "./workflow-files.js"; import { createWorkflowReplay } from "./workflow-replay.js"; @@ -415,19 +415,20 @@ export async function runWorkflowWorker( }, resolveNestedSource: async (ref) => { if (typeof ref === "string") { - const named = await resolveNamedWorkflowScript({ + const named = await resolveNamedWorkflowScriptResult({ name: ref, workspaceRoot: claimed.workspaceRoot, stateDir: config.stateDir, }); - return named.source; + if (named.isErr()) throw named.error; + return named.value.source; } - return ( - await readProjectWorkflowScriptFile({ - scriptPath: ref.scriptPath, - workspaceRoot: claimed.workspaceRoot, - }) - ).source; + const nested = await readProjectWorkflowScriptFileResult({ + scriptPath: ref.scriptPath, + workspaceRoot: claimed.workspaceRoot, + }); + if (nested.isErr()) throw nested.error; + return nested.value.source; }, }); diff --git a/src/workflow-files.test.ts b/src/workflow-files.test.ts index 4cd1ecd7..37f445e5 100644 --- a/src/workflow-files.test.ts +++ b/src/workflow-files.test.ts @@ -3,17 +3,20 @@ import { mkdir, mkdtemp, rm, symlink, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { - parseWorkflowArgFlags, - persistWorkflowScript, - readProjectWorkflowScriptFile, - resolveNamedWorkflowScript, - resolveWorkflowScriptFromPathOrName, - WorkflowPathError, + parseWorkflowArgFlagsResult, + persistWorkflowScriptResult, + readProjectWorkflowScriptFileResult, + resolveNamedWorkflowScriptResult, + resolveWorkflowScriptFromPathOrNameResult, } from "./workflow-files.js"; +import { + InvalidWorkflowInputError, + NamedWorkflowNotFoundError, +} from "./workflow-errors.js"; import { hashSource } from "./workflow-script.js"; { - const { args, rest } = parseWorkflowArgFlags([ + const parsed = parseWorkflowArgFlagsResult([ "--arg", "n=1", "--arg", @@ -21,55 +24,66 @@ import { hashSource } from "./workflow-script.js"; "--follow", "extra", ]); - assert.deepEqual(args, { n: 1, files: ["a.ts"] }); - assert.deepEqual(rest, ["--follow", "extra"]); + assert.equal(parsed.isOk(), true); + if (parsed.isOk()) { + assert.deepEqual(parsed.value.args, { n: 1, files: ["a.ts"] }); + assert.deepEqual(parsed.value.rest, ["--follow", "extra"]); + } } { const dir = await mkdtemp(join(tmpdir(), "wf-files-")); - const path = await persistWorkflowScript({ + const persisted = await persistWorkflowScriptResult({ stateDir: dir, runId: "wfr_test", source: "export const meta = { name: 'x', description: 'd' }\nreturn 1\n", preferredName: "demo", }); + assert.equal(persisted.isOk(), true); + if (!persisted.isOk()) throw persisted.error; + const path = persisted.value; assert.match(path.replaceAll("\\", "/"), /workflow-scripts\/wfr_test\/demo\.js$/); - const file = await resolveWorkflowScriptFromPathOrName({ + const file = await resolveWorkflowScriptFromPathOrNameResult({ file: path, workspaceRoot: dir, }); - assert.equal(file.origin, "file"); - assert.equal(file.scriptHash, hashSource(file.source)); + assert.equal(file.isOk(), true); + if (!file.isOk()) throw file.error; + assert.equal(file.value.origin, "file"); + assert.equal(file.value.scriptHash, hashSource(file.value.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({ + const named = await resolveNamedWorkflowScriptResult({ name: "named", workspaceRoot: dir, }); - assert.equal(named.origin, "named"); - assert.match(named.source, /named/); - assert.equal( - ( - await readProjectWorkflowScriptFile({ - scriptPath: join(dir, ".devspace", "workflows", "named.js"), - workspaceRoot: dir, - }) - ).nameHint, - "named", - ); - await assert.rejects( - () => - readProjectWorkflowScriptFile({ - scriptPath: path, - workspaceRoot: dir, - }), - /must be inside/, - ); + assert.equal(named.isOk(), true); + if (!named.isOk()) throw named.error; + assert.equal(named.value.origin, "named"); + assert.match(named.value.source, /named/); + + const projectRead = await readProjectWorkflowScriptFileResult({ + scriptPath: join(dir, ".devspace", "workflows", "named.js"), + workspaceRoot: dir, + }); + assert.equal(projectRead.isOk(), true); + if (!projectRead.isOk()) throw projectRead.error; + assert.equal(projectRead.value.nameHint, "named"); + + const outsideProject = await readProjectWorkflowScriptFileResult({ + scriptPath: path, + workspaceRoot: dir, + }); + assert.equal(outsideProject.isErr(), true); + if (outsideProject.isErr()) { + assert.equal(InvalidWorkflowInputError.is(outsideProject.error), true); + assert.match(outsideProject.error.message, /must be inside/); + } if (process.platform !== "win32") { const outside = await mkdtemp(join(tmpdir(), "wf-files-outside-")); @@ -83,14 +97,15 @@ import { hashSource } from "./workflow-script.js"; outsideScript, join(dir, ".devspace", "workflows", "escape.js"), ); - await assert.rejects( - () => - readProjectWorkflowScriptFile({ - scriptPath: "escape.js", - workspaceRoot: dir, - }), - /resolves outside/, - ); + const escaped = await readProjectWorkflowScriptFileResult({ + scriptPath: "escape.js", + workspaceRoot: dir, + }); + assert.equal(escaped.isErr(), true); + if (escaped.isErr()) { + assert.equal(InvalidWorkflowInputError.is(escaped.error), true); + assert.match(escaped.error.message, /resolves outside/); + } } finally { await rm(outside, { recursive: true, force: true }); } @@ -101,15 +116,23 @@ import { hashSource } from "./workflow-script.js"; join(dir, "workflows", "legacy.js"), "export const meta = { name: 'legacy', description: 'd' }\nreturn 3\n", ); - await assert.rejects( - () => resolveNamedWorkflowScript({ name: "legacy", workspaceRoot: dir }), - WorkflowPathError, - ); + const legacy = await resolveNamedWorkflowScriptResult({ + name: "legacy", + workspaceRoot: dir, + }); + assert.equal(legacy.isErr(), true); + if (legacy.isErr()) { + assert.equal(NamedWorkflowNotFoundError.is(legacy.error), true); + } - await assert.rejects( - () => resolveNamedWorkflowScript({ name: "missing", workspaceRoot: dir }), - WorkflowPathError, - ); + const missing = await resolveNamedWorkflowScriptResult({ + name: "missing", + workspaceRoot: dir, + }); + assert.equal(missing.isErr(), true); + if (missing.isErr()) { + assert.equal(NamedWorkflowNotFoundError.is(missing.error), true); + } await rm(dir, { recursive: true, force: true }); } diff --git a/src/workflow-files.ts b/src/workflow-files.ts index 5aae5fe5..90f3016d 100644 --- a/src/workflow-files.ts +++ b/src/workflow-files.ts @@ -1,6 +1,6 @@ -import { createHash, randomBytes } from "node:crypto"; +import { randomBytes } from "node:crypto"; import { mkdir, readFile, realpath, writeFile } from "node:fs/promises"; -import { basename, dirname, extname, isAbsolute, join, resolve } from "node:path"; +import { basename, 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"; @@ -13,13 +13,6 @@ import { } from "./workflow-errors.js"; import { isPathInsideRoot } from "./roots.js"; -export class WorkflowPathError extends Error { - constructor(message: string) { - super(message); - this.name = "WorkflowPathError"; - } -} - export interface ResolvedWorkflowScript { source: string; scriptPath: string; @@ -38,17 +31,6 @@ export type WorkflowFileResolveError = * 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 result = await persistWorkflowScriptResult(input); - if (result.isErr()) throw result.error; - return result.value; -} - export async function persistWorkflowScriptResult(input: { stateDir: string; runId: string; @@ -70,12 +52,6 @@ export async function persistWorkflowScriptResult(input: { }); } -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> { @@ -98,15 +74,6 @@ export async function readWorkflowScriptFileResult( }); } -export async function readProjectWorkflowScriptFile(input: { - scriptPath: string; - workspaceRoot: string; -}): Promise { - const result = await readProjectWorkflowScriptFileResult(input); - if (result.isErr()) throwPathCompatibilityError(result.error); - return result.value; -} - /** Resolve an explicit nested script only inside `/.devspace/workflows`. */ export async function readProjectWorkflowScriptFileResult(input: { scriptPath: string; @@ -158,16 +125,6 @@ export async function readProjectWorkflowScriptFileResult(input: { * 1. `/.devspace/workflows/.js` * 2. `/workflows/.js` (if stateDir provided) */ -export async function resolveNamedWorkflowScript(input: { - name: string; - 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; @@ -199,17 +156,6 @@ export async function resolveNamedWorkflowScriptResult(input: { return Result.err(new NamedWorkflowNotFoundError(name, candidates)); } -export async function resolveWorkflowScriptFromPathOrName(input: { - file?: string; - name?: string; - 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; @@ -245,15 +191,6 @@ export async function resolveWorkflowScriptFromPathOrNameResult(input: { ); } -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< @@ -314,18 +251,6 @@ function sanitizeSegment(value: string): string { .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); -} - function isFileNotFound(error: unknown): boolean { return Boolean( error && @@ -334,9 +259,3 @@ function isFileNotFound(error: unknown): boolean { (error as { code?: unknown }).code === "ENOENT", ); } - -function throwPathCompatibilityError(error: Error): never { - const compatible = new WorkflowPathError(error.message); - compatible.cause = error; - throw compatible; -} From 8cffabe538a27d65a4f23d0fd619f65977884227 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:10:56 +0530 Subject: [PATCH 4/7] fix(workflow): track phase in sandbox child via AsyncLocalStorage Host ALS on IPC handlers does not follow concurrent script chains. Inject phase ALS into the child vm, stamp agent opts.phase and log payloads across the bridge, and cover concurrent phase via sandbox e2e. --- src/workflow-api.ts | 22 +++++++++++-- src/workflow-sandbox-child.ts | 57 ++++++++++++++++++++++++++------ src/workflow-sandbox.test.ts | 62 +++++++++++++++++++++++++++++++++-- 3 files changed, 126 insertions(+), 15 deletions(-) diff --git a/src/workflow-api.ts b/src/workflow-api.ts index 86474f7f..071b0d59 100644 --- a/src/workflow-api.ts +++ b/src/workflow-api.ts @@ -612,6 +612,8 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { if (typeof title !== "string" || !title.trim()) { throw new WorkflowEngineError("internal", "phase(title) requires a non-empty string"); } + // In-process tests still use host ALS. Sandbox scripts track phase in the + // child and inject opts.phase / log payloads across IPC. phaseAls.enterWith(title); deps.journal.appendEvent({ runId: deps.runId, @@ -622,11 +624,27 @@ export function createWorkflowApi(deps: WorkflowApiDeps): WorkflowApi { }; const log = (...args: unknown[]): void => { - const message = args.map(String).join(" "); + let message: string; + let phaseTitle = phaseAls.getStore(); + if ( + args.length === 1 && + args[0] && + typeof args[0] === "object" && + !Array.isArray(args[0]) && + "message" in (args[0] as object) + ) { + const payload = args[0] as { message?: unknown; phase?: unknown }; + message = String(payload.message ?? ""); + if (typeof payload.phase === "string" && payload.phase.trim()) { + phaseTitle = payload.phase; + } + } else { + message = args.map(String).join(" "); + } deps.journal.appendEvent({ runId: deps.runId, type: "log", - phase: phaseAls.getStore(), + phase: phaseTitle, data: { message: truncate(message, WORKFLOW_LIMITS.eventDataJsonBytes) }, }); }; diff --git a/src/workflow-sandbox-child.ts b/src/workflow-sandbox-child.ts index f26c4de1..ffe90e2b 100644 --- a/src/workflow-sandbox-child.ts +++ b/src/workflow-sandbox-child.ts @@ -1,8 +1,12 @@ +import { AsyncLocalStorage } from "node:async_hooks"; import vm from "node:vm"; import { parseWorkflowScript } from "./workflow-script.js"; import type { JsonValue } from "./json-types.js"; import { WORKFLOW_MAX_ITEMS } from "./workflow-types.js"; +/** Phase context for concurrent script chains (must live in the child process). */ +const phaseAls = new AsyncLocalStorage(); + type SandboxMethod = "agent" | "workflow" | "phase" | "log"; interface StartMessage { @@ -79,7 +83,17 @@ async function execute(message: StartMessage): Promise { }); }; - const context = vm.createContext({ __workflowBridge: bridge }); + const context = vm.createContext({ + __workflowBridge: bridge, + __workflowPhaseAls: { + enterWith(title: string) { + phaseAls.enterWith(title); + }, + getStore() { + return phaseAls.getStore(); + }, + }, + }); installContextApi(context, message); const factory = parsed.script.runInContext(context, { timeout: 5_000, @@ -160,6 +174,8 @@ function installContextApi(context: vm.Context, message: StartMessage): void { if (typeof input?.stack === "string") error.stack = input.stack; return error; }; + const phaseAls = globalThis.__workflowPhaseAls; + delete globalThis.__workflowPhaseAls; const call = (method, callArgs) => new Promise((resolve, reject) => { bridge(method, callArgs).then( (payloadJson) => { @@ -176,15 +192,37 @@ function installContextApi(context: vm.Context, message: StartMessage): void { () => reject(new WorkflowEngineError("internal", "Workflow bridge call failed")), ); }); - const agent = (...callArgs) => call("agent", callArgs); + // Inject current ALS phase so host journal/agent rows stay correct even though + // host phase() only records events (host ALS is not on the script chain). + const agent = (prompt, opts = {}) => { + const inherited = typeof phaseAls?.getStore === "function" ? phaseAls.getStore() : undefined; + const nextOpts = + opts && typeof opts === "object" + ? { + ...opts, + phase: + typeof opts.phase === "string" && opts.phase.trim() + ? opts.phase + : inherited, + } + : inherited + ? { phase: inherited } + : opts; + return call("agent", [prompt, nextOpts]); + }; const workflow = (...callArgs) => call("workflow", callArgs); const phase = (title) => { if (typeof title !== "string" || !title.trim()) { throw new WorkflowEngineError("internal", "phase(title) requires a non-empty string"); } + phaseAls.enterWith(title); return bridge("phase", [title]); }; - const log = (...callArgs) => bridge("log", callArgs); + const emitLog = (...callArgs) => { + const message = callArgs.map(stringifyConsoleArg).join(" "); + const inherited = typeof phaseAls?.getStore === "function" ? phaseAls.getStore() : undefined; + return bridge("log", [{ message, phase: inherited }]); + }; const parallel = async (tasks) => { if (!Array.isArray(tasks)) { throw new WorkflowEngineError("internal", "parallel(thunks) requires an array of functions"); @@ -244,20 +282,19 @@ function installContextApi(context: vm.Context, message: StartMessage): void { if (typeof value === "string") return value; try { return JSON.stringify(value); } catch { return String(value); } }; - const consoleLine = (...callArgs) => log(callArgs.map(stringifyConsoleArg).join(" ")); const console = Object.freeze({ - log: consoleLine, - warn: consoleLine, - error: consoleLine, - info: consoleLine, - debug: consoleLine, + log: emitLog, + warn: emitLog, + error: emitLog, + info: emitLog, + debug: emitLog, }); Object.defineProperties(globalThis, { agent: { value: Object.freeze(agent), writable: false, configurable: false }, workflow: { value: Object.freeze(workflow), writable: false, configurable: false }, phase: { value: Object.freeze(phase), writable: false, configurable: false }, - log: { value: Object.freeze(log), writable: false, configurable: false }, + log: { value: Object.freeze(emitLog), writable: false, configurable: false }, parallel: { value: Object.freeze(parallel), writable: false, configurable: false }, pipeline: { value: Object.freeze(pipeline), writable: false, configurable: false }, args: { value: Object.freeze(args), writable: false, configurable: false }, diff --git a/src/workflow-sandbox.test.ts b/src/workflow-sandbox.test.ts index 35f7de27..e33dd618 100644 --- a/src/workflow-sandbox.test.ts +++ b/src/workflow-sandbox.test.ts @@ -8,13 +8,26 @@ import { import type { WorkflowSandboxApi } from "./workflow-sandbox.js"; import { runWorkflowSandbox, WorkflowDeterminismError } from "./workflow-sandbox.js"; -function api(meta: WorkflowMeta, logs?: string[]): WorkflowSandboxApi { +function api( + meta: WorkflowMeta, + logs?: string[], + hooks?: { + agent?: WorkflowSandboxApi["agent"]; + phaseTitles?: string[]; + }, +): WorkflowSandboxApi { return { - agent: async () => "", + agent: hooks?.agent ?? (async () => ""), parallel: async () => [], pipeline: async () => [], - phase: () => {}, + phase: (title: string) => { + hooks?.phaseTitles?.push(title); + }, log: (msg: unknown) => { + if (msg && typeof msg === "object" && !Array.isArray(msg) && "message" in msg) { + logs?.push(String((msg as { message: unknown }).message)); + return; + } logs?.push(String(msg)); }, args: undefined as unknown, @@ -226,4 +239,47 @@ return agent.constructor('return process.version')() ); } +// Child-owned phase ALS: concurrent chains inject distinct opts.phase over IPC. +{ + const seen: Array<{ prompt: string; phase?: string }> = []; + const phaseTitles: string[] = []; + const hostApi = api( + { name: "phase-ipc", description: "d" }, + undefined, + { + phaseTitles, + agent: async (prompt: string, opts?: { phase?: string }) => { + seen.push({ prompt, phase: opts?.phase }); + await new Promise((r) => setTimeout(r, 20)); + return `ok:${prompt}`; + }, + }, + ); + // Host parallel is unused; child implements parallel. Agent is bridged. + const result = await runWorkflowSandbox({ + parsed: parseWorkflowScript(` +export const meta = { name: 'phase-ipc', description: 'd' } +return await parallel([ + async () => { + phase('A') + log('in-a') + return await agent('from-a') + }, + async () => { + phase('B') + log('in-b') + return await agent('from-b') + }, +]) +`), + api: hostApi, + }); + assert.deepEqual(result, ["ok:from-a", "ok:from-b"]); + assert.deepEqual(new Set(phaseTitles), new Set(["A", "B"])); + const a = seen.find((row) => row.prompt === "from-a"); + const b = seen.find((row) => row.prompt === "from-b"); + assert.equal(a?.phase, "A"); + assert.equal(b?.phase, "B"); +} + console.log("workflow-sandbox.test.ts: ok"); From 94e198f323e79190b727b47b17923d23659140b8 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:16:36 +0530 Subject: [PATCH 5/7] refactor(workflow): extract shared launch and worker modules CLI and MCP both start runs through launchWorkflowRun. Detached worker spawn and execution live in workflow-worker; live provider ordering is shared via workflow-providers. Adds launch unit coverage without spawn. --- package.json | 2 +- src/workflow-cli.ts | 376 +++++++++--------------------------- src/workflow-engine.test.ts | 12 +- src/workflow-launch.test.ts | 66 +++++++ src/workflow-launch.ts | 276 ++++++++++++++++++++++++++ src/workflow-providers.ts | 12 ++ src/workflow-tools.ts | 170 +++++----------- src/workflow-worker.ts | 191 ++++++++++++++++++ 8 files changed, 692 insertions(+), 413 deletions(-) create mode 100644 src/workflow-launch.test.ts create mode 100644 src/workflow-launch.ts create mode 100644 src/workflow-providers.ts create mode 100644 src/workflow-worker.ts diff --git a/package.json b/package.json index f81a50ba..6e1aad7a 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/open-workspace-capabilities.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-capabilities.test.ts && tsx src/local-agent-catalog.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-resolution.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-lifecycle.test.ts && tsx src/workflow-view.test.ts && tsx src/workflow-ui.test.ts && tsx src/workflow-tui.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/open-workspace-capabilities.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-capabilities.test.ts && tsx src/local-agent-catalog.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-resolution.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-lifecycle.test.ts && tsx src/workflow-view.test.ts && tsx src/workflow-ui.test.ts && tsx src/workflow-tui.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-launch.test.ts && tsx src/workflow-replay.test.ts && tsx src/workflow-schema.test.ts", "typecheck": "tsc -p tsconfig.json --noEmit" }, "keywords": [], diff --git a/src/workflow-cli.ts b/src/workflow-cli.ts index be4f667f..6100332c 100644 --- a/src/workflow-cli.ts +++ b/src/workflow-cli.ts @@ -1,53 +1,34 @@ -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 { parseJsonText, type JsonObject, type JsonValue } from "./json-types.js"; -import { runLocalAgentProviderResult } from "./local-agent-adapters.js"; -import { getLocalAgentProviderAvailabilitySnapshot } from "./local-agent-availability.js"; -import { - isLocalAgentProvider, - loadLocalAgentProfiles, - LOCAL_AGENT_PROVIDERS, - type LocalAgentProvider, -} from "./local-agent-profiles.js"; -import { executeWorkflow, mapEngineErrorKind } from "./workflow-engine.js"; -import { - parseWorkflowArgFlagsResult, - persistWorkflowScriptResult, - readProjectWorkflowScriptFileResult, - readWorkflowScriptFileResult, - resolveNamedWorkflowScriptResult, - resolveWorkflowScriptFromPathOrNameResult, -} from "./workflow-files.js"; -import { createWorkflowReplay } from "./workflow-replay.js"; +import { parseWorkflowArgFlagsResult } from "./workflow-files.js"; import { cancelWorkflowRun, reapStaleWorkflows, } from "./workflow-lifecycle.js"; -import { parseWorkflowScript } from "./workflow-script.js"; import { createWorkflowStore, type WorkflowStore } from "./workflow-store.js"; import { - WORKFLOW_HEARTBEAT_MS, WORKFLOW_LIMITS, - resolveWorkflowConcurrency, type WorkflowEventRecord, type WorkflowAgentCallRecord, type WorkflowRunRecord, - type WorkflowRunSource, } from "./workflow-types.js"; import { parseWorkflowEventPayload } from "./workflow-contracts.js"; import { InvalidWorkflowInputError, WorkflowNotFoundError, - WorkflowStoredDataError, } from "./workflow-errors.js"; import { - createWorkflowWorktreeFactory, - resolveWorkspaceHead, -} from "./workflow-worktrees.js"; + launchWorkflowRun, + type LaunchWorkflowSource, +} from "./workflow-launch.js"; +import { + runWorkflowWorker, + spawnWorkflowWorker, + spawnWorkflowWorkerFromCli, +} from "./workflow-worker.js"; + +export { runWorkflowWorker, spawnWorkflowWorker, spawnWorkflowWorkerFromCli }; export async function runWorkflowCommand( args: string[], @@ -55,9 +36,11 @@ export async function runWorkflowCommand( ): Promise { const [subcommand, ...rest] = args; if (!config.workflows) { - throw new Error( - "Dynamic workflows are disabled. Set DEVSPACE_WORKFLOWS=1 to enable the experimental feature.", - ); + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: + "Dynamic workflows are disabled. Set DEVSPACE_WORKFLOWS=1 to enable the experimental feature.", + }); } switch (subcommand) { case "run": @@ -94,7 +77,10 @@ export async function runWorkflowCommand( printWorkflowHelp(); return; default: - throw new Error(`Unknown workflow command: ${subcommand}`); + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: `Unknown workflow command: ${subcommand}`, + }); } } @@ -140,107 +126,71 @@ async function runWorkflowRun(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]"); + if (!runId) { + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: "Usage: devspace workflow status [--follow]", + }); + } const store = createWorkflowStore(config); try { @@ -264,7 +214,12 @@ async function runWorkflowStatus(args: string[], config: ServerConfig): Promise< async function runWorkflowCancel(args: string[], config: ServerConfig): Promise { const runId = args[0]; - if (!runId) throw new Error("Usage: devspace workflow cancel "); + if (!runId) { + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: "Usage: devspace workflow cancel ", + }); + } const store = createWorkflowStore(config); try { reapStaleWorkflows(store); @@ -291,7 +246,12 @@ async function runWorkflowList(config: ServerConfig): Promise { async function runWorkflowCalls(args: string[], config: ServerConfig): Promise { const runId = args[0]; - if (!runId) throw new Error("Usage: devspace workflow calls "); + if (!runId) { + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: "Usage: devspace workflow calls ", + }); + } const store = createWorkflowStore(config); try { if (!store.getRun(runId)) throw new WorkflowNotFoundError(runId); @@ -310,181 +270,27 @@ async function runWorkflowCall(args: string[], config: ServerConfig): Promise "); + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: "Usage: devspace workflow call ", + }); } const store = createWorkflowStore(config); try { if (!store.getRun(runId)) throw new WorkflowNotFoundError(runId); const call = store.getAgentCall(runId, callIndex); - if (!call) throw new Error(`Unknown workflow agent call: ${runId}#${callIndex}`); - console.log(JSON.stringify(formatCallDetail(call), null, 2)); - } 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 claim = store.claimRunResult(runId, process.pid); - if (claim.isErr()) { - store.close(); - throw claim.error; - } - const claimed = claim.value; - - 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 availableProviders = resolveAvailableProviders(); - const agentProfiles = await loadLocalAgentProfiles(config, claimed.workspaceRoot); - const concurrency = resolveWorkflowConcurrency( - parsed.meta.concurrency, - availableParallelism(), - ); - - let argsValue: JsonValue | undefined; - try { - argsValue = parseJsonText(claimed.argsJson); - if (argsValue === null) argsValue = undefined; - } catch (cause) { - throw new WorkflowStoredDataError(`${claimed.id}.argsJson`, cause); - } - - 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, - availableProviders, - agentProfiles, - 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 providerRun = await runLocalAgentProviderResult(input.provider, { - prompt: input.prompt, - workspace: input.workspace, - providerSessionId: input.providerSessionId, - model: input.model, - effort: input.effort, - writeMode: "allowed", - schema: input.schema, - }); - if (providerRun.isErr()) throw providerRun.error; - const providerResult = providerRun.value; - return { - finalResponse: providerResult.finalResponse, - providerSessionId: providerResult.providerSessionId ?? undefined, - structured: providerResult.structured, - }; - }, - resolveNestedSource: async (ref) => { - if (typeof ref === "string") { - const named = await resolveNamedWorkflowScriptResult({ - name: ref, - workspaceRoot: claimed.workspaceRoot, - stateDir: config.stateDir, - }); - if (named.isErr()) throw named.error; - return named.value.source; - } - const nested = await readProjectWorkflowScriptFileResult({ - scriptPath: ref.scriptPath, - workspaceRoot: claimed.workspaceRoot, - }); - if (nested.isErr()) throw nested.error; - return nested.value.source; - }, - }); - - 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, 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 }); - } catch { - // terminal race + if (!call) { + throw new InvalidWorkflowInputError({ + code: "invalid_argument", + message: `Unknown workflow agent call: ${runId}#${callIndex}`, + }); } + console.log(JSON.stringify(formatCallDetail(call), null, 2)); } 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 (;;) { @@ -601,12 +407,6 @@ function safeParseJson(text: string): unknown { } } -function resolveAvailableProviders(): LocalAgentProvider[] { - 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[]; @@ -661,7 +461,3 @@ function collectArgTokens(args: string[]): string[] { function sleep(ms: number): Promise { return new Promise((resolveSleep) => setTimeout(resolveSleep, ms)); } - -function isJsonObject(value: JsonValue): value is JsonObject { - return typeof value === "object" && value !== null && !Array.isArray(value); -} diff --git a/src/workflow-engine.test.ts b/src/workflow-engine.test.ts index 7378b046..83091f2b 100644 --- a/src/workflow-engine.test.ts +++ b/src/workflow-engine.test.ts @@ -538,15 +538,19 @@ import type { LocalAgentProfile } from "./local-agent-profiles.js"; }), }); - await assert.rejects( - () => - api.agent("large structured", { + const oversizedStructured = () => + (api.agent as (prompt: string, opts?: object) => Promise)( + "large structured", + { schema: { type: "object", properties: { big: { type: "string" } }, required: ["big"], }, - }), + }, + ); + await assert.rejects( + oversizedStructured, (error: unknown) => error instanceof WorkflowEngineError && error.kind === "result_too_large", ); diff --git a/src/workflow-launch.test.ts b/src/workflow-launch.test.ts new file mode 100644 index 00000000..3f40f919 --- /dev/null +++ b/src/workflow-launch.test.ts @@ -0,0 +1,66 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm, writeFile, mkdir } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { WorkflowStore } from "./workflow-store.js"; +import { launchWorkflowRun } from "./workflow-launch.js"; + +{ + const dir = await mkdtemp(join(tmpdir(), "wf-launch-")); + const store = new WorkflowStore(dir); + const launched = await launchWorkflowRun({ + store, + config: { stateDir: dir }, + workspaceRoot: dir, + source: { + kind: "inline", + script: `export const meta = { name: 'launch-demo', description: 'd' }\nreturn 1\n`, + }, + args: { n: 1 }, + cliEntry: "/tmp/devspace-cli-not-used", + spawn: false, + }); + assert.equal(launched.isOk(), true); + if (!launched.isOk()) throw launched.error; + assert.equal(launched.value.run.name, "launch-demo"); + assert.equal(launched.value.run.status, "starting"); + assert.match(launched.value.run.scriptPath.replaceAll("\\", "/"), /workflow-scripts\//); + assert.equal(launched.value.run.argsJson, JSON.stringify({ n: 1 })); + + await mkdir(join(dir, ".devspace", "workflows"), { recursive: true }); + await writeFile( + join(dir, ".devspace", "workflows", "named-wf.js"), + `export const meta = { name: 'named-wf', description: 'd' }\nreturn 2\n`, + ); + const named = await launchWorkflowRun({ + store, + config: { stateDir: dir }, + workspaceRoot: dir, + source: { kind: "named", name: "named-wf" }, + cliEntry: "/tmp/devspace-cli-not-used", + spawn: false, + }); + assert.equal(named.isOk(), true); + if (!named.isOk()) throw named.error; + assert.equal(named.value.source, "named"); + assert.equal(named.value.run.name, "named-wf"); + + const resumed = await launchWorkflowRun({ + store, + config: { stateDir: dir }, + workspaceRoot: dir, + source: { kind: "resume", runId: launched.value.run.id }, + cliEntry: "/tmp/devspace-cli-not-used", + spawn: false, + }); + assert.equal(resumed.isOk(), true); + if (!resumed.isOk()) throw resumed.error; + assert.equal(resumed.value.source, "resume"); + assert.equal(resumed.value.run.resumedFromRunId, launched.value.run.id); + assert.equal(resumed.value.run.argsJson, JSON.stringify({ n: 1 })); + + store.close(); + await rm(dir, { recursive: true, force: true }); +} + +console.log("workflow-launch.test.ts: ok"); diff --git a/src/workflow-launch.ts b/src/workflow-launch.ts new file mode 100644 index 00000000..23d44144 --- /dev/null +++ b/src/workflow-launch.ts @@ -0,0 +1,276 @@ +import type { ServerConfig } from "./config.js"; +import { parseJsonText, type JsonObject, type JsonValue } from "./json-types.js"; +import { + persistWorkflowScriptResult, + readWorkflowScriptFileResult, + resolveNamedWorkflowScriptResult, + resolveWorkflowScriptFromPathOrNameResult, +} from "./workflow-files.js"; +import { parseWorkflowScript } from "./workflow-script.js"; +import type { WorkflowStore } from "./workflow-store.js"; +import type { WorkflowRunRecord, WorkflowRunSource } from "./workflow-types.js"; +import { + InvalidWorkflowInputError, + WorkflowNotFoundError, + WorkflowStoredDataError, + type WorkflowOperationError, +} from "./workflow-errors.js"; +import { resolveWorkspaceHead } from "./workflow-worktrees.js"; +import { spawnWorkflowWorker } from "./workflow-worker.js"; +import { Result, type Result as BetterResult } from "better-result"; +import type { WorkflowScriptError } from "./workflow-script.js"; +import type { WorkflowFileWriteError } from "./workflow-errors.js"; +import type { WorkflowRunTransitionError } from "./workflow-store.js"; + +export type LaunchWorkflowSource = + | { kind: "inline"; script: string; filename?: string } + | { kind: "file"; path: string } + | { kind: "named"; name: string } + | { + kind: "resume"; + runId: string; + /** Optional replacement source while resuming. */ + override?: + | { kind: "inline"; script: string; filename?: string } + | { kind: "file"; path: string } + | { kind: "named"; name: string }; + }; + +export interface LaunchWorkflowRunInput { + store: WorkflowStore; + config: Pick; + workspaceRoot: string; + workspaceId?: string; + source: LaunchWorkflowSource; + args?: JsonValue; + /** Absolute path to cli entry used to spawn `workflow __worker`. */ + cliEntry: string; + /** When false, create the run row but do not spawn (tests). Default true. */ + spawn?: boolean; +} + +export type LaunchWorkflowError = + | WorkflowOperationError + | WorkflowScriptError + | WorkflowFileWriteError + | WorkflowRunTransitionError; + +export interface LaunchWorkflowRunResult { + run: WorkflowRunRecord; + parsedName: string; + scriptHash: string; + source: WorkflowRunSource; +} + +/** + * Shared CLI/MCP start path: resolve script → parse → create run → persist → spawn. + */ +export async function launchWorkflowRun( + input: LaunchWorkflowRunInput, +): Promise> { + try { + const resolved = await resolveLaunchSource(input); + if (resolved.isErr()) return resolved; + + const { + sourceText, + scriptHash, + nameHint, + runSource, + priorRunId, + filename, + args, + } = resolved.value; + + const parsed = parseWorkflowScript(sourceText, { filename }); + const baseSha = await resolveWorkspaceHead(input.workspaceRoot); + const preferredName = parsed.meta.name || nameHint; + + const run = input.store.createRun({ + name: preferredName, + source: runSource, + scriptPath: "pending", + scriptHash, + workspaceRoot: input.workspaceRoot, + workspaceId: input.workspaceId, + argsJson: JSON.stringify(args === undefined ? null : args), + resumedFromRunId: priorRunId, + baseSha, + }); + + const persisted = await persistWorkflowScriptResult({ + stateDir: input.config.stateDir, + runId: run.id, + source: sourceText, + preferredName, + }); + if (persisted.isErr()) return persisted; + + const updated = input.store.setScriptPathResult(run.id, persisted.value); + if (updated.isErr()) return updated; + + if (input.spawn !== false) { + spawnWorkflowWorker(run.id, input.cliEntry); + } + + return Result.ok({ + run: updated.value, + parsedName: preferredName, + scriptHash, + source: runSource, + }); + } catch (error) { + if (isLaunchError(error)) return Result.err(error); + throw error; + } +} + +interface ResolvedLaunch { + sourceText: string; + scriptHash: string; + nameHint: string; + runSource: WorkflowRunSource; + priorRunId?: string; + filename: string; + args: JsonValue | undefined; +} + +async function resolveLaunchSource( + input: LaunchWorkflowRunInput, +): Promise> { + const { source, store, config, workspaceRoot } = input; + let args = input.args; + + if (source.kind === "resume") { + const priorResult = store.getRunResult(source.runId); + if (priorResult.isErr()) return priorResult; + const prior = priorResult.value; + if (!prior) return Result.err(new WorkflowNotFoundError(source.runId)); + + let sourceText: string; + let scriptHash: string; + let nameHint: string; + let filename: string; + + if (source.override?.kind === "inline") { + sourceText = source.override.script; + const overrideParsed = parseWorkflowScript(sourceText, { + filename: source.override.filename ?? "workflow:inline", + }); + scriptHash = overrideParsed.scriptHash; + nameHint = overrideParsed.meta.name; + filename = source.override.filename ?? "workflow:inline"; + } else if (source.override?.kind === "named") { + const named = await resolveNamedWorkflowScriptResult({ + name: source.override.name, + workspaceRoot, + stateDir: config.stateDir, + }); + if (named.isErr()) return named; + sourceText = named.value.source; + scriptHash = named.value.scriptHash; + nameHint = named.value.nameHint; + filename = named.value.scriptPath; + } else if (source.override?.kind === "file") { + const file = await readWorkflowScriptFileResult(source.override.path); + if (file.isErr()) return file; + sourceText = file.value.source; + scriptHash = file.value.scriptHash; + nameHint = file.value.nameHint; + filename = file.value.scriptPath; + } else { + const priorScript = await readWorkflowScriptFileResult(prior.scriptPath); + if (priorScript.isErr()) return priorScript; + sourceText = priorScript.value.source; + scriptHash = priorScript.value.scriptHash; + nameHint = prior.name; + filename = prior.scriptPath; + } + + if (args === undefined && prior.argsJson && prior.argsJson !== "null") { + try { + args = parseJsonText(prior.argsJson); + } catch (cause) { + return Result.err(new WorkflowStoredDataError(`${prior.id}.argsJson`, cause)); + } + } + + return Result.ok({ + sourceText, + scriptHash, + nameHint, + runSource: "resume", + priorRunId: prior.id, + filename, + args, + }); + } + + if (source.kind === "inline") { + const parsed = parseWorkflowScript(source.script, { + filename: source.filename ?? "workflow:inline", + }); + return Result.ok({ + sourceText: source.script, + scriptHash: parsed.scriptHash, + nameHint: parsed.meta.name, + runSource: "inline", + filename: source.filename ?? "workflow:inline", + args, + }); + } + + if (source.kind === "named") { + const named = await resolveNamedWorkflowScriptResult({ + name: source.name, + workspaceRoot, + stateDir: config.stateDir, + }); + if (named.isErr()) return named; + return Result.ok({ + sourceText: named.value.source, + scriptHash: named.value.scriptHash, + nameHint: named.value.nameHint, + runSource: "named", + filename: named.value.scriptPath, + args, + }); + } + + if (source.kind === "file") { + const file = await resolveWorkflowScriptFromPathOrNameResult({ + file: source.path, + workspaceRoot, + stateDir: config.stateDir, + }); + if (file.isErr()) return file; + return Result.ok({ + sourceText: file.value.source, + scriptHash: file.value.scriptHash, + nameHint: file.value.nameHint, + runSource: file.value.origin === "named" ? "named" : "inline", + filename: file.value.scriptPath, + args, + }); + } + + return Result.err( + new InvalidWorkflowInputError({ + code: "missing_source", + message: "Provide a workflow script source", + }), + ); +} + +function isLaunchError(error: unknown): error is LaunchWorkflowError { + return ( + typeof error === "object" && + error !== null && + "name" in error && + (error as { name?: string }).name === "WorkflowScriptError" + ); +} + +export function isJsonObject(value: JsonValue): value is JsonObject { + return typeof value === "object" && value !== null && !Array.isArray(value); +} diff --git a/src/workflow-providers.ts b/src/workflow-providers.ts new file mode 100644 index 00000000..31836fd3 --- /dev/null +++ b/src/workflow-providers.ts @@ -0,0 +1,12 @@ +import { getLocalAgentProviderAvailabilitySnapshot } from "./local-agent-availability.js"; +import { + LOCAL_AGENT_PROVIDERS, + type LocalAgentProvider, +} from "./local-agent-profiles.js"; + +/** Live providers in stable product order for workflow agent() resolution. */ +export function resolveWorkflowLiveProviders(): LocalAgentProvider[] { + 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)); +} diff --git a/src/workflow-tools.ts b/src/workflow-tools.ts index 6efbb6f6..e648467a 100644 --- a/src/workflow-tools.ts +++ b/src/workflow-tools.ts @@ -5,12 +5,6 @@ import * as z from "zod/v4"; import type { ServerConfig } from "./config.js"; import { jsonValueSchema, parseJsonText, type JsonValue } from "./json-types.js"; import type { WorkspaceRegistry } from "./workspaces.js"; -import { - persistWorkflowScriptResult, - resolveNamedWorkflowScriptResult, - readWorkflowScriptFileResult, -} from "./workflow-files.js"; -import { parseWorkflowScript } from "./workflow-script.js"; import { createWorkflowStore } from "./workflow-store.js"; import { WORKFLOW_MCP_YIELD_MS, @@ -18,26 +12,23 @@ import { type WorkflowEventRecord, type WorkflowRunRecord, } from "./workflow-types.js"; -import { resolveWorkspaceHead } from "./workflow-worktrees.js"; -import { spawnWorkflowWorkerFromCli } from "./workflow-cli.js"; import { cancelWorkflowRun } from "./workflow-lifecycle.js"; -import { getLocalAgentProviderAvailabilitySnapshot } from "./local-agent-availability.js"; -import { - LOCAL_AGENT_PROVIDERS, - type LocalAgentProvider, -} from "./local-agent-profiles.js"; import { InvalidWorkflowInputError, isWorkflowOperationError, serializeWorkflowError, WorkflowNotFoundError, - WorkflowStoredDataError, } from "./workflow-errors.js"; import { loadWorkflowUiCallDetail, loadWorkflowUiProject, loadWorkflowUiRun, } from "./workflow-ui.js"; +import { + launchWorkflowRun, + type LaunchWorkflowSource, +} from "./workflow-launch.js"; +import { resolveWorkflowLiveProviders } from "./workflow-providers.js"; const WORKSPACE_APP_URI = "ui://devspace/workspace-app.html"; const WORKFLOW_UI_WAIT_MAX_MS = 30_000; @@ -106,109 +97,30 @@ export function registerWorkflowTools( }); } - let source: string; - let scriptHash: string; - let nameHint: string; - let priorRunId: string | undefined; - let priorScriptPath: string | undefined; - let runSource: "inline" | "named" | "resume" = "inline"; - - if (resumeFromRunId) { - const priorResult = store.getRunResult(resumeFromRunId); - if (priorResult.isErr()) throw priorResult.error; - const prior = priorResult.value; - if (!prior) throw new WorkflowNotFoundError(resumeFromRunId); - priorRunId = prior.id; - const overridePath = scriptPath; - if (script !== undefined) { - source = script; - const overrideParsed = parseWorkflowScript(source); - scriptHash = overrideParsed.scriptHash; - nameHint = overrideParsed.meta.name; - } else if (name) { - const resolvedResult = await resolveNamedWorkflowScriptResult({ - name, - workspaceRoot: workspace.root, - stateDir: config.stateDir, - }); - if (resolvedResult.isErr()) throw resolvedResult.error; - source = resolvedResult.value.source; - scriptHash = resolvedResult.value.scriptHash; - nameHint = resolvedResult.value.nameHint; - } else { - priorScriptPath = overridePath ?? prior.scriptPath; - const resolvedResult = await readWorkflowScriptFileResult(priorScriptPath); - if (resolvedResult.isErr()) throw resolvedResult.error; - source = resolvedResult.value.source; - scriptHash = resolvedResult.value.scriptHash; - nameHint = overridePath ? resolvedResult.value.nameHint : prior.name; - } - runSource = "resume"; - if (args === undefined && prior.argsJson && prior.argsJson !== "null") { - try { - args = parseJsonText(prior.argsJson); - } catch (cause) { - throw new WorkflowStoredDataError(`${prior.id}.argsJson`, cause); - } - } - } else if (name) { - 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; - runSource = "named"; - } else if (scriptPath) { - const resolvedResult = await readWorkflowScriptFileResult(scriptPath); - if (resolvedResult.isErr()) throw resolvedResult.error; - source = resolvedResult.value.source; - scriptHash = resolvedResult.value.scriptHash; - nameHint = resolvedResult.value.nameHint; - } else { - source = script!; - const parsed = parseWorkflowScript(source); - scriptHash = parsed.scriptHash; - nameHint = parsed.meta.name; - runSource = "inline"; - } - - const parsed = parseWorkflowScript(source); - const baseSha = await resolveWorkspaceHead(workspace.root); - const run = store.createRun({ - name: parsed.meta.name || nameHint, - source: runSource, - scriptPath: "pending", - scriptHash, + const source = buildMcpLaunchSource({ + script, + name, + scriptPath, + resumeFromRunId, + }); + const launched = await launchWorkflowRun({ + store, + config, workspaceRoot: workspace.root, workspaceId, - argsJson: JSON.stringify(args ?? null), - resumedFromRunId: priorRunId, - baseSha, - }); - - const persistedResult = await persistWorkflowScriptResult({ - stateDir: config.stateDir, - runId: run.id, source, - preferredName: parsed.meta.name || nameHint, + args, + cliEntry: fileURLToPath( + import.meta.url.replace(/workflow-tools\.(ts|js)$/, "cli.$1"), + ), }); - if (persistedResult.isErr()) throw persistedResult.error; - const 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"), - ); - spawnWorkflowWorkerFromCli(run.id, cliEntry); + if (launched.isErr()) { + if (isWorkflowOperationError(launched.error)) return workflowToolError(launched.error); + throw launched.error; + } const yieldMs = yieldTimeMs ?? 2_000; - const page = await yieldEvents(store, run.id, 0, yieldMs); + const page = await yieldEvents(store, launched.value.run.id, 0, yieldMs); return toolResult(page, "run_workflow"); } catch (error) { if (isWorkflowOperationError(error)) return workflowToolError(error); @@ -558,11 +470,33 @@ function sleep(ms: number): Promise { return new Promise((r) => setTimeout(r, ms)); } -/** Resolve live providers in stable product order for workflows. */ -export function resolveWorkflowEnabledProviders(): LocalAgentProvider[] { - const snapshot = getLocalAgentProviderAvailabilitySnapshot(); - const live = new Set( - snapshot.filter((row) => row.available).map((row) => row.name), - ); - return LOCAL_AGENT_PROVIDERS.filter((id): id is LocalAgentProvider => live.has(id)); +function buildMcpLaunchSource(input: { + script?: string; + name?: string; + scriptPath?: string; + resumeFromRunId?: string; +}): LaunchWorkflowSource { + if (input.resumeFromRunId) { + const override = + input.script !== undefined + ? ({ kind: "inline", script: input.script } as const) + : input.name + ? ({ kind: "named", name: input.name } as const) + : input.scriptPath + ? ({ kind: "file", path: input.scriptPath } as const) + : undefined; + return { kind: "resume", runId: input.resumeFromRunId, override }; + } + if (input.name) return { kind: "named", name: input.name }; + if (input.scriptPath) return { kind: "file", path: input.scriptPath }; + if (input.script !== undefined) return { kind: "inline", script: input.script }; + throw new InvalidWorkflowInputError({ + code: "missing_source", + message: "Provide script, name, scriptPath, or resumeFromRunId", + }); +} + +/** @deprecated Prefer resolveWorkflowLiveProviders from workflow-providers.js */ +export function resolveWorkflowEnabledProviders() { + return resolveWorkflowLiveProviders(); } diff --git a/src/workflow-worker.ts b/src/workflow-worker.ts new file mode 100644 index 00000000..619b5ba0 --- /dev/null +++ b/src/workflow-worker.ts @@ -0,0 +1,191 @@ +import { spawn } from "node:child_process"; +import { readFile } from "node:fs/promises"; +import { availableParallelism } from "node:os"; +import type { ServerConfig } from "./config.js"; +import { parseJsonText, type JsonValue } from "./json-types.js"; +import { runLocalAgentProviderResult } from "./local-agent-adapters.js"; +import { + isLocalAgentProvider, + loadLocalAgentProfiles, +} from "./local-agent-profiles.js"; +import { executeWorkflow, mapEngineErrorKind } from "./workflow-engine.js"; +import { + readProjectWorkflowScriptFileResult, + resolveNamedWorkflowScriptResult, +} from "./workflow-files.js"; +import { createWorkflowReplay } from "./workflow-replay.js"; +import { parseWorkflowScript } from "./workflow-script.js"; +import { createWorkflowStore } from "./workflow-store.js"; +import { + WORKFLOW_HEARTBEAT_MS, + WORKFLOW_LIMITS, + resolveWorkflowConcurrency, +} from "./workflow-types.js"; +import { WorkflowStoredDataError } from "./workflow-errors.js"; +import { createWorkflowWorktreeFactory } from "./workflow-worktrees.js"; +import { resolveWorkflowLiveProviders } from "./workflow-providers.js"; + +/** 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 claim = store.claimRunResult(runId, process.pid); + if (claim.isErr()) { + store.close(); + throw claim.error; + } + const claimed = claim.value; + + 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 availableProviders = resolveWorkflowLiveProviders(); + const agentProfiles = await loadLocalAgentProfiles(config, claimed.workspaceRoot); + const concurrency = resolveWorkflowConcurrency( + parsed.meta.concurrency, + availableParallelism(), + ); + + let argsValue: JsonValue | undefined; + try { + argsValue = parseJsonText(claimed.argsJson); + if (argsValue === null) argsValue = undefined; + } catch (cause) { + throw new WorkflowStoredDataError(`${claimed.id}.argsJson`, cause); + } + + 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, + availableProviders, + agentProfiles, + 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 providerRun = await runLocalAgentProviderResult(input.provider, { + prompt: input.prompt, + workspace: input.workspace, + providerSessionId: input.providerSessionId, + model: input.model, + effort: input.effort, + writeMode: "allowed", + schema: input.schema, + }); + if (providerRun.isErr()) throw providerRun.error; + const providerResult = providerRun.value; + return { + finalResponse: providerResult.finalResponse, + providerSessionId: providerResult.providerSessionId ?? undefined, + structured: providerResult.structured, + }; + }, + resolveNestedSource: async (ref) => { + if (typeof ref === "string") { + const named = await resolveNamedWorkflowScriptResult({ + name: ref, + workspaceRoot: claimed.workspaceRoot, + stateDir: config.stateDir, + }); + if (named.isErr()) throw named.error; + return named.value.source; + } + const nested = await readProjectWorkflowScriptFileResult({ + scriptPath: ref.scriptPath, + workspaceRoot: claimed.workspaceRoot, + }); + if (nested.isErr()) throw nested.error; + return nested.value.source; + }, + }); + + 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, 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 }); + } catch { + // terminal race + } + } finally { + clearInterval(heartbeat); + store.close(); + } +} + +export function spawnWorkflowWorker(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(); +} + +/** @deprecated Use spawnWorkflowWorker */ +export const spawnWorkflowWorkerFromCli = spawnWorkflowWorker; From 8337f2c23402e6d6e1a8bf93f4d915e5ab1bd785 Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:17:38 +0530 Subject: [PATCH 6/7] refactor(workflow): drop consume-once compatible_key and unused schema helpers Prefix-only replay never emits compatible_key. Narrow contracts/types to same_index and remove unused schemaAwareRunProvider wrappers. --- src/workflow-api.ts | 4 ++-- src/workflow-contracts.ts | 2 +- src/workflow-schema.ts | 45 +-------------------------------------- src/workflow-store.ts | 6 +++--- src/workflow-types.ts | 2 +- src/workflow-ui.ts | 2 +- src/workflow-view.ts | 2 +- 7 files changed, 10 insertions(+), 53 deletions(-) diff --git a/src/workflow-api.ts b/src/workflow-api.ts index 071b0d59..525fba03 100644 --- a/src/workflow-api.ts +++ b/src/workflow-api.ts @@ -76,7 +76,7 @@ export interface WorkflowReplayHit { structuredJson?: string; returnValueJson: string; providerSessionId?: string; - replayMatch: "same_index" | "compatible_key"; + replayMatch: "same_index"; replayedFromRunId: string; replayedFromCallIndex: number; } @@ -125,7 +125,7 @@ export interface WorkflowJournal { phase?: string; isolation?: AgentIsolationMode; worktreePath?: string; - replayMatch?: "same_index" | "compatible_key"; + replayMatch?: "same_index"; replayedFromRunId?: string; replayedFromCallIndex?: number; replayReason?: string; diff --git a/src/workflow-contracts.ts b/src/workflow-contracts.ts index d4641962..b06841d0 100644 --- a/src/workflow-contracts.ts +++ b/src/workflow-contracts.ts @@ -206,7 +206,7 @@ export const workflowEventPayloadSchemas = { callIndex: z.number().int().nonnegative(), cacheKey: z.string(), provider: localAgentProviderSchema, - replayMatch: z.enum(["same_index", "compatible_key"]), + replayMatch: z.enum(["same_index"]), replayedFromRunId: z.string(), replayedFromCallIndex: z.number().int().nonnegative(), }) diff --git a/src/workflow-schema.ts b/src/workflow-schema.ts index 0249c2dd..5bdd856e 100644 --- a/src/workflow-schema.ts +++ b/src/workflow-schema.ts @@ -2,7 +2,7 @@ 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 type { WorkflowProviderRunResult } from "./workflow-api.js"; import { classifyAgentProviderError, isProviderSchemaUnsupportedError, @@ -233,49 +233,6 @@ function toSchemaIssues( })); } -/** Helper for wiring into agent(): wrap a one-shot provider as retrying schema runner. */ -export function schemaAwareRunProvider( - runProvider: WorkflowRunProvider, - schema: JsonSchema, - base: Parameters[0], - onRetry?: EnforceSchemaInput["onRetry"], -): Promise { - return enforceAgentSchema({ - schema, - prompt: base.prompt, - provider: base.provider, - onRetry, - run: (prompt, options) => - runProvider({ - ...base, - prompt, - providerSessionId: options.providerSessionId, - ...(options.mode === "native" ? { schema } : {}), - }), - }); -} - -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 5f59df77..5c30fa51 100644 --- a/src/workflow-store.ts +++ b/src/workflow-store.ts @@ -61,7 +61,7 @@ export interface BeginAgentCallInput { phase?: string; isolation?: AgentIsolationMode; worktreePath?: string; - replayMatch?: "same_index" | "compatible_key"; + replayMatch?: "same_index"; replayedFromRunId?: string; replayedFromCallIndex?: number; replayReason?: string; @@ -918,8 +918,8 @@ function rowToAgentCall(row: WorkflowAgentCallRow): WorkflowAgentCallRecord { error: row.error ?? undefined, errorKind: (row.error_kind as WorkflowErrorKind | null) ?? undefined, replayMatch: - row.replay_match === "same_index" || row.replay_match === "compatible_key" - ? row.replay_match + row.replay_match === "same_index" + ? "same_index" : undefined, replayedFromRunId: row.replayed_from_run_id ?? undefined, replayedFromCallIndex: row.replayed_from_call_index ?? undefined, diff --git a/src/workflow-types.ts b/src/workflow-types.ts index 4f46e1d7..2e3996b7 100644 --- a/src/workflow-types.ts +++ b/src/workflow-types.ts @@ -134,7 +134,7 @@ export interface WorkflowAgentCallRecord { returnValueJson?: string; error?: string; errorKind?: WorkflowErrorKind; - replayMatch?: "same_index" | "compatible_key"; + replayMatch?: "same_index"; replayedFromRunId?: string; replayedFromCallIndex?: number; replayReason?: string; diff --git a/src/workflow-ui.ts b/src/workflow-ui.ts index b97b526a..a2c4b81e 100644 --- a/src/workflow-ui.ts +++ b/src/workflow-ui.ts @@ -38,7 +38,7 @@ export interface WorkflowCallDetailView { worktreePath?: string; dirty?: boolean; fromCache: boolean; - replayMatch?: "same_index" | "compatible_key"; + replayMatch?: "same_index"; replayedFromRunId?: string; replayedFromCallIndex?: number; replayReason?: string; diff --git a/src/workflow-view.ts b/src/workflow-view.ts index 218bc246..919d5347 100644 --- a/src/workflow-view.ts +++ b/src/workflow-view.ts @@ -35,7 +35,7 @@ export interface WorkflowCallView { worktreePath?: string; dirty?: boolean; fromCache: boolean; - replayMatch?: "same_index" | "compatible_key"; + replayMatch?: "same_index"; replayedFromRunId?: string; replayedFromCallIndex?: number; replayReason?: string; From e55f20eeba87d5dab564d8bafb050ea900cfe05f Mon Sep 17 00:00:00 2001 From: Waishnav Date: Mon, 27 Jul 2026 23:17:54 +0530 Subject: [PATCH 7/7] docs(workflow): align resume docs with prefix-only replay Update primitives-spec and the dynamic-workflows skill so resume semantics and oversized return failure match the implementation. --- docs/dynamic-workflow/devspace/primitives-spec.md | 9 +++++---- skills/dynamic-workflows/SKILL.md | 4 ++++ 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/docs/dynamic-workflow/devspace/primitives-spec.md b/docs/dynamic-workflow/devspace/primitives-spec.md index 0f23ea6a..219408f3 100644 --- a/docs/dynamic-workflow/devspace/primitives-spec.md +++ b/docs/dynamic-workflow/devspace/primitives-spec.md @@ -55,7 +55,7 @@ ChatGPT ── MCP tools ──► engine ── agent() ──► adapter | `budget` | Shared host token hard ceiling | **Stub** `{ total: null, spent:0, remaining: Infinity }` | | `workflow()` | Nested name/scriptPath; depth 1; shared caps | Same spirit | | Determinism bans | Date.now / Math.random / bare new Date | Same | -| Resume | Prefix cache by prompt+opts | Index+key + consume-once cacheKey fallback | +| Resume | Prefix cache by prompt+opts | Deterministic call-index prefix only (first miss closes replay) | | File diffs per stage | **Not a primitive** | Same — no auto-diff | --- @@ -614,11 +614,12 @@ Adapters: no individual abort API — accepted; group-kill is backstop. |---|---| | New run | `--resume` / `resumeFromRunId` creates new run with `resumedFromRunId`. | | Cache key | `sha256(canonicalJson({ prompt, provider, model, effort, schema, isolation }))` | -| Match | (1) same callIndex + key (2) on first miss, consume-once by key (fan-out order). | +| Match | Same callIndex + cache key while the prefix remains open. | +| Close | First failed, interrupted, changed, missing, corrupt, worktree, or unpersisted result executes live and closes replay for later calls. | | Record | Cache hits written as new rows `from_cache=1` so chains chain. | | Determinism | Bans make prompt construction stable if args fixed. | -Document CC divergence (consume-once) in skill. +Document prefix-only resume (no consume-once key fallback) in skill. --- @@ -748,7 +749,7 @@ return pipeline( | log / args / budget | `workflow-api.ts` | freeze, stub budget | | workflow nest | `workflow-api.ts` + engine | depth 1, shared journal | | store | `workflow-store.ts` | seq, reap, cancel | -| replay | `workflow-replay.ts` | index+key, consume-once | +| replay | `workflow-replay.ts` | deterministic call-index prefix | | CLI | `cli.ts` | run/status/cancel/ls/__worker | | MCP | `workflow-tools.ts` | yield, survive disconnect | | skill | `skills/dynamic-workflows` | education | diff --git a/skills/dynamic-workflows/SKILL.md b/skills/dynamic-workflows/SKILL.md index de4ab181..49487cbc 100644 --- a/skills/dynamic-workflows/SKILL.md +++ b/skills/dynamic-workflows/SKILL.md @@ -117,6 +117,10 @@ calls assume their existing filesystem effects are still present. Worktree calls are never reused unless their exact worktree can be restored, so they currently end the reusable prefix and run live. +Return values must fit the replay budget (~1 MiB JSON). Oversized returns fail +the `agent()` call with `result_too_large` — prefer summaries or paths to large +artifacts on disk. + ### Cancel `workflow cancel` sets a cooperative flag; worker aborts then hard-kills if needed.