From c921419284ec90053ae0b8308b1dad2ef58bd39e Mon Sep 17 00:00:00 2001 From: Waishnav Date: Tue, 21 Jul 2026 17:42:38 +0530 Subject: [PATCH 1/3] feat(workflow): add dynamic-workflows skill and seed both defaults Always include bundled skills root when subagents enabled so seeding subagent-delegation no longer hides later skills. Seed dynamic-workflows alongside subagent-delegation on init. Co-Authored-By: Claude --- skills/dynamic-workflows/SKILL.md | 132 ++++++++++++++++++++++++++++++ src/config.test.ts | 7 +- src/skills.ts | 25 +++--- src/user-config.ts | 26 ++++-- 4 files changed, 171 insertions(+), 19 deletions(-) create mode 100644 skills/dynamic-workflows/SKILL.md diff --git a/skills/dynamic-workflows/SKILL.md b/skills/dynamic-workflows/SKILL.md new file mode 100644 index 00000000..e7ab685d --- /dev/null +++ b/skills/dynamic-workflows/SKILL.md @@ -0,0 +1,132 @@ +--- +name: dynamic-workflows +description: Orchestrate multi-agent coding workflows via DevSpace Dynamic Workflows (CLI or MCP). +--- + +# Dynamic Workflows + +Use this skill when the user wants multi-step, multi-agent orchestration — fan-out +review, migrate-and-verify, research panels — **not** a single subagent turn. + +## Entry points + +| Host | Surface | +|---|---| +| Coding agent (Claude Code, Codex, pi, …) | CLI + this skill | +| ChatGPT / MCP client | MCP tools `run_workflow` / `workflow_status` / `workflow_cancel` | + +```bash +devspace workflow run --file path/to/script.js [--arg k=v]... [--follow] +devspace workflow run --name review-auth [--follow] +devspace workflow run --resume +devspace workflow status [--follow] +devspace workflow cancel +devspace workflow ls +``` + +Named scripts: `.devspace/workflows/.js` or `workflows/.js`. + +## Script shape + +```js +export const meta = { + name: 'review-auth', + description: 'Fan-out review of auth changes', + phases: [{ title: 'Review' }, { title: 'Synthesize' }], + // optional DevSpace: + // defaultProvider: 'codex', + // concurrency: 4, +} + +phase('Review') +const findings = await parallel([ + () => agent('Review for correctness…', { label: 'correctness' }), + () => agent('Review for security…', { label: 'security' }), +]) +phase('Synthesize') +const summary = await agent(`Synthesize: ${JSON.stringify(findings)}`) +return { summary, findings } +``` + +### Primitives + +| API | Notes | +|---|---| +| `agent(prompt, opts?)` | Throws on failure. `opts`: `label`, `phase`, `schema`, `model`, `effort`, `provider`, `isolation: 'worktree'` | +| `parallel(thunks)` | Barrier; throw → `null` slot | +| `pipeline(items, ...stages)` | Per-item chains; no cross-item barrier | +| `phase(title)` / `log(msg)` | Progress; journaled | +| `args` | Run input (object preferred) | +| `budget` | Stub: `total: null`, `remaining(): Infinity` — do not loop on budget alone | +| `workflow(name\|{scriptPath}, args?)` | Nested, depth 1, shared call index | + +**No `writeMode`.** Teach read-only vs write in the prompt. Use `isolation: 'worktree'` when parallel mutators would conflict (git required). + +### Determinism bans + +`Date.now()`, `Math.random()`, and `new Date()` without args throw. Pass timestamps via `args` if needed. + +### Schema + +```js +const out = await agent('Return JSON findings', { + schema: { + type: 'object', + properties: { bugs: { type: 'array', items: { type: 'string' } } }, + required: ['bugs'], + }, +}) +// out is validated object; engine retries ≤2 on invalid JSON +``` + +### Providers + +Default: first **enabled ∩ available** provider (`agentProviders.enabled` in config, else all live providers in product order). Override with `opts.provider` or `meta.defaultProvider`. + +### Resume + +`devspace workflow run --resume ` creates a **new** run that replays completed agent calls by cache key (callIndex+key, then consume-once by key). + +### Cancel + +`workflow cancel` sets a cooperative flag; worker aborts then hard-kills if needed. + +## When to use CLI vs MCP + +- **CLI**: host agent can shell; prefer for long runs + `--follow`. +- **MCP**: ChatGPT plans; call `run_workflow`, then `workflow_status` until terminal. Disconnecting MCP does **not** kill the worker. + +## Worked mini-examples + +**1. Parallel review** + +```js +export const meta = { name: 'p-review', description: 'Two reviewers' } +const [a, b] = await parallel([ + () => agent('Correctness review of the diff', { label: 'corr' }), + () => agent('Security review of the diff', { label: 'sec' }), +]) +return { a, b } +``` + +**2. Pipeline with schema** + +```js +export const meta = { name: 'pipe', description: 'Find then fix plan' } +return await pipeline( + args.files, + (file) => agent(`List bugs in ${file}`, { schema: { type: 'object', properties: { bugs: { type: 'array', items: { type: 'string' } } }, required: ['bugs'] } }), + (findings, file) => agent(`Plan fixes for ${file}: ${JSON.stringify(findings)}`), +) +``` + +**3. Isolation for parallel writers** + +```js +export const meta = { name: 'iso', description: 'Parallel mutators' } +await parallel([ + () => agent('Implement feature A in isolation', { isolation: 'worktree', label: 'a' }), + () => agent('Implement feature B in isolation', { isolation: 'worktree', label: 'b' }), +]) +// dirty worktrees preserved; compose via return text / shared follow-up +``` diff --git a/src/config.test.ts b/src/config.test.ts index 0b4f99a8..d4c7b90f 100644 --- a/src/config.test.ts +++ b/src/config.test.ts @@ -39,9 +39,14 @@ assert.equal(resolveSubagentsFlag({}, { DEVSPACE_SUBAGENTS: "1" }), true); const seededConfigDir = mkdtempSync(join(tmpdir(), "devspace-seeded-skills-test-")); const seededSkillPaths = ensureDevspaceDefaultSkills({ DEVSPACE_CONFIG_DIR: seededConfigDir }); -assert.deepEqual(seededSkillPaths, [join(seededConfigDir, "skills", "subagent-delegation", "SKILL.md")]); +assert.deepEqual(seededSkillPaths, [ + join(seededConfigDir, "skills", "subagent-delegation", "SKILL.md"), + join(seededConfigDir, "skills", "dynamic-workflows", "SKILL.md"), +]); assert.equal(existsSync(seededSkillPaths[0]), true); +assert.equal(existsSync(seededSkillPaths[1]), true); assert.match(readFileSync(seededSkillPaths[0], "utf8"), /name: subagent-delegation/); +assert.match(readFileSync(seededSkillPaths[1], "utf8"), /name: dynamic-workflows/); assert.deepEqual(ensureDevspaceDefaultSkills({ DEVSPACE_CONFIG_DIR: seededConfigDir }), []); assert.throws( diff --git a/src/skills.ts b/src/skills.ts index c1f146a9..36e413b6 100644 --- a/src/skills.ts +++ b/src/skills.ts @@ -22,26 +22,25 @@ export interface SkillReadResolution { } const SUBAGENT_DELEGATION_NAME = "subagent-delegation"; -const SUBAGENT_DELEGATION_SKILL = join(SUBAGENT_DELEGATION_NAME, "SKILL.md"); +const DYNAMIC_WORKFLOWS_NAME = "dynamic-workflows"; function bundledSkillsDir(): string { return fileURLToPath(new URL("../skills", import.meta.url)); } -function hasSubagentDelegationSkill(skillDir: string): boolean { - return existsSync(join(skillDir, SUBAGENT_DELEGATION_SKILL)); -} - +/** + * Always include the bundled skills root when subagents are enabled. + * Previously the whole dir was dropped if the user had seeded + * subagent-delegation — that hid later skills (e.g. dynamic-workflows). + * User/devspace copies still win via earlier path order + name collisions. + */ export function effectiveSkillPaths(config: ServerConfig, cwd: string): string[] { - const bundledSkills = bundledSkillsDir(); const defaultPathCandidates = [ join(homedir(), ".agents", "skills"), resolve(cwd, ".agents", "skills"), config.devspaceSkillsDir, join(config.agentDir, "skills"), - config.subagents && !hasSubagentDelegationSkill(config.devspaceSkillsDir) - ? bundledSkills - : undefined, + config.subagents ? bundledSkillsDir() : undefined, ]; const defaultPaths = defaultPathCandidates.filter( (path): path is string => path !== undefined && existsSync(path), @@ -73,11 +72,15 @@ export function loadWorkspaceSkills(config: ServerConfig, cwd: string): LoadedSk if (config.subagents) return result; + const gated = new Set([SUBAGENT_DELEGATION_NAME, DYNAMIC_WORKFLOWS_NAME]); return { - skills: result.skills.filter((skill) => skill.name !== SUBAGENT_DELEGATION_NAME), + skills: result.skills.filter((skill) => !gated.has(skill.name)), diagnostics: result.diagnostics.filter((diagnostic) => { const collision = diagnostic.collision; - return !(collision?.resourceType === "skill" && collision.name === SUBAGENT_DELEGATION_NAME); + return !( + collision?.resourceType === "skill" && + gated.has(collision.name) + ); }), }; } diff --git a/src/user-config.ts b/src/user-config.ts index c8da90b5..5371b04b 100644 --- a/src/user-config.ts +++ b/src/user-config.ts @@ -8,6 +8,7 @@ import { import { homedir } from "node:os"; import { dirname, join, resolve } from "node:path"; import { expandHomePath } from "./roots.js"; +import type { AgentProvidersConfig } from "./workflow-types.js"; export interface DevspaceUserConfig { host?: string; @@ -19,6 +20,8 @@ export interface DevspaceUserConfig { worktreeRoot?: string; agentDir?: string; subagents?: boolean; + /** Ordered enable-list for local agent providers used by workflows/subagents. */ + agentProviders?: AgentProvidersConfig; } export interface DevspaceAuthConfig { @@ -97,14 +100,23 @@ export function generateOwnerToken(): string { return randomBytes(32).toString("base64url"); } -export function ensureDevspaceDefaultSkills(env: NodeJS.ProcessEnv = process.env): string[] { - const targetPath = join(devspaceSkillsDir(env), "subagent-delegation", "SKILL.md"); - if (existsSync(targetPath)) return []; +const DEFAULT_SKILLS = ["subagent-delegation", "dynamic-workflows"] as const; - const sourcePath = new URL("../skills/subagent-delegation/SKILL.md", import.meta.url); - mkdirSync(dirname(targetPath), { recursive: true }); - writeFileSync(targetPath, readFileSync(sourcePath, "utf8"), { mode: 0o644 }); - return [targetPath]; +export function ensureDevspaceDefaultSkills(env: NodeJS.ProcessEnv = process.env): string[] { + const seeded: string[] = []; + for (const name of DEFAULT_SKILLS) { + const targetPath = join(devspaceSkillsDir(env), name, "SKILL.md"); + if (existsSync(targetPath)) continue; + const sourcePath = new URL(`../skills/${name}/SKILL.md`, import.meta.url); + try { + mkdirSync(dirname(targetPath), { recursive: true }); + writeFileSync(targetPath, readFileSync(sourcePath, "utf8"), { mode: 0o644 }); + seeded.push(targetPath); + } catch { + // skill may not exist in package yet; skip + } + } + return seeded; } export function resolveSubagentsFlag( From 349b52ce3c928c518970905449186a2fadd181aa Mon Sep 17 00:00:00 2001 From: Waishnav Date: Tue, 21 Jul 2026 17:42:38 +0530 Subject: [PATCH 2/3] feat(agents): add agentProviders config + init/doctor probe MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Load ordered enable-list from config/env, probe available providers on init and doctor, and filter workflow CLI providers by enabled∩live. Co-Authored-By: Claude --- src/cli.ts | 70 ++++++++++++++++++++++++++++++++++++++++++++- src/config.ts | 47 ++++++++++++++++++++++++++++++ src/workflow-cli.ts | 11 +++++-- 3 files changed, 124 insertions(+), 4 deletions(-) diff --git a/src/cli.ts b/src/cli.ts index 9cc41360..4789a324 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -15,11 +15,13 @@ import { runLocalAgentProvider } from "./local-agent-adapters.js"; import { isLocalAgentProvider, loadLocalAgentProfiles, + LOCAL_AGENT_PROVIDERS, type LocalAgentProfile, } from "./local-agent-profiles.js"; import { assertLocalAgentProviderAvailable, formatLocalAgentProviderAvailabilitySummary, + getLocalAgentProviderAvailabilitySnapshot, } from "./local-agent-availability.js"; import { formatAvailableLocalAgentTargets, @@ -169,12 +171,18 @@ async function runInit({ force }: { force: boolean }): Promise { validate: validateRequiredPublicBaseUrl, })); + const subagents = resolveSubagentsFlag(files.config); + const agentProviders = + subagents === true + ? probeAndBuildAgentProviders(files.config.agentProviders) + : files.config.agentProviders; const config: DevspaceUserConfig = { host: files.config.host ?? "127.0.0.1", port, allowedRoots, publicBaseUrl, - subagents: resolveSubagentsFlag(files.config), + subagents, + ...(agentProviders ? { agentProviders } : {}), }; const auth = { ownerToken: files.auth.ownerToken ?? generateOwnerToken(), @@ -277,11 +285,71 @@ async function runDoctor(): Promise { console.log(`Public MCP URL: ${new URL("/mcp", config.publicBaseUrl).toString()}`); console.log(`Allowed roots: ${config.allowedRoots.join(", ")}`); console.log(`Allowed hosts: ${config.allowedHosts.join(", ")}`); + console.log(`Subagents: ${config.subagents ? "enabled" : "disabled"}`); + if (config.subagents) { + const snapshot = getLocalAgentProviderAvailabilitySnapshot(); + console.log( + `Agent providers (live): ${formatLocalAgentProviderAvailabilitySummary(snapshot)}`, + ); + if (config.agentProviders) { + console.log( + `Agent providers (enabled): ${ + config.agentProviders.enabled.length + ? config.agentProviders.enabled.join(", ") + : "(empty — no providers)" + }`, + ); + if (config.agentProviders.detectedAt) { + console.log(`Agent providers last probe: ${config.agentProviders.detectedAt}`); + } + } else { + console.log("Agent providers (config): missing (compat = all available)"); + } + + // Refresh lastProbe write-back when subagents on and config exists + if (files.configExists) { + const refreshed = probeAndBuildAgentProviders(files.config.agentProviders); + writeDevspaceConfig({ + ...files.config, + agentProviders: { + // keep user enable-list if set; only refresh probe metadata + available adds when empty + enabled: + files.config.agentProviders?.enabled ?? refreshed.enabled, + detectedAt: refreshed.detectedAt, + lastProbe: refreshed.lastProbe, + }, + }); + console.log(`Agent providers probe written to ${files.configPath}`); + } + } } catch (error) { console.log(`Config status: ${error instanceof Error ? error.message : String(error)}`); } } +/** Probe PATH and build AgentProvidersConfig (available ids in product order). */ +function probeAndBuildAgentProviders( + existing?: DevspaceUserConfig["agentProviders"], +): NonNullable { + const snapshot = getLocalAgentProviderAvailabilitySnapshot(); + const available = new Set( + snapshot.filter((row) => row.available).map((row) => row.name), + ); + const enabled = + existing?.enabled && existing.enabled.length > 0 + ? existing.enabled.filter((id) => LOCAL_AGENT_PROVIDERS.includes(id as never)) + : LOCAL_AGENT_PROVIDERS.filter((id) => available.has(id)); + return { + enabled, + detectedAt: new Date().toISOString(), + lastProbe: snapshot.map((row) => ({ + id: row.name, + available: row.available, + detail: row.reason, + })), + }; +} + function runConfigCommand(args: string[]): void { const [subcommand, key, ...rest] = args; const files = loadDevspaceFiles(); diff --git a/src/config.ts b/src/config.ts index 4fc1bcbb..65d10a3e 100644 --- a/src/config.ts +++ b/src/config.ts @@ -4,6 +4,12 @@ import { expandHomePath } from "./roots.js"; import type { LoggingConfig, LogFormat, LogLevel } from "./logger.js"; import type { OAuthConfig } from "./oauth-provider.js"; import { devspaceAgentsDir, devspaceSkillsDir, loadDevspaceFiles } from "./user-config.js"; +import type { AgentProvidersConfig } from "./workflow-types.js"; +import { + isLocalAgentProvider, + LOCAL_AGENT_PROVIDERS, + type LocalAgentProvider, +} from "./local-agent-profiles.js"; export type ToolMode = "minimal" | "full" | "codex"; export type WidgetMode = "off" | "changes" | "full"; @@ -26,6 +32,12 @@ export interface ServerConfig { devspaceSkillsDir: string; devspaceAgentsDir: string; subagents: boolean; + /** + * Resolved enable-list for agent providers. + * Missing user config → undefined (compat: all live providers). + * Explicit empty → no providers. + */ + agentProviders?: AgentProvidersConfig; agentDir: string; logging: LoggingConfig; } @@ -234,11 +246,46 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): ServerConfig { env.DEVSPACE_SUBAGENTS === undefined ? files.config.subagents === true : parseBoolean(env.DEVSPACE_SUBAGENTS), + agentProviders: parseAgentProvidersConfig( + env.DEVSPACE_AGENT_PROVIDERS, + files.config.agentProviders, + ), agentDir: resolve(expandHomePath(env.DEVSPACE_AGENT_DIR ?? files.config.agentDir ?? defaultAgentDir())), logging: parseLoggingConfig(env), }; } +/** + * Env `DEVSPACE_AGENT_PROVIDERS=codex,claude` replaces enabled list. + * Missing config → undefined (compat all-available). + * Explicit enabled: [] stays empty. + */ +export function parseAgentProvidersConfig( + envValue: string | undefined, + fileConfig: AgentProvidersConfig | undefined, +): AgentProvidersConfig | undefined { + if (envValue !== undefined) { + const enabled = envValue + .split(",") + .map((entry) => entry.trim()) + .filter((entry): entry is LocalAgentProvider => isLocalAgentProvider(entry)); + return { enabled }; + } + if (!fileConfig) return undefined; + const enabled = (fileConfig.enabled ?? []).filter((id): id is LocalAgentProvider => + isLocalAgentProvider(id), + ); + return { + enabled, + detectedAt: fileConfig.detectedAt, + lastProbe: fileConfig.lastProbe, + }; +} + +export function defaultAgentProvidersOrder(): LocalAgentProvider[] { + return [...LOCAL_AGENT_PROVIDERS]; +} + function parsePublicBaseUrl(value: string): string { const parsed = new URL(value); parsed.hash = ""; diff --git a/src/workflow-cli.ts b/src/workflow-cli.ts index 620cd50c..3a7062db 100644 --- a/src/workflow-cli.ts +++ b/src/workflow-cli.ts @@ -281,7 +281,7 @@ export async function runWorkflowWorker( try { const source = await readFile(claimed.scriptPath, "utf8"); const parsed = parseWorkflowScript(source, { filename: claimed.scriptPath }); - const enabledProviders = resolveEnabledProviders(); + const enabledProviders = resolveEnabledProviders(config.agentProviders); const concurrency = resolveWorkflowConcurrency( parsed.meta.concurrency, availableParallelism(), @@ -474,10 +474,15 @@ function formatRunLine( return `${run.id} ${run.status} ${run.name}${err}`; } -function resolveEnabledProviders(): string[] { +function resolveEnabledProviders( + agentProviders?: ServerConfig["agentProviders"], +): string[] { const snapshot = getLocalAgentProviderAvailabilitySnapshot(); const live = new Set(snapshot.filter((row) => row.available).map((row) => row.name)); - return LOCAL_AGENT_PROVIDERS.filter((id) => live.has(id)); + if (!agentProviders) { + return LOCAL_AGENT_PROVIDERS.filter((id) => live.has(id)); + } + return agentProviders.enabled.filter((id) => live.has(id as never)); } function splitFlags(args: string[]): { From e1ec0499c729ad1727cd65b26213438fc334db0d Mon Sep 17 00:00:00 2001 From: Waishnav Date: Tue, 21 Jul 2026 17:42:38 +0530 Subject: [PATCH 3/3] feat(workflow): register MCP run/status/cancel tools Gate on config.subagents. run_workflow spawns the same detached worker as CLI; status long-polls journal events; cancel requests cooperative stop. Co-Authored-By: Claude --- src/server.ts | 5 + src/workflow-tools.ts | 317 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 322 insertions(+) create mode 100644 src/workflow-tools.ts diff --git a/src/server.ts b/src/server.ts index e47aeb1e..c44f5032 100644 --- a/src/server.ts +++ b/src/server.ts @@ -47,6 +47,7 @@ import { formatPathForPrompt } from "./skills.js"; import { createWorkspaceStore } from "./workspace-store.js"; import { formatAgentsPath, WorkspaceRegistry } from "./workspaces.js"; import { summarizeLocalAgentProfile } from "./local-agent-profiles.js"; +import { registerWorkflowTools } from "./workflow-tools.js"; import { formatLocalAgentProviderAvailabilitySummary, getLocalAgentProviderAvailabilitySnapshot, @@ -1594,6 +1595,10 @@ function createMcpServer( registerCodexProcessTools(server, config, workspaces, processSessions); } + if (config.subagents) { + registerWorkflowTools(server, config, workspaces); + } + return server; } diff --git a/src/workflow-tools.ts b/src/workflow-tools.ts new file mode 100644 index 00000000..8551701f --- /dev/null +++ b/src/workflow-tools.ts @@ -0,0 +1,317 @@ +import { fileURLToPath } from "node:url"; +import type { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; +import { registerAppTool } from "@modelcontextprotocol/ext-apps/server"; +import * as z from "zod/v4"; +import type { ServerConfig } from "./config.js"; +import type { WorkspaceRegistry } from "./workspaces.js"; +import { + persistWorkflowScript, + resolveNamedWorkflowScript, + readWorkflowScriptFile, +} from "./workflow-files.js"; +import { parseWorkflowScript } from "./workflow-script.js"; +import { createWorkflowStore } from "./workflow-store.js"; +import { + WORKFLOW_MCP_YIELD_MS, + WORKFLOW_LIMITS, + type AgentProvidersConfig, + type WorkflowEventRecord, + type WorkflowRunRecord, +} from "./workflow-types.js"; +import { resolveWorkspaceHead } from "./workflow-worktrees.js"; +import { spawnWorkflowWorkerFromCli } from "./workflow-cli.js"; +import { getLocalAgentProviderAvailabilitySnapshot } from "./local-agent-availability.js"; +import { + isLocalAgentProvider, + LOCAL_AGENT_PROVIDERS, + type LocalAgentProvider, +} from "./local-agent-profiles.js"; + +const WORKFLOW_API_CHEATSHEET = ` +Workflow scripts (JS only): + export const meta = { name, description, phases?, defaultProvider?, concurrency? } + agent(prompt, { label?, phase?, schema?, model?, effort?, provider?, isolation?: 'worktree' }) + parallel(thunks) → Array // barrier; throw → null + pipeline(items, ...stages) // no cross-item barrier + phase(title); log(msg); args; budget (stub) + workflow(name | { scriptPath }, args?) // nest depth 1 +Bans: Date.now(), Math.random(), new Date() without args. +No writeMode — teach RO vs write in prompts; isolation contains writes. +`.trim(); + +export function registerWorkflowTools( + server: McpServer, + config: ServerConfig, + workspaces: WorkspaceRegistry, +): void { + if (!config.subagents) return; + + registerAppTool( + server, + "run_workflow", + { + title: "Run workflow", + description: + `Start a DevSpace Dynamic Workflow in an open workspace. Prefer named scripts or short inline scripts. ` + + `Poll with workflow_status until terminal. Cancel with workflow_cancel. ${WORKFLOW_API_CHEATSHEET}`, + inputSchema: { + workspaceId: z.string().describe("Workspace id from open_workspace."), + script: z + .string() + .optional() + .describe("Inline workflow script source (export const meta = …)."), + name: z.string().optional().describe("Named workflow under .devspace/workflows/.js"), + resumeFromRunId: z.string().optional().describe("Prior run id to resume (new run + cache)."), + args: z.unknown().optional().describe("Args object/array passed to script as `args`."), + yieldTimeMs: z + .number() + .int() + .min(0) + .max(WORKFLOW_MCP_YIELD_MS) + .optional() + .describe(`Ms to wait for early completion (default 2000, max ${WORKFLOW_MCP_YIELD_MS}).`), + }, + annotations: { readOnlyHint: false }, + _meta: {}, + }, + async ({ workspaceId, script, name, resumeFromRunId, args, yieldTimeMs }) => { + const workspace = workspaces.getWorkspace(workspaceId); + const store = createWorkflowStore(config); + try { + const provided = [script, name, resumeFromRunId].filter((v) => v !== undefined); + if (provided.length !== 1) { + throw new Error("Provide exactly one of script, name, or resumeFromRunId"); + } + + 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 prior = store.getRun(resumeFromRunId); + if (!prior) throw new Error(`Unknown run: ${resumeFromRunId}`); + priorRunId = prior.id; + priorScriptPath = prior.scriptPath; + const resolved = await readWorkflowScriptFile(prior.scriptPath); + source = resolved.source; + scriptHash = prior.scriptHash; + nameHint = prior.name; + runSource = "resume"; + if (args === undefined && prior.argsJson && prior.argsJson !== "null") { + try { + args = JSON.parse(prior.argsJson); + } catch { + // keep undefined + } + } + } else if (name) { + const resolved = await resolveNamedWorkflowScript({ + name, + workspaceRoot: workspace.root, + stateDir: config.stateDir, + }); + source = resolved.source; + scriptHash = resolved.scriptHash; + nameHint = resolved.nameHint; + runSource = "named"; + } 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: priorScriptPath ?? "pending", + scriptHash, + workspaceRoot: workspace.root, + workspaceId, + argsJson: JSON.stringify(args ?? null), + resumedFromRunId: priorRunId, + baseSha, + }); + + const persisted = + priorScriptPath ?? + (await persistWorkflowScript({ + stateDir: config.stateDir, + runId: run.id, + source, + preferredName: parsed.meta.name || nameHint, + })); + if (!priorScriptPath) store.setScriptPath(run.id, persisted); + + const cliEntry = fileURLToPath( + import.meta.url.replace(/workflow-tools\.(ts|js)$/, "cli.$1"), + ); + spawnWorkflowWorkerFromCli(run.id, cliEntry); + + const yieldMs = yieldTimeMs ?? 2_000; + const page = await yieldEvents(store, run.id, 0, yieldMs); + return toolResult(page); + } finally { + store.close(); + } + }, + ); + + registerAppTool( + server, + "workflow_status", + { + title: "Workflow status", + description: "Drain events for a workflow run; optional long-poll yield.", + inputSchema: { + runId: z.string(), + sinceSeq: z.number().int().min(0).optional(), + yieldTimeMs: z + .number() + .int() + .min(0) + .max(WORKFLOW_MCP_YIELD_MS) + .optional() + .describe(`Long-poll ms (default 0, max ${WORKFLOW_MCP_YIELD_MS}).`), + }, + annotations: { readOnlyHint: true }, + _meta: {}, + }, + async ({ runId, sinceSeq, yieldTimeMs }) => { + const store = createWorkflowStore(config); + try { + if (!store.getRun(runId)) throw new Error(`Unknown workflow run: ${runId}`); + const page = await yieldEvents(store, runId, sinceSeq ?? 0, yieldTimeMs ?? 0); + return toolResult(page); + } finally { + store.close(); + } + }, + ); + + registerAppTool( + server, + "workflow_cancel", + { + title: "Cancel workflow", + description: "Request cooperative cancel of a running workflow.", + inputSchema: { + runId: z.string(), + }, + annotations: { readOnlyHint: false }, + _meta: {}, + }, + async ({ runId }) => { + const store = createWorkflowStore(config); + try { + const run = store.requestCancel(runId); + if (run.pid && (run.status === "running" || run.status === "starting")) { + try { + process.kill(run.pid, "SIGTERM"); + } catch { + // already gone + } + } + const latest = store.getRun(runId)!; + return { + content: [{ type: "text" as const, text: JSON.stringify({ runId, status: latest.status }) }], + structuredContent: { runId, status: latest.status }, + }; + } finally { + store.close(); + } + }, + ); + + } + +async function yieldEvents( + store: ReturnType, + runId: string, + sinceSeq: number, + yieldMs: number, +): Promise<{ + run: WorkflowRunRecord; + events: WorkflowEventRecord[]; + nextSeq: number; + terminal: boolean; +}> { + const deadline = Date.now() + Math.min(yieldMs, WORKFLOW_MCP_YIELD_MS); + let cursor = sinceSeq; + let events: WorkflowEventRecord[] = []; + let terminal = false; + let run = store.getRun(runId)!; + + for (;;) { + const page = store.drainEvents(runId, cursor, WORKFLOW_LIMITS.eventDrainDefault); + events = events.concat(page.events); + cursor = page.nextSeq; + terminal = page.terminal; + run = page.run; + if (terminal || Date.now() >= deadline) break; + await sleep(250); + } + + return { run, events, nextSeq: cursor, terminal }; +} + +function toolResult(page: { + run: WorkflowRunRecord; + events: WorkflowEventRecord[]; + nextSeq: number; + terminal: boolean; +}) { + const payload = { + runId: page.run.id, + status: page.run.status, + events: page.events.map((e) => ({ + seq: e.seq, + type: e.type, + phase: e.phase, + label: e.label, + dataJson: e.dataJson, + })), + nextSeq: page.nextSeq, + result: page.run.resultJson ? safeJson(page.run.resultJson) : undefined, + error: page.run.error, + errorKind: page.run.errorKind, + }; + return { + content: [{ type: "text" as const, text: JSON.stringify(payload, null, 2) }], + structuredContent: payload, + }; +} + +function safeJson(text: string): unknown { + try { + return JSON.parse(text); + } catch { + return text; + } +} + +function sleep(ms: number): Promise { + return new Promise((r) => setTimeout(r, ms)); +} + +/** Resolve enabled ∩ live providers for workflows. */ +export function resolveWorkflowEnabledProviders( + agentProviders: AgentProvidersConfig | undefined, +): LocalAgentProvider[] { + const snapshot = getLocalAgentProviderAvailabilitySnapshot(); + const live = new Set( + snapshot.filter((row) => row.available).map((row) => row.name), + ); + if (!agentProviders) { + return LOCAL_AGENT_PROVIDERS.filter((id) => live.has(id)); + } + return agentProviders.enabled.filter( + (id): id is LocalAgentProvider => isLocalAgentProvider(id) && live.has(id), + ); +}