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/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/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. diff --git a/src/workflow-api.ts b/src/workflow-api.ts index cd0d60da..525fba03 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) @@ -73,7 +76,7 @@ export interface WorkflowReplayHit { structuredJson?: string; returnValueJson: string; providerSessionId?: string; - replayMatch: "same_index" | "compatible_key"; + replayMatch: "same_index"; replayedFromRunId: string; replayedFromCallIndex: number; } @@ -122,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; @@ -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 // --------------------------------------------------------------------------- @@ -466,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"); @@ -489,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, @@ -623,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, @@ -633,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) }, }); }; @@ -792,23 +799,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-cli.ts b/src/workflow-cli.ts index bfd00367..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, - readProjectWorkflowScriptFile, - readWorkflowScriptFileResult, - resolveNamedWorkflowScript, - 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,180 +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 resolveNamedWorkflowScript({ - name: ref, - workspaceRoot: claimed.workspaceRoot, - stateDir: config.stateDir, - }); - return named.source; - } - return ( - await readProjectWorkflowScriptFile({ - scriptPath: ref.scriptPath, - workspaceRoot: claimed.workspaceRoot, - }) - ).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 (;;) { @@ -600,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[]; @@ -660,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-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-engine.test.ts b/src/workflow-engine.test.ts index fc4bf57d..83091f2b 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,23 @@ import type { LocalAgentProfile } from "./local-agent-profiles.js"; }), }); - assert.deepEqual( - await api.agent("large structured", { - schema: { - type: "object", - properties: { big: { type: "string" } }, - required: ["big"], + const oversizedStructured = () => + (api.agent as (prompt: string, opts?: object) => Promise)( + "large structured", + { + schema: { + type: "object", + properties: { big: { type: "string" } }, + required: ["big"], + }, }, - }), - { big }, + ); + await assert.rejects( + oversizedStructured, + (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); 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", )<{ 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; -} 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-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"); 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-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-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; 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;