/** * dynamic-workflow-context.ts — WorkflowCtx facade for dynamic-workflow scripts (P2). * * Spec: research-findings/goal-workflow/00-SPEC.md §3.2 * Plan: 07-PLAN.md v3 P2 + §0b G4 + §0c C4/C5/C7. * * The `ctx` object passed to a `.dwf.ts` script's `export default async function(ctx)`. * Capability-locked: exposes ONLY the documented methods (no raw manifest/process/require). * The script host (dynamic-workflow-runner.ts) loads the script via jiti in plain module * scope with a FROZEN WorkflowCtx. v1 has NO vm sandbox (review H-2): the script CAN * reach `process`/`require`/`import` directly — the frozen ctx is a contract surface, * not a security boundary. `.dwf.ts` = postinstall-equivalent trust. isolated-vm v1.5. * * `agent()` resolution (§0b G4): 4-tier precedence * 1. opts.agent (explicit name) — bypasses team lookup * 2. team.roles.find(r => r.name === role)?.agent → allAgents lookup * 3. allAgents(discoverAgents(cwd)).find(a => a.name === role) (role name == agent name) * 4. synthesize minimal AgentConfig (source:"dynamic", systemPrompt:"You are {role}.") * * Isolation (§0b G3 / report 05 §C.4): worker output → artifact file; `agent()` returns * structured data + writes a side artifact. The script holds results in JS vars; only * `setResult()` reaches the main context. */ import { randomBytes } from "node:crypto"; import type { TSchema } from "@sinclair/typebox"; import type { AgentConfig } from "../../agents/agent-config.ts"; import { allAgents, discoverAgents } from "../../agents/discover-agents.ts"; import { appendMailboxMessage, readMailbox } from "../../state/coordination/mailbox.ts"; import { appendEvent } from "../../state/event-log/event-log.ts"; import { writeArtifact } from "../../state/stores/artifact-store.ts"; import type { TeamRunManifest } from "../../state/types.ts"; import type { TeamConfig } from "../../teams/team-config.ts"; import { logInternalError } from "../../utils/internal-error.ts"; import { safeAbort } from "../../utils/safe-abort.ts"; import { cleanupAgentWorktreeAsync, prepareAgentWorktreeAsync } from "../../worktree/worktree-manager.ts"; import type { DwfCheckpointState } from "../dwf-state-store.ts"; import { parsePiJsonOutput } from "../output/pi-json-output.ts"; import { extractStructuredResult } from "../output/result-extractor.ts"; import { executeWithRetry } from "../recovery/retry-executor.ts"; import { runWorker } from "../run-worker.ts"; import { mapConcurrent } from "../scheduling/parallel-utils.ts"; import { Semaphore } from "../scheduling/semaphore.ts"; import { estimateTokens } from "../task-runner/prompt-builder.ts"; import { renderPlanTemplate } from "./plan-templates.ts"; export interface AgentCallOpts { prompt: string; /** Role name (resolved via G4 4-tier chain) OR explicit agent name. */ role?: string; /** Explicit agent name — bypasses team-role lookup (tier 1). */ agent?: string; description?: string; model?: string; skill?: string[] | false; maxTurns?: number; graceTurns?: number; /** Dependency artifact paths injected into the agent prompt. */ inputs?: string[]; /** Disable ALL tools for this call (Pi `--no-tools`, §0c C6). Use for pure-judgment / * verdict steps where the agent must answer directly without exploring, e.g. * `ctx.review()`'s JSON-verdict call. Without this, role-based tools (read/grep/bash) * apply and the model may loop exploring instead of answering. */ disableTools?: boolean; /** Override the resolved agent's system prompt. Use when the call needs a different * persona/output-format than the role's defined agent — e.g. `ctx.review()` needs a * JSON-verdict judge, but the user's reviewer.md agent is a markdown code-reviewer. * When set, the resolved agent's systemPrompt is replaced entirely. */ systemPrompt?: string; /** Round-13 P0-3: optional TypeBox schema. When set, the call's output is validated * against the schema after extraction. Validation failure yields ok:false with a * structured `error` and undefined `structured` field. Forward-compatible: when * undefined, behavior is identical to the regex-based extractor. */ schema?: TSchema; /** round-17 P2-4: spawn this agent in an isolated git worktree. * Useful when parallel agents modify files concurrently (avoids conflicts). The * worktree is created from HEAD, the agent runs there, and on completion the * diff is captured as an artifact before cleanup. Default false. * If worktree creation fails (no git repo, dirty leader), the agent runs in the * normal cwd and a warning is logged via ctx.log(). Backward compatible — * omitting it is identical to `false`. */ worktree?: boolean; } export interface AgentResult { ok: boolean; text: string; structured?: unknown; usage?: { input?: number; output?: number; cost?: number; turns?: number }; runId?: string; taskId?: string; artifactPath?: string; error?: string; durationMs?: number; } /** round-14 P1-2: per-workflow token budget. Frozen read-only surface exposed as ctx.budget. */ export interface WorkflowBudget { /** Configured budget, or null when unbounded. */ total: number | null; /** Tokens consumed so far (accumulated from each ctx.agent() run's usage). */ spent(): number; /** Tokens remaining; Infinity when total is null. */ remaining(): number; } /** SDD-3 W-C G13: default cap on ctx.agent() invocations per DWF run. * * Rationale: the largest real workflows observed (distill pipelines, adaptive * plans) use ~10-40 agent calls; 200 gives ≥5× headroom so a legitimate run * never trips it. Cost-wise 200 calls × ~5k tokens/call bounds a no-budget run * to roughly a million tokens — an explainable blast radius. For real spawns * the 30-min script timeout usually binds first; this cap's job is the * no-budget runaway class, which previously had NO bound on the number of * spawned steps and died only via the blind timeout (upgrade plan §2 G13). */ export const DEFAULT_MAX_AGENT_CALLS = 200; /** SDD-3 W-C G13: structured termination error thrown when the agent-call cap trips. * Propagates out of ctx.agent() (and therefore out of the .dwf.ts script unless * the script catches it) so the runner fails the run with a structured reason * (dwf.failed) instead of waiting out the blind script timeout. */ export class DwfAgentCallCapError extends Error { /** Effective cap that was tripped. */ readonly limit: number; /** Completed agent() invocations when the cap tripped. */ readonly used: number; constructor(limit: number, used: number) { super(`dynamic workflow agent-call cap reached (${used}/${limit} calls; raise the run's maxAgentCalls to allow more)`); this.name = "DwfAgentCallCapError"; this.limit = limit; this.used = used; } } export interface WorkflowCtx { cwd: string; runId: string; goal?: string; /** Spawn one agent, await result. Concurrency enforced by ctx.semaphore. */ agent(opts: AgentCallOpts): Promise; /** Bounded fan-out preserving order (wraps mapConcurrent). */ fanOut(items: T[], limit: number, fn: (item: T, i: number) => Promise): Promise; /** Pipeline: sequential per-item stages, parallel across items (bounded by * ctx.semaphore). Each item passes through all stages in order; different * items may run concurrently. A failed stage yields `null` for that item * (logged via ctx.log) and other items continue. Aborts propagate. * round-16 (P2-1). */ pipeline( items: TItem[], ...stages: Array<(previous: TResult, original: TItem, index: number) => Promise | TResult> ): Promise<(TResult | null)[]>; /** Run a reviewer agent over an artifact; parse {outcome, feedback}. §3.2. */ review( taskId: string, reviewerRole?: string, opts?: { content?: string; artifactPath?: string; disableTools?: boolean; }, ): Promise<{ outcome: "accept" | "reject" | "changes_requested"; feedback: string; }>; /** Re-run a task with feedback (wraps executeWithRetry). */ retry(taskId: string, opts?: { feedback?: string }): Promise; /** Send a mailbox message to another agent/leader. */ mail( to: string, body: string, opts?: { kind?: string; taskId?: string; replyTo?: string; replyDeadline?: number; }, ): string; /** Block until N mailbox replies arrive or deadline. ~10 LOC net-new (report 05 §G.4). */ gatherReplies(messageIds: string[], deadlineMs: number): Promise; /** Render a built-in plan template (full-implementation / standard-review). */ renderTemplate(name: string, vars: Record): unknown; /** Persistent variables (revived intermediate-store). */ vars: Record; /** Mark the final result. ONLY this artifact reaches the main context. */ setResult(artifactPath: string, meta?: Record): void; /** Mark the start of a named workflow phase. Emits a `dwf.phase_started` event * (and a `dwf.phase_completed` for the previous phase, if any) to the run's * events.jsonl. Idempotent on the same title — calling twice with the same * title is a no-op. Phase titles are in-memory only; the events log is the * durable source of truth for phase boundaries. */ phase(title: string): void; /** round-14 P1-3: append a workflow-level log line. Persists to events.jsonl * as a `dwf.log` event and keeps a bounded in-memory copy (capped at 1000). */ log(message: unknown): void; /** round-14 P1-2: per-workflow token budget. ctx.agent() auto-rejects with * ok:false once exhausted. */ budget: WorkflowBudget; /** round-14 P1-5: typed workflow arguments. Reads the value passed via * MakeWorkflowCtxOptions.args (sourced from manifest.args). Defaults to {} * when unset. */ args(): T; semaphore: Semaphore; /** Abort signal (cancel/stop). */ signal: AbortSignal; } export interface MakeWorkflowCtxOptions { concurrency?: number; signal: AbortSignal; team?: TeamConfig; modelOverride?: string; /** round-14 P1-2: per-workflow token budget. null/undefined = unbounded. */ tokenBudget?: number | null; /** SDD-3 W-C G13: cap on ctx.agent() invocations (new spawns AND cached * replays) per run. null/undefined = DEFAULT_MAX_AGENT_CALLS — a run is * always bounded even with no tokenBudget configured. */ maxAgentCalls?: number | null; /** round-14 P1-5: typed workflow arguments (sourced from manifest.args). Defaults to {}. */ args?: unknown; /** round-18 P2-3: checkpoint state to hydrate ctx with on resume. When provided, * the ctx starts with the resumed vars/phases/logs/spent/agentCount instead of * empty defaults. Omit (or undefined) for a fresh run — backward compatible. */ resumedState?: DwfCheckpointState; /** round-18 P2-3: callback invoked after each `ctx.agent()` call completes * (success OR fail). The runner wires this to `DwfStore.save()` so a crash after * an agent call leaves a durable checkpoint. Best-effort — failures are swallowed * so checkpointing can never crash the workflow. */ onCheckpoint?: (state: DwfCheckpointState) => void; } /** * Resolve a role/agent name to a full AgentConfig (§0b G4 4-tier precedence). * Module-local — NOT promoted to a shared module (keeps P2 isolated from the * load-bearing team-runner path). */ export function resolveAgentForRole( roleName: string | undefined, opts: { explicitAgent?: string; team?: TeamConfig; cwd: string }, ): AgentConfig { const cwd = opts.cwd; // Tier 1: explicit agent name. if (opts.explicitAgent) { const found = allAgents(discoverAgents(cwd)).find((a) => a.name === opts.explicitAgent); if (found) return found; // Fall through to synthesize if the named agent doesn't exist (P2-friendly). } // Tier 2: team.roles[].agent lookup. if (opts.team) { const role = opts.team.roles.find((r) => r.name === roleName); if (role) { const byAgentName = allAgents(discoverAgents(cwd)).find((a) => a.name === role.agent); if (byAgentName) return byAgentName; } } // Tier 3: discoverAgents by role name (role name == agent name). if (roleName) { const byRoleName = allAgents(discoverAgents(cwd)).find((a) => a.name === roleName); if (byRoleName) return byRoleName; } // Tier 4: synthesize a minimal AgentConfig. const name = opts.explicitAgent ?? roleName ?? "executor"; return synthesizeAgentConfig(name); } /** Synthesize a minimal AgentConfig (§0c C7: source:"dynamic", not "synthetic"). */ export function synthesizeAgentConfig(name: string, model?: string): AgentConfig { return { name, description: `Synthesized agent for dynamic workflow (${name}).`, source: "dynamic", filePath: ``, systemPrompt: `You are ${name}, an agent in a dynamic pi-crew workflow. Use the provided tools (read, grep, find, ls, bash) to investigate the target and produce concrete written findings. Do not return an empty response — always output substantive content for your task.`, model, // Round-N fix: give synthesized agents the STANDARD file-investigation toolkit // (matches real built-in agents). Previously tools:[] relied on pi-args default // behavior, which left most synthesized agents (explorer/analyst/critic/executor) // unable to read the codebase → 10/11 agents returned empty/!ok in distill-dwf runs. tools: ["read", "grep", "find", "ls", "bash"], inheritProjectContext: false, inheritSkills: false, }; } /** Build the WorkflowCtx facade. Capability-locked: only documented methods exposed. */ export function makeWorkflowCtx(manifest: TeamRunManifest, opts: MakeWorkflowCtxOptions): WorkflowCtx { const concurrency = Math.max(1, opts.concurrency ?? 4); const semaphore = new Semaphore(concurrency); let finalResult: { artifactPath: string; meta?: Record } | undefined; // SDD-3 W-C G13: effective agent-call cap. Always ≥ 1 — omitting every knob // still leaves the run bounded (default DEFAULT_MAX_AGENT_CALLS). A cap error // aborts capController below so in-flight children are killed and pipeline() // rethrows the cap error instead of swallowing it into a null item. const maxAgentCalls = typeof opts.maxAgentCalls === "number" && Number.isFinite(opts.maxAgentCalls) && opts.maxAgentCalls >= 1 ? Math.floor(opts.maxAgentCalls) : DEFAULT_MAX_AGENT_CALLS; const capController = new AbortController(); // G13: combine the caller's signal with the cap controller so EITHER source // (external abort OR cap trip) reaches runWorker and ctx.signal consumers. // AbortSignal.any is available since Node 20.3 (project requires Node 22+); // an already-aborted input yields an already-aborted combined signal. const ctxSignal = AbortSignal.any([opts.signal, capController.signal]); // round-18 P2-3: agent invocation counter. Hydrated from a resumed checkpoint so a // resumed run keeps an accurate count; incremented in agent()'s finally block. let agentCount = opts.resumedState ? opts.resumedState.agentCount : 0; // round-12 P0-1: in-memory phase state, exposed via non-enumerable getter like __finalResult. // The events log is the durable source of truth for phase boundaries. // round-18 P2-3: hydrate phaseState from a resumed checkpoint (backward compatible when unset). const phaseState: { currentPhase: string | undefined; phases: string[] } = opts.resumedState ? { currentPhase: opts.resumedState.currentPhase, phases: [...opts.resumedState.phases], } : { currentPhase: undefined, phases: [] }; let phaseCapWarned = false; // round-14 P1-2/P1-3/P1-5: closure-scoped runtime state shared by budget/log/args. // Mirrors the pi-dynamic-workflows RuntimeState pattern (workflow.ts:state). // round-18 P2-3: hydrate spent/logs from a resumed checkpoint (backward compatible when unset). const wfState: { spent: number; logs: string[]; args: unknown } = { spent: opts.resumedState?.spent ?? 0, logs: opts.resumedState ? [...opts.resumedState.logs].slice(0, 1000) : [], args: opts.args ?? {}, }; // PERS-1: per-agent-call idempotency cache. Hydrated from resumedState so a DWF // resume skips already-completed agent calls (avoids duplicating artifacts/mailbox/tokens). // Uses Record (not Map) for JSON serialization safety. const completedAgentCalls: Record = opts.resumedState?.completedAgentCalls ?? {}; // round-14 P1-2: frozen budget surface. The closures read wfState.spent so the // object stays live after Object.freeze(ctx). total is a snapshot primitive. const budget = Object.freeze({ total: opts.tokenBudget ?? null, spent: () => wfState.spent, remaining: () => (opts.tokenBudget == null ? Infinity : Math.max(0, opts.tokenBudget - wfState.spent)), } satisfies WorkflowBudget); const ctx: WorkflowCtx = { cwd: manifest.cwd, runId: manifest.runId, goal: manifest.goal, signal: ctxSignal, semaphore, async agent(call: AgentCallOpts): Promise { await semaphore.acquire(); const started = Date.now(); // round-17 P2-4: declared before the try so the finally can clean it up // regardless of which return/throw path is taken. let worktreePath: string | undefined; let worktreeBranch: string | undefined; // BDG-2 + G14 (SDD-3): content-based reserve estimate — replaces the flat // ESTIMATE=4096. estimateTokens() (chars/4, the in-tree heuristic exported // from runtime/task-runner/prompt-builder.ts — the same estimator // pre-execution.ts uses) runs over the call's prompt + optional // systemPrompt override, floored at 512 tokens for the fixed per-call // overhead (schema/JSON directive, artifact headers) not present in the // prompt string. Short calls stop over-reserving; very long prompts // reserve proportionally so N concurrent calls cannot blow past the // budget between reserve and the post-run adjust. Every BDG-2 site // below (check/reserve/un-reserve/adjust/refund) uses THIS value, so // refund always matches what was reserved. const estimate = Math.max(512, estimateTokens((call.prompt?.length ?? 0) + (call.systemPrompt?.length ?? 0))); let reserved = false; try { // SDD-3 W-C G13: agent-call cap — the structured bound for no-budget runs. // Checked BEFORE the PERS-1 cache lookup on purpose: cached replays are // free but must still count, so a spin loop over the same prompt cannot // bypass the cap (the default 200 leaves ≥5× headroom over the largest // real workflows, so resume replays never trip it legitimately). Counts // completed invocations; up to concurrency-1 extra calls may already be // in flight when it trips (they are killed by the cap abort below). if (agentCount >= maxAgentCalls) { // Durable record (dwf.log is the registered workflow-log event type). appendEvent(manifest.eventsPath, { type: "dwf.log", runId: manifest.runId, data: { message: `agent-call cap reached: ${agentCount}/${maxAgentCalls} — terminating workflow (raise maxAgentCalls to allow more)`, limit: maxAgentCalls, used: agentCount, }, }); // Kill in-flight children + make pipeline()/gatherReplies() see the abort. safeAbort(capController, "dwf-agent-call-cap"); throw new DwfAgentCallCapError(maxAgentCalls, agentCount); } // PERS-1: per-agent-call idempotency. Compute a deterministic call ID // from the call arguments and check if this call was already completed // (e.g. during a previous run before a crash). If so, return the cached // result without re-spawning the agent. const callId = JSON.stringify({ role: call.role, agent: call.agent, prompt: call.prompt, model: call.model, schema: call.schema ? "yes" : "no", }); const cached = completedAgentCalls[callId]; if (cached) { return { ok: true, text: cached.text, usage: cached.usage, durationMs: 0, }; } // BDG-2: reserve-then-adjust budget. Before spawning, estimate the cost and // reserve it by adding to wfState.spent. This prevents N concurrent calls from // all seeing the same remaining budget and overspending. if (budget.total !== null && budget.remaining() < estimate) { return { ok: false, text: "", error: "workflow token budget exhausted", durationMs: 0, }; } // Reserve the estimate before spawning. wfState.spent += estimate; reserved = true; const agentConfig = resolveAgentForRole(call.role, { explicitAgent: call.agent, team: opts.team, cwd: manifest.cwd, }); // §0c C6: per-call disableTools override. When set, force Pi `--no-tools` so the // agent answers directly without exploring. Applied AFTER role resolution so it // wins over any role-defined tools. let effectiveAgent = call.disableTools === true ? { ...agentConfig, disableTools: true, tools: [] } : agentConfig; // Per-call systemPrompt override (replaces the resolved agent's persona/output-format). // Used by ctx.review() to force a JSON-verdict judge instead of the role's markdown reviewer. // Round-13 P0-3: when a schema is provided, append a JSON-output instruction so // the model returns parseable JSON instead of prose. Schema name is intentionally // generic — we don't reveal TypeBox internal types. // // Smoke-test fix: when BOTH schema AND an explicit call.systemPrompt are set, // the call.systemPrompt is the caller's intended persona (e.g. a JSON-verdict // judge). It MUST be used as the base for the JSON instruction — otherwise the // role's persona leaks through and the model returns prose, failing schema // validation. Previously call.systemPrompt was silently dropped when a schema // was present, which confused models into returning text like "hello". if (call.schema !== undefined) { const base = call.systemPrompt ?? effectiveAgent.systemPrompt; effectiveAgent = { ...effectiveAgent, systemPrompt: composeSchemaSystemPrompt(base, call.schema), }; } else if (call.systemPrompt !== undefined) { effectiveAgent = { ...effectiveAgent, systemPrompt: call.systemPrompt, }; } const task = composeAgentTask(call); // round-17 P2-4: worktree isolation per agent. When requested, spawn the // agent in an isolated git worktree so parallel file-modifying agents // don't clobber each other. Falls back to the normal cwd (with a warning) // when worktree creation is unavailable (no git repo, dirty leader). let agentCwd = manifest.cwd; if (call.worktree === true) { const wt = await prepareAgentWorktreeAsync(manifest, `dwf-agent-${Date.now()}-${randomBytes(4).toString("hex")}`); if (wt?.worktreePath) { agentCwd = wt.cwd; worktreePath = wt.worktreePath; worktreeBranch = wt.branch; ctx.log(`worktree: agent isolated at ${wt.worktreePath}`); } else { ctx.log("worktree: creation unavailable — falling back to normal cwd"); } } const childResult = await runWorker({ cwd: agentCwd, task, agent: effectiveAgent, model: call.model ?? opts.modelOverride ?? agentConfig.model, skillPaths: undefined, // skills resolved via agent config + team-role plumbing maxTurns: call.maxTurns, graceTurns: call.graceTurns, signal: ctxSignal, artifactsRoot: manifest.artifactsRoot, runId: manifest.runId, role: call.role ?? call.agent, }); if (childResult.exitCode !== 0 || childResult.error) { // BDG-2: un-reserve the estimate on spawn failure. wfState.spent -= estimate; reserved = false; return { ok: false, text: "", error: childResult.error ?? `exit ${childResult.exitCode}`, durationMs: Date.now() - started, }; } const parsed = parsePiJsonOutput(childResult.stdout); // round-14 P1-2: accumulate this run's token usage into the workflow budget. // BDG-2: adjust the reserve — subtract the estimate, add the actual usage. // This correctly reduces spent when actualUsage < estimate. wfState.spent += (parsed.usage?.input ?? 0) + (parsed.usage?.output ?? 0) - estimate; reserved = false; let text = childResult.rawFinalText || parsed.finalText || ""; // Round-11 test fix: parsePiJsonOutput only extracts text from pi event stream // ({type:"message_end", message:{role:"assistant", content:[...]}}). When the // agent emits plain JSON, plain text, or a different format, finalText is empty. // Fallback to a more permissive extraction that handles multiple output shapes. if (!text.trim()) { text = extractTextFallback(childResult.stdout); } // Round-13 P0-3: schema validation post-extraction. The schema option is // additive — when undefined the call site is unchanged. With a schema, // extracted.error means the worker output didn't match expected shape and // the script should treat the result as failed (ok:false, error set). const extracted = extractStructuredResult(text, call.schema); // Write a side artifact for audit/isolation (§0b G3). const rel = `wf/${Date.now()}-${randomBytes(4).toString("hex")}.md`; const artifact = writeArtifact(manifest.artifactsRoot, { kind: "result", relativePath: rel, content: text, producer: "dynamic-workflow", }); if (call.schema !== undefined && !extracted.structured) { // BDG-2: un-reserve was already done above (reserved = false after adjust). return { ok: false, text, usage: parsed.usage, artifactPath: artifact.path, error: extracted.error ?? "structured output does not match schema", durationMs: Date.now() - started, }; } // PERS-1: cache the successful result for idempotency on resume. completedAgentCalls[callId] = { text, usage: { input: parsed.usage?.input ?? 0, output: parsed.usage?.output ?? 0 }, }; return { ok: true, text, structured: extracted.structured ? extracted.data : undefined, usage: parsed.usage, artifactPath: artifact.path, durationMs: Date.now() - started, }; } catch (error) { // BDG-2: un-reserve the estimate on failure (only if still reserved). if (reserved) wfState.spent -= estimate; // SDD-3 W-C G13: cap errors must TERMINATE the script with the // structured reason — degrading to ok:false would let scripts (and the // review()/retry()/pipeline fallbacks) swallow the cap and keep looping. if (error instanceof DwfAgentCallCapError) throw error; logInternalError("dynamic-workflow-context.agent", error, `runId=${manifest.runId}`); return { ok: false, text: "", error: error instanceof Error ? error.message : String(error), durationMs: Date.now() - started, }; } finally { // round-17 P2-4: clean up the worktree after the agent completes (success // OR failure). Captures the diff as an artifact before removal. Best-effort // — a leak must never crash the workflow. if (worktreePath) { try { await cleanupAgentWorktreeAsync(manifest, worktreePath, worktreeBranch); } catch (cleanupError) { logInternalError("dynamic-workflow-context.worktree-cleanup", cleanupError, `worktreePath=${worktreePath}`); } } // round-18 P2-3: checkpoint AFTER the agent completes (success or fail) so a // crash between agent calls leaves durable state to resume from. The counter is // incremented here (after the call) so the checkpoint reflects the call that ran. agentCount++; if (opts.onCheckpoint) { try { opts.onCheckpoint({ runId: manifest.runId, vars: ctx.vars, phases: phaseState.phases, currentPhase: phaseState.currentPhase, logs: wfState.logs.slice(0, 1000), spent: wfState.spent, agentCount, completedAgentCalls, updatedAt: new Date().toISOString(), }); } catch (checkpointError) { logInternalError("dynamic-workflow-context.checkpoint", checkpointError, `runId=${manifest.runId}`); } } semaphore.release(); } }, async fanOut(items: T[], limit: number, fn: (item: T, i: number) => Promise): Promise { return mapConcurrent(items, Math.max(1, limit), fn); }, async pipeline( items: TItem[], ...stages: Array<(previous: TResult, original: TItem, index: number) => Promise | TResult> ): Promise<(TResult | null)[]> { if (!Array.isArray(items)) { throw new TypeError("pipeline() expects an array as the first argument"); } if (stages.length === 0 || stages.some((s) => typeof s !== "function")) { throw new TypeError("pipeline() stages must be functions"); } if (items.length === 0) return []; // Parallel across items, bounded by the workflow concurrency (mirrors fanOut). // Per-item stages run sequentially. A failed stage yields null for that item // (logged via ctx.log) and the remaining items continue. Aborts propagate. return mapConcurrent(items, concurrency, async (item, index): Promise => { let value: unknown = item; for (const stage of stages) { try { value = await stage(value as TResult, item, index); } catch (error) { // G13: ctxSignal is aborted on cap trip — rethrow so the structured cap // error terminates the workflow instead of degrading to a null item. if (ctxSignal.aborted) throw error; ctx.log(`pipeline[${index}] failed: ${error instanceof Error ? error.message : String(error)}`); return null; } } return value as TResult; }); }, async review( taskId: string, reviewerRole = "reviewer", reviewOpts?: { content?: string; artifactPath?: string; disableTools?: boolean; }, ): Promise<{ outcome: "accept" | "reject" | "changes_requested"; feedback: string; }> { // review() is a VERDICT step: it must produce a parseable JSON {outcome, feedback}, not a // free-form markdown review. The resolved reviewer agent (e.g. ~/.pi/agent/agents/reviewer.md) // has tools (read/grep/bash) + a markdown-output system prompt. Without disableTools, the // reviewer explores the repo looking for the task's work, loops, and gets killed (exit 143) // before producing JSON — leaving text="" and the fallback verdict. Default: disableTools so // the reviewer judges the provided content (or taskId context) directly. const disableTools = reviewOpts?.disableTools !== false; // default true const workContext = reviewOpts?.content ? `\n\nWork to review:\n"""\n${reviewOpts.content}\n"""` : reviewOpts?.artifactPath ? `\n\nRead the work from artifact: ${reviewOpts.artifactPath}` : ""; const res = await ctx.agent({ role: reviewerRole, prompt: `You are reviewing the work for task '${taskId}'.${workContext}\n\nEvaluate the work and respond with ONLY a single JSON object, no prose, no markdown:\n{"outcome":"accept|reject|changes_requested","feedback":""}\n\n- "accept": work is complete and correct.\n- "reject": work is fundamentally wrong.\n- "changes_requested": work needs revision (explain what in feedback).`, maxTurns: 3, disableTools, systemPrompt: 'You are a JSON verdict judge. You output ONLY a single JSON object with keys "outcome" (one of accept/reject/changes_requested) and "feedback" (a concise explanation). Never output prose, markdown, or code fences. Begin your response with { and end with }.', }); const extracted = res.structured as { outcome?: string; feedback?: string } | undefined; if (extracted && typeof extracted.outcome === "string" && typeof extracted.feedback === "string") { const outcome = extracted.outcome === "accept" || extracted.outcome === "reject" || extracted.outcome === "changes_requested" ? extracted.outcome : "changes_requested"; return { outcome, feedback: extracted.feedback }; } // Fallback (round-11 runtime): many models (e.g. MiniMax-M3) ignore JSON-output // instructions and produce a prose review instead. Rather than report an // unparseable verdict, run a tiny judge call that converts the prose review into a // JSON verdict. This guarantees ctx.review() always returns a structured verdict // regardless of the reviewer's output format. Skipped when the reviewer produced // no text at all (genuine failure). if (res.text.trim()) { const judge = await ctx.agent({ role: reviewerRole, prompt: `Convert the following code review into a verdict JSON. Read the review and decide the outcome.\n\nREVIEW:\n"""\n${res.text.slice(0, 4000)}\n"""\n\nRespond with ONLY a JSON object:\n{"outcome":"accept|reject|changes_requested","feedback":""}\n- accept: review found no real issues.\n- reject: review found critical/fundamental problems.\n- changes_requested: review found issues that need fixing.`, maxTurns: 1, disableTools: true, systemPrompt: "You output ONLY a single JSON object with keys outcome and feedback. Begin with { and end with }. Never output prose.", }); const judged = judge.structured as { outcome?: string; feedback?: string } | undefined; if (judged && typeof judged.outcome === "string" && typeof judged.feedback === "string") { const outcome = judged.outcome === "accept" || judged.outcome === "reject" || judged.outcome === "changes_requested" ? judged.outcome : "changes_requested"; return { outcome, feedback: judged.feedback }; } } // Tier-3 sentiment fallback (round-11): when neither the reviewer nor the judge // produced JSON (common with MiniMax-M3, GLM, which ignore JSON-output // instructions), classify the outcome from the REVIEWER's prose sentiment. We use // the reviewer's text (not the judge's terse output) because the original review is // the richest sentiment signal. This keeps outcome ACCURATE (accept vs reject vs // changes_requested) even when no JSON is ever produced — without it, outcome was // always the hardcoded 'changes_requested' default (e.g. correct code was // misclassified as needing changes). if (res.text.trim()) { return { outcome: classifyReviewOutcome(res.text), feedback: res.text, }; } return { outcome: "changes_requested", feedback: res.text || "(reviewer produced no parseable verdict)", }; }, async retry(taskId: string, retryOpts?: { feedback?: string }): Promise { return executeWithRetry( async () => ctx.agent({ role: "executor", prompt: `Re-do task '${taskId}'.${retryOpts?.feedback ? ` Feedback: ${retryOpts.feedback}` : ""}`, }), { maxAttempts: 3, backoffMs: 0, jitterRatio: 0, exponentialFactor: 1, }, ); }, mail( to: string, body: string, mailOpts?: { kind?: string; taskId?: string; replyTo?: string; replyDeadline?: number; }, ): string { const msg = appendMailboxMessage(manifest, { direction: "outbox", from: "dynamic-workflow", to, body, kind: (mailOpts?.kind as never) ?? "message", taskId: mailOpts?.taskId, replyTo: mailOpts?.replyTo, replyDeadline: mailOpts?.replyDeadline, }); return msg.id; }, async gatherReplies(messageIds: string[], deadlineMs: number): Promise { const deadline = Date.now() + deadlineMs; while (Date.now() < deadline) { const inbox = readMailbox(manifest, "inbox"); const got = inbox.filter((m) => m.replyTo && messageIds.includes(m.replyTo)); if (got.length >= messageIds.length) return got; await new Promise((r) => setTimeout(r, 500)); if (ctxSignal.aborted) return inbox.filter((m) => m.replyTo && messageIds.includes(m.replyTo)); } return readMailbox(manifest, "inbox").filter((m) => m.replyTo && messageIds.includes(m.replyTo)); }, renderTemplate(name: string, vars: Record): unknown { return renderPlanTemplate(name, vars); }, vars: opts.resumedState ? { ...opts.resumedState.vars } : ({} as Record), setResult(artifactPath: string, meta?: Record): void { finalResult = { artifactPath, meta }; }, phase(title: string): void { if (typeof title !== "string" || title.length === 0) { throw new TypeError("ctx.phase(title) requires a non-empty string title."); } // Idempotency: same phase title → no event, no state change. if (title === phaseState.currentPhase) return; // Close out the previous open phase BEFORE the new one opens. // REVIEW FIX (2026-09-10): reverted M2a buffered conversion — phase // transitions are low-frequency AND read back synchronously (tests, // checkpoint resume; dwf-setresult rounds 12/14/18 assert the events // file immediately after the run). if (phaseState.currentPhase !== undefined) { appendEvent(manifest.eventsPath, { type: "dwf.phase_completed", runId: manifest.runId, data: { phase: phaseState.currentPhase }, }); } phaseState.currentPhase = title; // Dedup append with hard cap to bound memory; events still flow. if (!phaseState.phases.includes(title)) { if (phaseState.phases.length < 100) { phaseState.phases.push(title); } else if (!phaseCapWarned) { phaseCapWarned = true; logInternalError( "dynamic-workflow-context.phase-cap", new Error( "Phase list cap of 100 reached; further phases still emit events but are not added to the in-memory phases[] list. Use the events log as the durable source of truth.", ), `runId=${manifest.runId}`, ); } } // REVIEW FIX (2026-09-10): reverted M2a buffered conversion — see phase_completed. appendEvent(manifest.eventsPath, { type: "dwf.phase_started", runId: manifest.runId, data: { phase: title }, }); }, budget, log(message: unknown): void { // round-14 P1-3: stringify non-strings, keep a bounded in-memory copy, and // always emit a dwf.log event (the events log is the durable source of truth). const text = typeof message === "string" ? message : JSON.stringify(message); if (wfState.logs.length < 1000) { wfState.logs.push(text); } // REVIEW FIX (2026-09-10): reverted M2a buffered conversion — see phase(). appendEvent(manifest.eventsPath, { type: "dwf.log", runId: manifest.runId, data: { message: text }, }); }, args(): T { // round-14 P1-5: typed workflow args sourced from manifest (via opts.args). return wfState.args as T; }, }; // Attach the final-result slot via a non-enumerable getter so the runner can read it // without exposing a mutation surface on the ctx the script sees. Object.defineProperty(ctx, "__finalResult", { get: () => finalResult, enumerable: false, }); // round-12 P0-1: phase state is read-only from the runner; the script can only mutate // it via ctx.phase(title), which is the documented public surface. Object.defineProperty(ctx, "__phaseState", { get: () => phaseState, enumerable: false, }); // round-14 P1-3: in-memory log buffer is read-only from the runner; the script can only // append via ctx.log(message). The events log remains the durable source of truth. Object.defineProperty(ctx, "__logs", { get: () => wfState.logs, enumerable: false, }); // round-18 P2-3: agent invocation counter is read-only from the runner. The script can // only advance it via ctx.agent() (incremented in agent()'s finally). Exposed so // getWorkflowCheckpoint() can report an accurate count. Object.defineProperty(ctx, "__agentCount", { get: () => agentCount, enumerable: false, }); // SDD-3 W-C G13: effective run limits are read-only from the runner (mirrors // __agentCount). Exposed so getWorkflowLimits() can report the effective cap. Object.defineProperty(ctx, "__limits", { get: () => ({ maxAgentCalls }), enumerable: false, }); // PERS-1: completedAgentCalls is read-only from the runner; the agent() method // is the only writer. Exposed so getWorkflowCheckpoint() can include it. Object.defineProperty(ctx, "__completedAgentCalls", { get: () => completedAgentCalls, enumerable: false, }); return ctx; } /** Read the final result set by the script (runner-only; not part of the public ctx surface). */ export function getWorkflowFinalResult(ctx: WorkflowCtx): { artifactPath: string; meta?: Record } | undefined { return ( ctx as unknown as { __finalResult?: { artifactPath: string; meta?: Record; }; } ).__finalResult; } /** Read the in-memory phase state set by the script (runner-only; not part of the public ctx surface). */ export function getWorkflowPhaseState(ctx: WorkflowCtx): { currentPhase: string | undefined; phases: string[] } | undefined { return ( ctx as unknown as { __phaseState?: { currentPhase: string | undefined; phases: string[]; }; } ).__phaseState; } /** Read the in-memory log buffer appended by ctx.log() (runner-only; not part of the public ctx surface). * Capped at 1000 entries — the events log (dwf.log) is the durable source of truth. */ export function getWorkflowLogs(ctx: WorkflowCtx): string[] | undefined { return (ctx as unknown as { __logs?: string[] }).__logs; } /** SDD-3 W-C G13: read the effective run limits (runner-only; not part of the public * ctx surface). Mirrors getWorkflowFinalResult/getWorkflowPhaseState — exposed for * tests and diagnostics so the default cap is observable without dispatching calls. */ export function getWorkflowLimits(ctx: WorkflowCtx): { maxAgentCalls: number } | undefined { return (ctx as unknown as { __limits?: { maxAgentCalls: number } }).__limits; } /** round-18 P2-3: snapshot the current DWF checkpoint state (runner-only; not part of the public * ctx surface). Mirrors getWorkflowFinalResult/getWorkflowPhaseState. The runner relies on the * `onCheckpoint` callback for accurate per-agent-call checkpoints (it captures the closure value * at call time); this helper is a best-effort snapshot for inspection/debugging. */ export function getWorkflowCheckpoint(ctx: WorkflowCtx): DwfCheckpointState { const phaseState = getWorkflowPhaseState(ctx); const logs = getWorkflowLogs(ctx); return { runId: ctx.runId, vars: ctx.vars, phases: phaseState?.phases ?? [], currentPhase: phaseState?.currentPhase, logs: logs ?? [], spent: ctx.budget.spent(), agentCount: (ctx as unknown as { __agentCount?: number }).__agentCount ?? 0, completedAgentCalls: (ctx as unknown as { __completedAgentCalls?: Record }) .__completedAgentCalls ?? {}, updatedAt: new Date().toISOString(), }; } /** Compose the agent task: prompt + optional dependency-input context block. */ function composeAgentTask(call: AgentCallOpts): string { let base = call.prompt; if (call.inputs?.length) { const block = call.inputs.map((p) => `- ${p}`).join("\n"); base = `${base}\n\n## Inputs (artifact paths)\n${block}`; } // Round-13 P0-3: when a schema is requested, append a JSON-output directive. // The directive lives at the END of the prompt so it wins over any conflicting // persona instruction in the agent's system prompt. if (call.schema !== undefined) { base = `${base}\n\n## Output format\nRespond with ONLY a single JSON object that matches the schema described in your instructions. Begin your response with { and end with }. Do not wrap the JSON in a code fence. Do not add any prose before or after the JSON.`; } return base; } /** * Round-13 P0-3: compose a system-prompt suffix that asks the agent to output a * structured JSON object matching the schema's required shape. We don't expose * the TypeBox internal type — we describe the SHAPE so the model can match it. */ function composeSchemaSystemPrompt(base: string | undefined, schema: TSchema): string { const shape = describeSchemaShape(schema, 0); const intro = "You are a structured-output assistant. "; const instruction = `When responding, output ONLY a single JSON object matching this shape (no prose, no markdown fences, no commentary): ${shape}. Begin your response with { and end with }.`; if (typeof base === "string" && base.length > 0) { return `${base}\n\n${intro}${instruction}`; } return `${intro}${instruction}`; } /** * Walk a TypeBox schema recursively and produce a human-readable shape description. * Depth-limited to avoid runaway expansion on deeply nested schemas. */ function describeSchemaShape(schema: unknown, depth: number): string { if (depth > 4) return "{...}"; if (!schema || typeof schema !== "object") return "any"; const obj = schema as Record; // TypeBox: every schema has a `type` discriminator or a `kind` field. const kind = obj.kind as string | undefined; const type = obj.type as string | undefined; if (kind === "object" || type === "object") { const properties = obj.properties; if (!properties || typeof properties !== "object") return "{}"; const required = Array.isArray(obj.required) ? new Set(obj.required as string[]) : new Set(); const props = Object.entries(properties as Record) .map(([key, sub]) => { const mark = required.has(key) ? "" : "?"; return `"${key}"${mark}: ${describeSchemaShape(sub, depth + 1)}`; }) .join(", "); return `{${props}}`; } if (kind === "array" || type === "array") { const items = obj.items; return `[${describeSchemaShape(items, depth + 1)}]`; } if (type === "string") return "string"; if (type === "number" || type === "integer") return "number"; if (type === "boolean") return "boolean"; if (type === "null") return "null"; // Union/Enum fallbacks. if (Array.isArray(obj.anyOf)) return obj.anyOf.map((s) => describeSchemaShape(s, depth + 1)).join(" | "); if (Array.isArray(obj.oneOf)) return obj.oneOf.map((s) => describeSchemaShape(s, depth + 1)).join(" | "); if (Array.isArray(obj.enum)) return obj.enum.map((v) => JSON.stringify(v)).join(" | "); return "any"; } /** * Classify a review outcome from prose when no JSON was produced (round-11 tier-3/4 fallback). * Scans the reviewer's prose for sentiment signals to decide accept / reject / changes_requested. * This keeps the outcome ACCURATE for models that ignore JSON-output instructions. * * Decision order: reject (critical issues) → accept (explicit approval) → changes_requested (default). * reject is checked first because a review can mention both "correctly" (describing existing code) * AND "critical bug" (the verdict) — the verdict signal must win. */ // Precompiled review-outcome signals (module scope — compiled once, not per call). // RR-021 WI-1.1: these were previously string literals like "\baccept\b" compiled // per call — and "\b" in a TS string literal is a BACKSPACE character (U+0008), // so those two signals never matched anything. Regex literals here give real word // boundaries. No /g flags → no shared lastIndex state across .test() calls. const REJECT_SIGNAL_REGEXES: readonly RegExp[] = [ /\breject\b/, /fundamentally/, /completely broken/, /totally broken/, /critical bug/, /critical issue/, /critical flaw/, /security vulnerability/, /does not work/, /doesn't work/, /will not work/, /fails to/, /unacceptable/, /must not be merged/, /do not merge/, /wrong approach/, /logically incorrect/, /incorrectly implements/, /returns the opposite/, /subtraction instead of addition/, /opposite of its intended/, ]; const ACCEPT_SIGNAL_REGEXES: readonly RegExp[] = [ /\baccept\b/, /looks good/, /well done/, /no issues/, /no real issues/, /no problems/, /no concerns/, /nothing to change/, /ready to merge/, /lgtm/, /ship it/, /correctly implements/, /correctly returns/, /works as expected/, /works correctly/, /no bugs/, /no defects/, /meets all requirements/, /all requirements met/, /passes all/, /is correct/, /are correct/, /no changes needed/, /no changes required/, /no further changes/, /nothing more to/, /complete and correct/, /sound implementation/, ]; // RR-021 WI-1.1 negation guard: phrases that NEGATE the bare "accept" signal. // "I cannot accept this" is a complaint, not approval — the accept hit must be // discarded so the outcome falls through to changes_requested (or reject if a // reject signal is also present). const NEGATED_ACCEPT_REGEXES: readonly RegExp[] = [ /\bcannot\s+accept\b/, /\bcan'?t\s+accept\b/, /\bdon'?t\s+accept\b/, /\bwon'?t\s+accept\b/, /\bdo\s+not\s+accept\b/, /\bwill\s+not\s+accept\b/, /\bunable\s+to\s+accept\b/, /\bnot\s+accept(?:ed|ing|able)?\b/, ]; export function classifyReviewOutcome(prose: string): "accept" | "reject" | "changes_requested" { const text = prose.toLowerCase(); // Reject-first precedence: a review can mention both "correctly" (describing // existing code) AND "critical bug" (the verdict) — the verdict signal must win. const hasReject = REJECT_SIGNAL_REGEXES.some((re) => re.test(text)); if (hasReject) return "reject"; const hasAccept = ACCEPT_SIGNAL_REGEXES.some((re) => re.test(text)); if (hasAccept && !NEGATED_ACCEPT_REGEXES.some((re) => re.test(text))) return "accept"; return "changes_requested"; } /** * Round-11 test fix: permissive text extraction for ctx.agent(). * parsePiJsonOutput only handles the canonical pi event stream. When the child emits * a different shape, finalText is empty. This fallback walks the JSON tree looking * for any text-shaped string at any depth, then returns the longest one (typically * the final assistant response). */ export function extractTextFallback(stdout: string): string { const trimmed = stdout.trim(); if (!trimmed) return ""; const candidates: string[] = []; const collect = (value: unknown): void => { if (typeof value === "string") { const t = value.trim(); // Skip very short strings and JSON-ish strings if (t.length >= 2 && !t.startsWith("{") && !t.startsWith("[") && !/^[\d.]+$/.test(t)) { candidates.push(t); } } else if (Array.isArray(value)) { for (const item of value) collect(item); } else if (value && typeof value === "object") { for (const v of Object.values(value as Record)) collect(v); } }; // 1. Try parsing each line as JSON, walk tree for (const line of trimmed.split("\n")) { const lineTrim = line.trim(); if (!lineTrim.startsWith("{")) continue; try { const obj = JSON.parse(lineTrim); collect(obj); } catch { /* skip */ } } // 2. If nothing from JSON, try plain text (longest non-empty line that's not JSON) if (candidates.length === 0) { for (const line of trimmed.split("\n")) { const l = line.trim(); if (l.length >= 3 && !l.startsWith("{") && !l.startsWith("[") && !l.startsWith("=")) candidates.push(l); } } // 3. Return the longest candidate (typically the final answer) if (candidates.length === 0) return ""; candidates.sort((a, b) => b.length - a.length); return candidates[0]; }