diff --git a/packages/core/src/__tests__/step-parsers.test.ts b/packages/core/src/__tests__/step-parsers.test.ts new file mode 100644 index 0000000000..44d0f1d738 --- /dev/null +++ b/packages/core/src/__tests__/step-parsers.test.ts @@ -0,0 +1,317 @@ +import { describe, it, expect, afterEach, beforeEach } from "vitest"; +import { writeFile } from "node:fs/promises"; +import { join } from "node:path"; + +import { createTaskStoreTestHarness } from "./store-test-helpers.js"; +import { + StepParserRegistry, + StepParserRegistrationError, + getStepParser, + listStepParsers, + registerStepParser, + unregisterStepParser, + parseStepHeadings, + parseJsonSteps, + __resetStepParserRegistryForTests, + type StepParser, +} from "../step-parsers.js"; + +describe("step-parsers registry (U12, KTD-12)", () => { + afterEach(() => { + __resetStepParserRegistryForTests(); + }); + + describe("step-headings built-in (byte-identical to legacy)", () => { + const headings = () => getStepParser("step-headings")!; + + it("is registered as a built-in", () => { + expect(getStepParser("step-headings")).toBeDefined(); + expect(listStepParsers().map((p) => p.id)).toContain("step-headings"); + }); + + it("parses unannotated headings byte-identically to the legacy regex", () => { + const content = `## Steps + +### Step 0: Preflight + +- [ ] x + +### Step 1: Implementation + +### Step 2: Testing +`; + expect(headings().parse(content).steps).toEqual([ + { name: "Preflight" }, + { name: "Implementation" }, + { name: "Testing" }, + ]); + }); + + it("matches the legacy regex output exactly for varied unannotated headings", () => { + const content = [ + "### Step 0: A", + "### Step 12: Multi word title", + "### Step 3 — dash but no annotation: Real Name", + "### Step 4: trailing spaces here ", + "### Step 5 no colon at all", + "not a step heading: ignored", + ].join("\n"); + const legacy: { name: string }[] = []; + const re = /^###\s+Step\s+\d+[^:]*:\s*(.+)$/gm; + let m: RegExpExecArray | null; + while ((m = re.exec(content)) !== null) { + legacy.push({ name: m[1].trim() }); + } + expect(headings().parse(content).steps).toEqual(legacy); + }); + + it("parses (depends: 1,2) into 0-indexed dependsOn", () => { + expect(headings().parse("### Step 3 (depends: 1,2): Title").steps).toEqual([ + { name: "Title", dependsOn: [0, 1] }, + ]); + }); + + it("dedupes and sorts depends values", () => { + expect(headings().parse("### Step 5 (depends: 3,1,3,2): T").steps).toEqual([ + { name: "T", dependsOn: [0, 1, 2] }, + ]); + }); + + it("empty depends list yields no dependsOn", () => { + expect(headings().parse("### Step 2 (depends: ): T").steps).toEqual([ + { name: "T" }, + ]); + }); + + it("falls back deterministically on a malformed depends annotation", () => { + expect(headings().parse("### Step 1 (depends: bad): Real Title").steps).toEqual([ + { name: "Real Title" }, + ]); + }); + + it("falls back deterministically when the annotation has no closing paren", () => { + expect(headings().parse("### Step 1 (depends: 1,2 oops: Title").steps).toEqual([ + { name: "1,2 oops: Title" }, + ]); + }); + + it("the extracted parseStepHeadings still yields TaskStep[] with status", () => { + // The store-facing function keeps the `status: "pending"` field. + expect(parseStepHeadings("### Step 0: Preflight")).toEqual([ + { name: "Preflight", status: "pending" }, + ]); + }); + }); + + describe("json-steps built-in", () => { + const json = () => getStepParser("json-steps")!; + + it("is registered as a built-in", () => { + expect(getStepParser("json-steps")).toBeDefined(); + }); + + it("parses a happy-path array of {name, depends}", () => { + const content = JSON.stringify([ + { name: "Plan" }, + { name: "Implement", depends: [1] }, + { name: "Test", depends: [1, 2] }, + ]); + expect(json().parse(content).steps).toEqual([ + { name: "Plan" }, + { name: "Implement", dependsOn: [0] }, + { name: "Test", dependsOn: [0, 1] }, + ]); + }); + + it("converts 1-indexed depends to 0-indexed dependsOn, deduped and sorted", () => { + const content = JSON.stringify([{ name: "X", depends: [3, 1, 3, 2] }]); + expect(json().parse(content).steps).toEqual([ + { name: "X", dependsOn: [0, 1, 2] }, + ]); + }); + + it("trims names and omits dependsOn when depends is empty", () => { + const content = JSON.stringify([{ name: " Spaced ", depends: [] }]); + expect(json().parse(content).steps).toEqual([{ name: "Spaced" }]); + }); + + it("parseJsonSteps is exported directly and matches the registry parser", () => { + const content = JSON.stringify([{ name: "A" }]); + expect(parseJsonSteps(content)).toEqual(json().parse(content)); + }); + + it("throws a descriptive error on non-JSON input", () => { + expect(() => json().parse("not json {")).toThrow(/not valid JSON/); + }); + + it("throws when the document is not an array", () => { + expect(() => json().parse(JSON.stringify({ name: "X" }))).toThrow( + /must be a JSON array/, + ); + }); + + it("throws when a step is missing its name", () => { + expect(() => json().parse(JSON.stringify([{ foo: "bar" }]))).toThrow( + /index 0 must have a non-empty string 'name'/, + ); + }); + + it("throws when a step name is blank", () => { + expect(() => json().parse(JSON.stringify([{ name: " " }]))).toThrow( + /non-empty string 'name'/, + ); + }); + + it("throws when depends is not an array", () => { + expect(() => + json().parse(JSON.stringify([{ name: "X", depends: 1 }])), + ).toThrow(/'depends' must be an array/); + }); + + it("throws when depends contains a non-positive-integer", () => { + expect(() => + json().parse(JSON.stringify([{ name: "X", depends: [0] }])), + ).toThrow(/positive integers/); + expect(() => + json().parse(JSON.stringify([{ name: "X", depends: ["1"] }])), + ).toThrow(/positive integers/); + }); + + it("throws when an entry is not an object", () => { + expect(() => json().parse(JSON.stringify(["just a string"]))).toThrow( + /index 0 must be an object/, + ); + }); + }); + + describe("registry semantics", () => { + it("rejects overwriting a built-in with a non-builtin id", () => { + const reg = new StepParserRegistry(); + reg.register({ id: "step-headings", parse: () => ({ steps: [] }) }, { builtin: true }); + expect(() => + reg.register({ id: "step-headings", parse: () => ({ steps: [] }) }), + ).toThrowError(StepParserRegistrationError); + try { + reg.register({ id: "step-headings", parse: () => ({ steps: [] }) }); + } catch (e) { + expect((e as StepParserRegistrationError).reason).toBe( + "builtin-namespace-protected", + ); + } + }); + + it("rejects a duplicate registration", () => { + const reg = new StepParserRegistry(); + const parser: StepParser = { + id: "plugin:acme:custom", + parse: () => ({ steps: [] }), + }; + reg.register(parser); + expect(() => reg.register(parser)).toThrowError(StepParserRegistrationError); + }); + + it("enforces the plugin id shape for non-builtins", () => { + const reg = new StepParserRegistry(); + const bad = ["custom", "plugin:acme", "plugin::custom", "plugin:Acme:Custom", "other:acme:custom"]; + for (const id of bad) { + expect(() => reg.register({ id, parse: () => ({ steps: [] }) })).toThrowError( + StepParserRegistrationError, + ); + } + // A well-formed namespaced id is accepted. + expect(() => + reg.register({ id: "plugin:acme:custom", parse: () => ({ steps: [] }) }), + ).not.toThrow(); + }); + + it("allows a built-in to use a non-namespaced id", () => { + const reg = new StepParserRegistry(); + expect(() => + reg.register({ id: "step-headings", parse: () => ({ steps: [] }) }, { builtin: true }), + ).not.toThrow(); + }); + + it("rejects an invalid definition (no id / no parse)", () => { + const reg = new StepParserRegistry(); + expect(() => reg.register({ id: "", parse: () => ({ steps: [] }) })).toThrowError( + StepParserRegistrationError, + ); + expect(() => + reg.register({ id: "plugin:acme:x" } as unknown as StepParser), + ).toThrowError(StepParserRegistrationError); + }); + + it("round-trips register/unregister for a plugin parser via the shared API", () => { + const id = "plugin:acme:json2"; + expect(getStepParser(id)).toBeUndefined(); + registerStepParser({ id, parse: () => ({ steps: [{ name: "ok" }] }) }); + expect(getStepParser(id)?.parse("").steps).toEqual([{ name: "ok" }]); + expect(unregisterStepParser(id)).toBe(true); + expect(getStepParser(id)).toBeUndefined(); + // Unregistering again (or a missing id) is a no-op false. + expect(unregisterStepParser(id)).toBe(false); + }); + + it("never unregisters a built-in", () => { + const reg = new StepParserRegistry(); + reg.register({ id: "step-headings", parse: () => ({ steps: [] }) }, { builtin: true }); + expect(reg.unregister("step-headings")).toBe(false); + expect(reg.has("step-headings")).toBe(true); + }); + + it("getStepParser returns undefined for an unknown id", () => { + expect(getStepParser("nope")).toBeUndefined(); + expect(getStepParser("plugin:acme:absent")).toBeUndefined(); + }); + }); + + describe("parseStepsFromPrompt-through-registry parity (KTD-12)", () => { + const harness = createTaskStoreTestHarness(); + + beforeEach(async () => { + await harness.beforeEach(); + }); + afterEach(async () => { + await harness.afterEach(); + }); + + const FIXTURES = [ + `## Steps + +### Step 0: Preflight + +### Step 1: Implementation + +### Step 2: Testing +`, + `# Task + +## Steps + +### Step 1: First + +### Step 2 (depends: 1): Second + +### Step 3 (depends: 1,2): Third +`, + `### Step 1 (depends: bad): Real Title`, + ]; + + it("store path equals the direct step-headings parser on the same content", async () => { + const store = harness.store(); + const rootDir = harness.rootDir(); + for (const content of FIXTURES) { + const task = await store.createTask({ description: "parity" }); + const dir = join(rootDir, ".fusion", "tasks", task.id); + await writeFile(join(dir, "PROMPT.md"), content); + + const viaStore = await store.parseStepsFromPrompt(task.id); + // Direct parser yields { name, dependsOn? }; the store path re-applies + // the `pending` status. Reconstruct the expected store shape from the + // direct parse to assert identical behavior through both paths. + const direct = parseStepHeadings(content); + expect(viaStore).toEqual(direct); + } + }); + }); +}); diff --git a/packages/core/src/__tests__/store-update-step-order.test.ts b/packages/core/src/__tests__/store-update-step-order.test.ts index adde221985..42117abdd8 100644 --- a/packages/core/src/__tests__/store-update-step-order.test.ts +++ b/packages/core/src/__tests__/store-update-step-order.test.ts @@ -54,4 +54,81 @@ describe("TaskStore.updateStep step-order guard", () => { expect(updated.steps[0].status).toBe("done"); expect(updated.log.some((entry) => entry.action.includes("Ignored done→in-progress regression"))).toBe(true); }); + + // ── U6: graph-source projection discipline (KTD-7/KTD-11) ────────────────── + + it("graph source: done is legal in dependency order even when an earlier step is pending", async () => { + // Step 2 depends only on the previous step (1) by default. With step 1 done, + // step 2 may go done under graph source even though step 0 is still pending — + // the legacy strict-index-order guard relaxes to dependency order. + const store = harness.store(); + const task = await harness.createTaskWithSteps(); + // Prime the step list, then give step 2 an explicit dependency on step 0 only + // (skipping step 1), so step 2 may go done with step 1 still pending. + await store.updateStep(task.id, 0, "in-progress"); + const primed = await store.getTask(task.id); + const steps = primed.steps.map((s, i) => (i === 2 ? { ...s, dependsOn: [0] } : { ...s })); + await store.updateTask(task.id, { steps }); + + await store.updateStep(task.id, 0, "done", { source: "graph" }); + const updated = await store.updateStep(task.id, 2, "done", { source: "graph" }); + + expect(updated.steps[2].status).toBe("done"); + // Step 1 was never touched and remains pending — strict index order would have + // suppressed the step-2 done write. + expect(updated.steps[1].status).toBe("pending"); + }); + + it("graph source: out-of-order done (unmet dependency) is suppressed AND audited loudly", async () => { + // Step 1's default dependency is step 0, which is still pending → suppressed. + const store = harness.store(); + const task = await harness.createTaskWithSteps(); + // Prime the step list (graph source bypasses PROMPT.md auto-init). + await store.updateStep(task.id, 1, "in-progress"); + + const updated = await store.updateStep(task.id, 1, "done", { source: "graph" }); + + // Suppressed: step 1's default dependency (step 0) is still pending, so the + // done write is rejected and step 1 keeps its prior (non-done) status. + expect(updated.steps[1].status).not.toBe("done"); + expect( + updated.log.some((e) => e.action.includes("Ignored dependency-order done for step 1")), + ).toBe(true); + // Graph suppression is surfaced loudly (not the legacy silent ignore). + expect( + updated.log.some((e) => e.action.includes("[integrity-warning] graph-source updateStep suppressed")), + ).toBe(true); + }); + + it("legacy source: silent out-of-order ignore behavior is unchanged (no integrity-warning)", async () => { + const store = harness.store(); + const task = await harness.createTaskWithSteps(); + + await store.updateStep(task.id, 0, "done"); + const updated = await store.updateStep(task.id, 2, "done"); // legacy, no source + + expect(updated.steps[2].status).toBe("pending"); + expect(updated.log.some((e) => e.action.includes("Ignored out-of-order done for step 2"))).toBe(true); + // Legacy stays silent — no integrity-warning emitted. + expect(updated.log.some((e) => e.action.includes("[integrity-warning]"))).toBe(false); + }); + + it("graph source: auto-reinit from PROMPT.md is bypassed (explicit indices only)", async () => { + // A fresh task with no JSON steps would, under legacy semantics, parse steps + // from PROMPT.md on the first updateStep. Graph source bypasses that — so an + // index into an unparsed (empty) step list is out of range and rejects. + const store = harness.store(); + const task = await store.createTask({ description: "graph reinit bypass" }); + // No PROMPT.md steps are written; task.steps starts empty. + + await expect(store.updateStep(task.id, 0, "in-progress", { source: "graph" })).rejects.toThrow( + /out of range/, + ); + + // Legacy path on the same empty task would attempt the PROMPT.md reinit + // instead of bypassing — proving the divergence is graph-source-only. (Here + // there is no PROMPT.md either, so legacy also has zero steps and rejects, + // but via the auto-init path rather than the bypass.) + await expect(store.updateStep(task.id, 0, "in-progress")).rejects.toThrow(/out of range/); + }); }); diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index da2bad2e74..574b113075 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -104,6 +104,26 @@ export { registerBuiltinTraits, } from "./builtin-traits.js"; export type { BuiltinTraitId } from "./builtin-traits.js"; +// Step-inversion U12 (KTD-12): step-parser registry + built-ins. +export { + StepParserRegistry, + StepParserRegistrationError, + getStepParserRegistry, + registerStepParser, + getStepParser, + listStepParsers, + unregisterStepParser, + registerBuiltinStepParsers, + parseStepHeadings, + parseJsonSteps, + __resetStepParserRegistryForTests, +} from "./step-parsers.js"; +export type { + StepParser, + StepParseResult, + ParsedStep, + StepParserRegistrationReason, +} from "./step-parsers.js"; export { registerDefaultWorkflowHooks, __resetDefaultWorkflowHooksForTests, diff --git a/packages/core/src/step-parsers.ts b/packages/core/src/step-parsers.ts new file mode 100644 index 0000000000..4962c3b63d --- /dev/null +++ b/packages/core/src/step-parsers.ts @@ -0,0 +1,372 @@ +/** + * Step-parser registry (U12, KTD-12). + * + * Step parsing becomes a graph-native node (`parse-steps`): a registry resolves + * a parser id to an implementation that reads an artifact's content and yields a + * canonical step list. Built-ins: + * - `step-headings` — the extracted `parseStepsFromPrompt` logic (the + * `### Step N:` regex + `(depends: …)` annotation from U1); legacy callers + * in `store.ts` delegate to this exact function (byte-identical parity). + * - `json-steps` — a structured `[{ name, depends? }]` JSON document for + * workflows that plan in JSON. + * + * The registry mirrors the trait-registry posture: built-ins are protected from + * override, and plugins register under namespaced ids + * (`plugin::`). This module is engine-free and must NOT + * import `store.ts` (store imports the extracted parser from here). + * + * Parsers may throw on malformed input; callers (the engine's parse-steps + * handler) map a throw to a routable `outcome:parse-error`. + */ + +import type { TaskStep } from "./types.js"; + +// ── Parser contract ────────────────────────────────────────────────────────── + +/** A parsed step as produced by a parser. `dependsOn` is 0-indexed (same + * convention as the headings `(depends: …)` annotation). */ +export interface ParsedStep { + name: string; + dependsOn?: number[]; +} + +/** The result of running a step parser over an artifact's content. */ +export interface StepParseResult { + steps: ParsedStep[]; +} + +/** A step parser. `parse` may throw on malformed input; the caller maps a throw + * to a routable parse-error outcome. */ +export interface StepParser { + id: string; + parse(content: string): StepParseResult; +} + +// ── Registration error ────────────────────────────────────────────────────── + +/** Named reason codes for a rejected step-parser registration. */ +export type StepParserRegistrationReason = + | "duplicate-id" + | "builtin-namespace-protected" + | "invalid-id" + | "invalid-definition"; + +export class StepParserRegistrationError extends Error { + readonly reason: StepParserRegistrationReason; + readonly parserId: string; + constructor(reason: StepParserRegistrationReason, parserId: string, message: string) { + super(message); + this.name = "StepParserRegistrationError"; + this.reason = reason; + this.parserId = parserId; + } +} + +// ── The registry ──────────────────────────────────────────────────────────── + +interface RegisteredParser { + parser: StepParser; + builtin: boolean; +} + +/** Validate a plugin-namespaced parser id: `plugin::` with + * each segment a non-empty `[a-z0-9-]+` token. */ +function isValidPluginParserId(id: string): boolean { + const parts = id.split(":"); + if (parts.length !== 3) return false; + if (parts[0] !== "plugin") return false; + const seg = /^[a-z0-9-]+$/; + return seg.test(parts[1]) && seg.test(parts[2]); +} + +export class StepParserRegistry { + private readonly parsers = new Map(); + + /** Register a parser. Built-in ids cannot be overridden by non-builtins; a + * non-builtin must use a `plugin::` id. */ + register(parser: StepParser, opts?: { builtin?: boolean }): void { + const builtin = opts?.builtin ?? false; + if (!parser || typeof parser.id !== "string" || parser.id === "") { + throw new StepParserRegistrationError( + "invalid-definition", + String(parser?.id), + "Step parser must have a non-empty string id", + ); + } + if (typeof parser.parse !== "function") { + throw new StepParserRegistrationError( + "invalid-definition", + parser.id, + `Step parser '${parser.id}' must have a parse() function`, + ); + } + + // Existing-id checks first (built-in protection, then duplicate) so a + // non-builtin trying to overwrite a built-in surfaces the protection reason + // rather than the id-shape reason. + const existing = this.parsers.get(parser.id); + if (existing) { + if (!builtin && existing.builtin) { + throw new StepParserRegistrationError( + "builtin-namespace-protected", + parser.id, + `Step parser id '${parser.id}' is a built-in parser and cannot be overridden by a non-builtin registration`, + ); + } + throw new StepParserRegistrationError( + "duplicate-id", + parser.id, + `Step parser id '${parser.id}' is already registered`, + ); + } + + if (!builtin && !isValidPluginParserId(parser.id)) { + throw new StepParserRegistrationError( + "invalid-id", + parser.id, + `Non-builtin step parser '${parser.id}' must use a namespaced id of the form 'plugin::'`, + ); + } + + this.parsers.set(parser.id, { parser, builtin }); + } + + getParser(id: string): StepParser | undefined { + return this.parsers.get(id)?.parser; + } + + has(id: string): boolean { + return this.parsers.has(id); + } + + listParsers(): StepParser[] { + return [...this.parsers.values()].map((r) => r.parser); + } + + /** Remove a parser. Built-ins are never removed (callers should only pass + * plugin-namespaced ids — e.g. for plugin teardown). Returns true if a + * non-builtin parser was present and removed. */ + unregister(id: string): boolean { + const existing = this.parsers.get(id); + if (!existing || existing.builtin) return false; + return this.parsers.delete(id); + } +} + +// ── Built-in: step-headings ─────────────────────────────────────────────────── + +/** + * Parse `### Step N:` headings into the task step list (step-inversion U1). + * + * Backward compatibility is exact: an UNannotated heading parses byte-identically + * to the legacy regex `^###\s+Step\s+\d+[^:]*:\s*(.+)$` (name = text after the + * first colon, trimmed). + * + * The annotation `### Step N (depends: 1,2): Title` is parsed explicitly (the + * legacy regex breaks on the colon inside `depends:`): depends values are + * 1-indexed step numbers in the document and are stored as 0-indexed indices on + * `dependsOn` (deduped, sorted, dropping values <= 0). + * + * Malformed `(depends: …)` annotations fall back deterministically: the heading + * is treated as `### Step N:` with the name starting after the FIRST colon + * following the closing paren (if present), else after the first colon — and no + * `dependsOn` is recorded. + */ +export function parseStepHeadings(content: string): TaskStep[] { + const steps: TaskStep[] = []; + // Legacy matcher — UNCHANGED from the original implementation, so unannotated + // headings (and every legacy edge case, including `[^:]*` spanning newlines) + // parse byte-identically. The full match (`m[0]`) is re-inspected only to layer + // the `(depends: …)` annotation on top. + const stepRegex = /^###\s+Step\s+\d+[^:]*:\s*(.+)$/gm; + // Well-formed annotation form: `### Step N (depends: …): name`. + const annotatedRegex = /^###\s+Step\s+\d+\s*\(depends:\s*([^)]*)\)\s*:\s*([^\n]+)$/; + + let match: RegExpExecArray | null; + while ((match = stepRegex.exec(content)) !== null) { + const full = match[0]; + + // No annotation present → byte-identical legacy behavior. + if (!full.includes("(depends:")) { + steps.push({ name: match[1].trim(), status: "pending" }); + continue; + } + + // 1) Well-formed depends annotation. + const annotated = annotatedRegex.exec(full); + if (annotated) { + const parsed = parseDependsList(annotated[1]); + const name = annotated[2].trim(); + if (parsed !== null) { + if (parsed.length > 0) steps.push({ name, status: "pending", dependsOn: parsed }); + else steps.push({ name, status: "pending" }); + continue; + } + } + + // 2) Annotation present but unparseable (bad values or no closing paren): + // deterministic fallback — name starts after the FIRST colon following the + // closing paren if present, else after the first colon. Operate on the + // first line of the match only (the heading line itself). + const line = full.split("\n")[0]; + const parenIdx = line.indexOf(")"); + const colonAfterParen = parenIdx >= 0 ? line.indexOf(":", parenIdx) : -1; + const colonIdx = colonAfterParen >= 0 ? colonAfterParen : line.indexOf(":"); + if (colonIdx >= 0) { + const fallbackName = line.slice(colonIdx + 1).trim(); + if (fallbackName) steps.push({ name: fallbackName, status: "pending" }); + } + } + return steps; +} + +/** Parse a `depends:` value list (1-indexed step numbers) into 0-indexed, + * deduped, sorted indices. Returns null if any token is not a positive integer. */ +function parseDependsList(raw: string): number[] | null { + const trimmed = raw.trim(); + if (trimmed === "") return []; + const tokens = trimmed.split(",").map((t) => t.trim()); + const out = new Set(); + for (const token of tokens) { + if (!/^\d+$/.test(token)) return null; + const n = Number(token); + if (!Number.isInteger(n) || n < 1) return null; + out.add(n - 1); + } + return [...out].sort((a, b) => a - b); +} + +// ── Built-in: json-steps ────────────────────────────────────────────────────── + +/** + * Parse a JSON document: an array of `{ name: string, depends?: number[] }`. + * `depends` values are 1-indexed step numbers in the document (same convention + * as the headings annotation), converted to 0-indexed `dependsOn` (deduped, + * sorted). Throws a descriptive error on any malformed input (not JSON, not an + * array, missing/blank name, bad depends). + */ +export function parseJsonSteps(content: string): StepParseResult { + let doc: unknown; + try { + doc = JSON.parse(content); + } catch (err) { + throw new Error( + `json-steps: content is not valid JSON: ${(err as Error).message}`, + ); + } + + if (!Array.isArray(doc)) { + throw new Error("json-steps: document must be a JSON array of step objects"); + } + + const steps: ParsedStep[] = []; + doc.forEach((entry, i) => { + if (typeof entry !== "object" || entry === null || Array.isArray(entry)) { + throw new Error(`json-steps: step at index ${i} must be an object`); + } + const obj = entry as Record; + const name = obj.name; + if (typeof name !== "string" || name.trim() === "") { + throw new Error( + `json-steps: step at index ${i} must have a non-empty string 'name'`, + ); + } + + const step: ParsedStep = { name: name.trim() }; + + if (obj.depends !== undefined) { + if (!Array.isArray(obj.depends)) { + throw new Error( + `json-steps: step at index ${i} 'depends' must be an array of positive integers`, + ); + } + const out = new Set(); + for (const raw of obj.depends) { + if (typeof raw !== "number" || !Number.isInteger(raw) || raw < 1) { + throw new Error( + `json-steps: step at index ${i} 'depends' must contain only positive integers (1-indexed step numbers); got ${JSON.stringify(raw)}`, + ); + } + out.add(raw - 1); + } + const dependsOn = [...out].sort((a, b) => a - b); + if (dependsOn.length > 0) step.dependsOn = dependsOn; + } + + steps.push(step); + }); + + return { steps }; +} + +// ── Built-in parser definitions ─────────────────────────────────────────────── + +const BUILTIN_STEP_PARSERS: StepParser[] = [ + { + id: "step-headings", + parse(content: string): StepParseResult { + // The headings parser yields TaskStep[]; map to the parser contract + // (dropping the `status` field, which the caller re-applies). + const steps = parseStepHeadings(content).map((s) => { + const out: ParsedStep = { name: s.name }; + if (s.dependsOn) out.dependsOn = s.dependsOn; + return out; + }); + return { steps }; + }, + }, + { + id: "json-steps", + parse: parseJsonSteps, + }, +]; + +/** Register the built-in step parsers into the given registry (defaults to the + * shared registry). Idempotent via `has`. */ +export function registerBuiltinStepParsers( + registry: StepParserRegistry = getStepParserRegistry(), +): void { + for (const parser of BUILTIN_STEP_PARSERS) { + if (registry.has(parser.id)) continue; + registry.register(parser, { builtin: true }); + } +} + +// ── Module-level default registry ─────────────────────────────────────────── + +let defaultRegistry: StepParserRegistry | undefined; + +export function getStepParserRegistry(): StepParserRegistry { + if (!defaultRegistry) { + defaultRegistry = new StepParserRegistry(); + registerBuiltinStepParsers(defaultRegistry); + } + return defaultRegistry; +} + +/** Test-only: reset the shared registry (so built-in registration can be + * re-exercised in isolation). */ +export function __resetStepParserRegistryForTests(): void { + defaultRegistry = undefined; +} + +// ── Convenience pass-throughs to the default registry ──────────────────────── + +export function registerStepParser(parser: StepParser, opts?: { builtin?: boolean }): void { + getStepParserRegistry().register(parser, opts); +} + +export function getStepParser(id: string): StepParser | undefined { + return getStepParserRegistry().getParser(id); +} + +export function listStepParsers(): StepParser[] { + return getStepParserRegistry().listParsers(); +} + +export function unregisterStepParser(id: string): boolean { + return getStepParserRegistry().unregister(id); +} + +// Register built-ins into the shared registry on import (idempotent via `has`). +registerBuiltinStepParsers(); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 6716e8f0ed..85d706e40d 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -56,6 +56,10 @@ import { // Side-effect import: registers the 14 built-in trait DEFINITIONS into the // shared trait registry on load (the flag-ON path resolves traits by id). import "./builtin-traits.js"; +// Step-inversion U12 (KTD-12): the legacy `parseStepsFromPrompt` path resolves +// the `step-headings` parser through the registry (proving the registry path), +// staying byte-identical with the direct extracted function. +import { getStepParser } from "./step-parsers.js"; import type { WorkflowDefinition, WorkflowDefinitionInput, @@ -808,86 +812,11 @@ const KNOWN_FILE_SCOPE_ROOT_FILES = new Set([ "agents.md", ]); -/** - * Parse `### Step N:` headings into the task step list (step-inversion U1). - * - * Backward compatibility is exact: an UNannotated heading parses byte-identically - * to the legacy regex `^###\s+Step\s+\d+[^:]*:\s*(.+)$` (name = text after the - * first colon, trimmed). - * - * The annotation `### Step N (depends: 1,2): Title` is parsed explicitly (the - * legacy regex breaks on the colon inside `depends:`): depends values are - * 1-indexed step numbers in the document and are stored as 0-indexed indices on - * `dependsOn` (deduped, sorted, dropping values <= 0). - * - * Malformed `(depends: …)` annotations fall back deterministically: the heading - * is treated as `### Step N:` with the name starting after the FIRST colon - * following the closing paren (if present), else after the first colon — and no - * `dependsOn` is recorded. - */ -export function parseStepHeadings(content: string): import("./types.js").TaskStep[] { - const steps: import("./types.js").TaskStep[] = []; - // Legacy matcher — UNCHANGED from the original implementation, so unannotated - // headings (and every legacy edge case, including `[^:]*` spanning newlines) - // parse byte-identically. The full match (`m[0]`) is re-inspected only to layer - // the `(depends: …)` annotation on top. - const stepRegex = /^###\s+Step\s+\d+[^:]*:\s*(.+)$/gm; - // Well-formed annotation form: `### Step N (depends: …): name`. - const annotatedRegex = /^###\s+Step\s+\d+\s*\(depends:\s*([^)]*)\)\s*:\s*([^\n]+)$/; - - let match: RegExpExecArray | null; - while ((match = stepRegex.exec(content)) !== null) { - const full = match[0]; - - // No annotation present → byte-identical legacy behavior. - if (!full.includes("(depends:")) { - steps.push({ name: match[1].trim(), status: "pending" }); - continue; - } - - // 1) Well-formed depends annotation. - const annotated = annotatedRegex.exec(full); - if (annotated) { - const parsed = parseDependsList(annotated[1]); - const name = annotated[2].trim(); - if (parsed !== null) { - if (parsed.length > 0) steps.push({ name, status: "pending", dependsOn: parsed }); - else steps.push({ name, status: "pending" }); - continue; - } - } - - // 2) Annotation present but unparseable (bad values or no closing paren): - // deterministic fallback — name starts after the FIRST colon following the - // closing paren if present, else after the first colon. Operate on the - // first line of the match only (the heading line itself). - const line = full.split("\n")[0]; - const parenIdx = line.indexOf(")"); - const colonAfterParen = parenIdx >= 0 ? line.indexOf(":", parenIdx) : -1; - const colonIdx = colonAfterParen >= 0 ? colonAfterParen : line.indexOf(":"); - if (colonIdx >= 0) { - const fallbackName = line.slice(colonIdx + 1).trim(); - if (fallbackName) steps.push({ name: fallbackName, status: "pending" }); - } - } - return steps; -} - -/** Parse a `depends:` value list (1-indexed step numbers) into 0-indexed, - * deduped, sorted indices. Returns null if any token is not a positive integer. */ -function parseDependsList(raw: string): number[] | null { - const trimmed = raw.trim(); - if (trimmed === "") return []; - const tokens = trimmed.split(",").map((t) => t.trim()); - const out = new Set(); - for (const token of tokens) { - if (!/^\d+$/.test(token)) return null; - const n = Number(token); - if (!Number.isInteger(n) || n < 1) return null; - out.add(n - 1); - } - return [...out].sort((a, b) => a - b); -} +// `parseStepHeadings` (the `### Step N:` parser, step-inversion U1) was extracted +// into `step-parsers.ts` as the `step-headings` built-in parser (U12, KTD-12). +// It is re-exported here for back-compat with callers/tests that import it from +// `store.ts`. `parseStepsFromPrompt` below delegates through the registry. +export { parseStepHeadings } from "./step-parsers.js"; export function isValidFileScopeEntry(token: string): boolean { const trimmed = token.trim(); @@ -7805,13 +7734,27 @@ export class TaskStore extends EventEmitter { id: string, stepIndex: number, status: import("./types.js").StepStatus, + options?: { source?: "graph" }, ): Promise { + // Step-inversion projection discipline (U6/KTD-7). A `source: "graph"` write + // is the workflow-graph executor projecting a foreach instance's lifecycle + // (in-progress / done / pending) onto Task.steps[] with EXPLICIT indices. Three + // behaviors diverge from the legacy (default) write: + // (a) the out-of-order-done guard relaxes from strict index order to + // DEPENDENCY order (a done write is legal when every dependsOn step — + // default: the immediately-preceding step — is done/skipped, KTD-11); + // (b) a guard that DOES suppress a graph write logs an audit warning loudly + // (legacy stays silent — a graph suppression is a projection bug); + // (c) the auto-reinit-from-PROMPT.md path is bypassed (the graph pinned the + // step count at foreach expansion; re-parsing here would desync, KTD-3). + const graphSource = options?.source === "graph"; return this.withTaskLock(id, async () => { const dir = this.taskDir(id); const task = await this.readTaskJson(dir); - // Auto-initialize steps from PROMPT.md if empty - if (task.steps.length === 0) { + // Auto-initialize steps from PROMPT.md if empty. Bypassed for graph-source + // writes (U6/KTD-3): the graph owns explicit indices pinned at expansion. + if (task.steps.length === 0 && !graphSource) { task.steps = await this.parseStepsFromPrompt(id); } @@ -7848,22 +7791,63 @@ export class TaskStore extends EventEmitter { } if (status === "done") { - for (let i = 0; i < stepIndex; i++) { - const priorStatus = task.steps[i].status; - if (priorStatus === "pending" || priorStatus === "in-progress") { - const ts = new Date().toISOString(); - task.updatedAt = ts; + // The set of predecessor steps that must be done/skipped before this step + // may go done. Legacy: strict index order (every earlier step). Graph: the + // step's dependsOn list (default = the immediately-preceding step when the + // annotation is absent — preserving sequential behavior, KTD-11). + let blockingIndex = -1; + let blockingStatus: import("./types.js").StepStatus | undefined; + if (graphSource) { + const deps = task.steps[stepIndex]?.dependsOn; + const depIndices = + Array.isArray(deps) && deps.length > 0 + ? deps + : stepIndex > 0 + ? [stepIndex - 1] + : []; + for (const i of depIndices) { + const priorStatus = task.steps[i]?.status; + if (priorStatus === "pending" || priorStatus === "in-progress") { + blockingIndex = i; + blockingStatus = priorStatus; + break; + } + } + } else { + for (let i = 0; i < stepIndex; i++) { + const priorStatus = task.steps[i].status; + if (priorStatus === "pending" || priorStatus === "in-progress") { + blockingIndex = i; + blockingStatus = priorStatus; + break; + } + } + } + if (blockingIndex !== -1) { + const ts = new Date().toISOString(); + task.updatedAt = ts; + const kind = graphSource ? "dependency-order" : "out-of-order"; + task.log.push({ + timestamp: ts, + action: + `Ignored ${kind} ${status} for step ${stepIndex} (${task.steps[stepIndex].name}) — ` + + `${graphSource ? "dependency" : "earlier"} step ${blockingIndex} (${task.steps[blockingIndex].name}) is still ${blockingStatus}`, + }); + // Graph-source suppression is a projection bug — surface it loudly in + // the activity log (U6) rather than the legacy silent ignore. + if (graphSource) { task.log.push({ timestamp: ts, action: - `Ignored out-of-order ${status} for step ${stepIndex} (${task.steps[stepIndex].name}) — ` + - `earlier step ${i} (${task.steps[i].name}) is still ${priorStatus}`, + `[integrity-warning] graph-source updateStep suppressed: step ${stepIndex} ` + + `(${task.steps[stepIndex].name}) → done blocked by unmet dependency ` + + `step ${blockingIndex} (${blockingStatus})`, }); - await this.atomicWriteTaskJson(dir, task); - if (this.isWatching) this.taskCache.set(id, { ...task }); - this.emit("task:updated", task); - return task; } + await this.atomicWriteTaskJson(dir, task); + if (this.isWatching) this.taskCache.set(id, { ...task }); + this.emit("task:updated", task); + return task; } } @@ -8795,7 +8779,19 @@ export class TaskStore extends EventEmitter { if (!existsSync(promptPath)) return []; const content = await readFile(promptPath, "utf-8"); - return parseStepHeadings(content); + // Step-inversion U12 (KTD-12): delegate to the registry's `step-headings` + // parser (resolved by id, not a direct import) so the registry path is + // proven and stays byte-identical to the extracted function. The parser + // yields `{ name, dependsOn? }`; re-apply the `pending` status here. + const parser = getStepParser("step-headings"); + if (!parser) { + throw new Error("Step parser 'step-headings' is not registered"); + } + return parser.parse(content).steps.map((s) => + s.dependsOn + ? { name: s.name, status: "pending" as const, dependsOn: s.dependsOn } + : { name: s.name, status: "pending" as const }, + ); } /** diff --git a/packages/engine/src/__tests__/workflow-graph-executor-parity.test.ts b/packages/engine/src/__tests__/workflow-graph-executor-parity.test.ts index e16d871c34..631138786b 100644 --- a/packages/engine/src/__tests__/workflow-graph-executor-parity.test.ts +++ b/packages/engine/src/__tests__/workflow-graph-executor-parity.test.ts @@ -40,9 +40,10 @@ describe("WorkflowGraphExecutor interpreter-parity", () => { schedule: async () => ({ outcome: "success" }), }; const legacyEvents = await runLegacy(seams)(); + type BaseSeam = "planning" | "execute" | "review" | "merge" | "schedule"; const executor = new WorkflowGraphExecutor({ seams, handlers: { prompt: async (node, ctx) => { const seam = String(node.config?.seam); - const result = await seams[seam as keyof WorkflowLegacySeams]!(ctx.task, ctx.context); + const result = await seams[seam as BaseSeam](ctx.task, ctx.context); events.push(`${seam}:${result.outcome}`); return result; } } }); diff --git a/packages/engine/src/__tests__/workflow-graph-foreach.test.ts b/packages/engine/src/__tests__/workflow-graph-foreach.test.ts index 86c14448ad..8734d1b0ad 100644 --- a/packages/engine/src/__tests__/workflow-graph-foreach.test.ts +++ b/packages/engine/src/__tests__/workflow-graph-foreach.test.ts @@ -460,6 +460,79 @@ describe("WorkflowGraphExecutor foreach (U3)", () => { expect(saved.some((s) => s.status === "completed")).toBe(true); expect(saved.every((s) => s.pinnedStepCount === 1)).toBe(true); }); + + // ── U6: projection discipline ────────────────────────────────────────────── + + it("projection-first ordering: step projection writes precede the completed instance row", async () => { + // The merge-blocker race (KTD-7) is closed by ordering: the step projection + // (updateStep) must be observable BEFORE the instance row flips to completed. + // We interleave both into one event log: the stepExecute seam stands in for + // the projection write; the persistence hook records the row status. + const events: string[] = []; + const seams = baseSeams({ + stepExecute: async (_t, ctx) => { + const active = ctx[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext; + events.push(`projection:done#${active.stepIndex}`); + return { outcome: "success", value: "step-done" }; + }, + }); + const executor = new WorkflowGraphExecutor({ + seams, + stepInstancePersistence: { + saveInstanceState: (s) => { + events.push(`row:${s.status}#${s.stepIndex}`); + }, + }, + }); + const result = await executor.run(taskWithSteps(1), settingsOn(), foreachIr(singleExecuteTemplate())); + + expect(result.outcome).toBe("success"); + const projectionIdx = events.indexOf("projection:done#0"); + const completedIdx = events.indexOf("row:completed#0"); + expect(projectionIdx).toBeGreaterThanOrEqual(0); + expect(completedIdx).toBeGreaterThanOrEqual(0); + // Projection (done) is observable before the instance row flips to completed. + expect(projectionIdx).toBeLessThan(completedIdx); + }); + + it("sets deferDoneToReview on the active instance when the template has a step-review node", async () => { + // U6/KTD-4: with a step-review node present, step-execute must NOT mark the + // step done (markDoneOnSuccess:false) — the active context flags this so the + // step-execute seam can pass the flag to runTaskStep. + let observedDefer: boolean | undefined; + let observedNoReviewDefer: boolean | undefined; + const seams = baseSeams({ + stepExecute: async (_t, ctx) => { + const active = ctx[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext; + observedDefer = active.deferDoneToReview; + return { outcome: "success", value: "step-done" }; + }, + stepReview: async () => ({ verdict: "APPROVE" as const }), + }); + const reviewTemplate = { + nodes: [ + { id: "exec", kind: "prompt" as const, config: { seam: "step-execute" } }, + { id: "review", kind: "step-review" as const, config: { type: "code" } }, + ], + edges: [{ from: "exec", to: "review", condition: "success" }], + }; + const executor = new WorkflowGraphExecutor({ seams }); + await executor.run(taskWithSteps(1), settingsOn(), foreachIr(reviewTemplate)); + expect(observedDefer).toBe(true); + + // Without a step-review node, deferDoneToReview is false (step-execute is the + // done authority). + const seamsNoReview = baseSeams({ + stepExecute: async (_t, ctx) => { + const active = ctx[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext; + observedNoReviewDefer = active.deferDoneToReview; + return { outcome: "success", value: "step-done" }; + }, + }); + const executor2 = new WorkflowGraphExecutor({ seams: seamsNoReview }); + await executor2.run(taskWithSteps(1), settingsOn(), foreachIr(singleExecuteTemplate())); + expect(observedNoReviewDefer).toBe(false); + }); }); // ── helpers ─────────────────────────────────────────────────────────────── diff --git a/packages/engine/src/__tests__/workflow-step-review.test.ts b/packages/engine/src/__tests__/workflow-step-review.test.ts new file mode 100644 index 0000000000..f5750b41da --- /dev/null +++ b/packages/engine/src/__tests__/workflow-step-review.test.ts @@ -0,0 +1,263 @@ +import { describe, expect, it, vi } from "vitest"; +import type { TaskDetail, TaskStep, WorkflowIr, WorkflowIrNode } from "@fusion/core"; + +import { WorkflowGraphExecutor } from "../workflow-graph-executor.js"; +import { + FOREACH_ACTIVE_CONTEXT_KEY, + SPLIT_ACTIVE_CONTEXT_KEY, + type ForeachActiveContext, + type StepReviewSeamResult, + type WorkflowLegacySeams, +} from "../workflow-node-handlers.js"; +import type { WorkflowStepInstanceState } from "../workflow-graph-foreach.js"; + +/** + * U5 — step-review node + verdict wiring (KTD-4). These scenarios exercise the + * real {@link createStepReviewHandler} (registered by default in the executor) + * driving a `seams.stepReview` fake, with the foreach sub-walk providing the + * `foreach:active` context, rework edges, and the RETHINK reset hook. + */ + +const settingsOn = () => ({ experimentalFeatures: { workflowGraphExecutor: true } }); + +function taskWithSteps(n: number): TaskDetail { + const steps: TaskStep[] = Array.from({ length: n }, (_, i) => ({ + name: `Step ${i + 1}`, + status: "pending" as const, + })); + return { id: "FN-REVIEW", steps } as unknown as TaskDetail; +} + +/** Base no-op seams with overrides. */ +function baseSeams(overrides: Partial): WorkflowLegacySeams { + const ok = async () => ({ outcome: "success" as const }); + return { planning: ok, execute: ok, review: ok, merge: ok, schedule: ok, ...overrides }; +} + +/** + * Build: start → foreach{ exec(step-execute) → review(step-review) } → end. + * Verdict edges from review: approve → exit (no edge = template exit), revise → + * rework to exec, rethink → rework to exec. Foreach exhaustion routes to a hold. + */ +function reviewForeachIr(opts: { config?: Record } = {}): WorkflowIr { + const template = { + nodes: [ + { id: "exec", kind: "prompt" as const, config: { seam: "step-execute" } }, + { id: "review", kind: "step-review" as const, config: { type: "code" } }, + ] as WorkflowIrNode[], + edges: [ + { from: "exec", to: "review", condition: "success" }, + // approve (and unavailable) have NO outgoing edge from review → template exit + // (instance done / advisory continuation). + { from: "review", to: "exec", condition: "outcome:revise", kind: "rework" as const }, + { from: "review", to: "exec", condition: "outcome:rethink", kind: "rework" as const }, + ], + }; + return { + version: "v2", + name: "review-test", + columns: [{ id: "work", name: "Work", traits: [] }], + nodes: [ + { id: "start", kind: "start" }, + { id: "fe", kind: "foreach", config: { source: "task-steps", template, ...(opts.config ?? {}) } }, + { id: "hold", kind: "prompt", config: {} }, + { id: "end", kind: "end" }, + ], + edges: [ + { from: "start", to: "fe" }, + { from: "fe", to: "end", condition: "success" }, + { from: "fe", to: "hold", condition: "outcome:rework-exhausted" }, + ], + }; +} + +describe("WorkflowGraphExecutor step-review (U5)", () => { + it("APPROVE marks the step done via the projection and routes the approve edge", async () => { + const doneMarks: Array<{ index: number; status: string }> = []; + const stepReview = vi.fn(async (): Promise => ({ verdict: "APPROVE" })); + const seams = baseSeams({ + stepExecute: async (_t, ctx) => { + const active = ctx[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext; + active.baselineSha = `base-${active.stepIndex}`; + // step-execute leaves the step in-progress (review decides done) — record + // that nothing was done here. + return { outcome: "success", value: "step-done", contextPatch: { [FOREACH_ACTIVE_CONTEXT_KEY]: active } }; + }, + stepReview: async (_t, _ctx, cfg) => { + const r = await stepReview(); + // Simulate the executor's APPROVE projection write. + if (r.verdict === "APPROVE" && !cfg.advisory) doneMarks.push({ index: 0, status: "done" }); + return r; + }, + }); + const executor = new WorkflowGraphExecutor({ seams }); + const result = await executor.run(taskWithSteps(1), settingsOn(), reviewForeachIr()); + + expect(result.outcome).toBe("success"); + expect(stepReview).toHaveBeenCalledTimes(1); + expect(doneMarks).toEqual([{ index: 0, status: "done" }]); + }); + + it("REVISE routes a rework edge without triggering a reset", async () => { + const resets: string[] = []; + let reviewCalls = 0; + const seams = baseSeams({ + stepExecute: async () => ({ outcome: "success", value: "step-done" }), + stepReview: async (): Promise => { + reviewCalls += 1; + return reviewCalls === 1 ? { verdict: "REVISE" } : { verdict: "APPROVE" }; + }, + }); + const executor = new WorkflowGraphExecutor({ + seams, + onReworkReset: async (active, reason) => { + resets.push(`${active.stepIndex}:${reason}`); + }, + }); + const result = await executor.run(taskWithSteps(1), settingsOn(), reviewForeachIr()); + + expect(result.outcome).toBe("success"); + expect(reviewCalls).toBe(2); // revise → rework → approve + expect(resets).toEqual([]); // REVISE never resets + }); + + it("RETHINK resets to baseline then re-executes the step", async () => { + const resets: Array<{ index: number; reason: string; baseline?: string }> = []; + let reviewCalls = 0; + const seams = baseSeams({ + stepExecute: async (_t, ctx) => { + const active = ctx[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext; + active.baselineSha = "base-rethink"; + active.checkpointId = "ckpt-1"; + return { outcome: "success", value: "step-done", contextPatch: { [FOREACH_ACTIVE_CONTEXT_KEY]: active } }; + }, + stepReview: async (): Promise => { + reviewCalls += 1; + return reviewCalls === 1 ? { verdict: "RETHINK" } : { verdict: "APPROVE" }; + }, + }); + const executor = new WorkflowGraphExecutor({ + seams, + onReworkReset: async (active, reason) => { + resets.push({ index: active.stepIndex, reason, baseline: active.baselineSha }); + }, + }); + const result = await executor.run(taskWithSteps(1), settingsOn(), reviewForeachIr()); + + expect(result.outcome).toBe("success"); + expect(reviewCalls).toBe(2); + expect(resets).toEqual([{ index: 0, reason: "rethink", baseline: "base-rethink" }]); + }); + + it("UNAVAILABLE retries inside the handler (cap 2) then routes outcome:unavailable", async () => { + let reviewCalls = 0; + const seams = baseSeams({ + stepExecute: async () => ({ outcome: "success", value: "step-done" }), + stepReview: async (): Promise => { + reviewCalls += 1; + return { verdict: "UNAVAILABLE" }; + }, + }); + const executor = new WorkflowGraphExecutor({ seams }); + const result = await executor.run(taskWithSteps(1), settingsOn(), reviewForeachIr()); + + // The handler retries up to the cap (3 invocations: initial + 2 retries). + expect(reviewCalls).toBe(3); + // value routed is "unavailable"; the IR has no unavailable edge from review, + // so the instance exits the template (advisory) and the foreach succeeds. + expect(result.outcome).toBe("success"); + }); + + it("persists the verdict into the instance row", async () => { + const saved: WorkflowStepInstanceState[] = []; + const seams = baseSeams({ + stepExecute: async () => ({ outcome: "success", value: "step-done" }), + stepReview: async (): Promise => ({ verdict: "APPROVE" }), + }); + const executor = new WorkflowGraphExecutor({ + seams, + stepInstancePersistence: { + saveInstanceState: (s) => { + saved.push({ ...s }); + }, + }, + }); + const result = await executor.run(taskWithSteps(1), settingsOn(), reviewForeachIr()); + + expect(result.outcome).toBe("success"); + // The final (completed) instance row carries the authoritative APPROVE verdict. + const completed = saved.filter((s) => s.status === "completed"); + expect(completed.length).toBeGreaterThan(0); + expect(completed[completed.length - 1].verdict).toBe("APPROVE"); + }); + + it("split-branch review is advisory-only: no authoritative verdict, no projection write", async () => { + // Simulate the split-active marker the executor sets around branches: the + // handler reads SPLIT_ACTIVE_CONTEXT_KEY from the shared context and flags the + // review advisory. We assert the seam was told advisory=true and that an + // advisory APPROVE does not write the projection. + const calls: Array<{ advisory: boolean | undefined }> = []; + const projectionWrites: number[] = []; + const seams = baseSeams({ + stepExecute: async () => ({ outcome: "success", value: "step-done" }), + stepReview: async (_t, _ctx, cfg) => { + calls.push({ advisory: cfg.advisory }); + if (cfg.type === "code" && !cfg.advisory) projectionWrites.push(1); + return { verdict: "APPROVE" }; + }, + }); + const executor = new WorkflowGraphExecutor({ seams }); + + // Build a foreach whose template puts the step-review behind a manual + // split-active marker on the shared context via a custom prelude node. + const template = { + nodes: [ + { id: "exec", kind: "prompt" as const, config: { seam: "step-execute" } }, + { id: "mark", kind: "prompt" as const, config: {} }, + { id: "review", kind: "step-review" as const, config: { type: "code" } }, + { id: "exit", kind: "prompt" as const, config: {} }, + ] as WorkflowIrNode[], + edges: [ + { from: "exec", to: "mark", condition: "success" }, + { from: "mark", to: "review", condition: "success" }, + { from: "review", to: "exit", condition: "outcome:approve" }, + ], + }; + const ir: WorkflowIr = { + version: "v2", + name: "advisory-test", + columns: [{ id: "work", name: "Work", traits: [] }], + nodes: [ + { id: "start", kind: "start" }, + { id: "fe", kind: "foreach", config: { source: "task-steps", template } }, + { id: "end", kind: "end" }, + ], + edges: [ + { from: "start", to: "fe" }, + { from: "fe", to: "end", condition: "success" }, + ], + }; + + // Custom handler for the "mark" node sets split:active on the shared context + // to simulate running inside a split branch window. + const exec = new WorkflowGraphExecutor({ + seams, + handlers: { + prompt: async (node, ctx) => { + if (node.config?.seam === "step-execute") return seams.stepExecute!(ctx.task, ctx.context); + if (node.id === "mark") { + ctx.context[SPLIT_ACTIVE_CONTEXT_KEY] = true; + return { outcome: "success" }; + } + return { outcome: "success" }; + }, + }, + }); + void executor; + const result = await exec.run(taskWithSteps(1), settingsOn(), ir); + + expect(result.outcome).toBe("success"); + expect(calls).toEqual([{ advisory: true }]); + expect(projectionWrites).toEqual([]); // advisory APPROVE never writes projection + }); +}); diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index b53ead6215..dbf4833764 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -18,6 +18,10 @@ import { } from "@fusion/core"; import { WorkflowGraphTaskRunner, type WorkflowGraphTaskRunResult } from "./workflow-graph-task-runner.js"; import type { WorkflowBranchPersistence, WorkflowBranchRunState } from "./workflow-graph-branches.js"; +import type { + WorkflowStepInstancePersistence, + WorkflowStepInstanceState, +} from "./workflow-graph-foreach.js"; import { observeWorkflowParity, WORKFLOW_INTERPRETER_DUAL_OBSERVE_FLAG } from "./workflow-parity-observer.js"; import { FOREACH_ACTIVE_CONTEXT_KEY, @@ -107,7 +111,7 @@ import type { StuckTaskDetector, StuckTaskEvent } from "./stuck-task-detector.js import type { PluginRunner } from "./plugin-runner.js"; import { isContextLimitError } from "./context-limit-detector.js"; import { StepSessionExecutor } from "./step-session-executor.js"; -import { resetStepToBaseline, runTaskStep } from "./step-runner.js"; +import { makeAncestryBlastRadiusGuard, resetStepToBaseline, runTaskStep } from "./step-runner.js"; import { acquireTaskWorktree } from "./worktree-acquisition.js"; import { resolveCapturedBaseCommitSha } from "./base-commit-capture.js"; import { installTaskWorktreeIdentityGuard } from "./worktree-hooks.js"; @@ -3215,6 +3219,20 @@ export class TaskExecutor { * Doubles as the re-entrancy guard for graph routing. */ private graphCompletionInterceptors = new Map void>(); + /** Step-inversion (KTD-2/KTD-8, U6/U8): tasks whose graph-owned step-execute + * driver has pinned step-session physics for the run. Forces the step-session + * path in execute() regardless of the `runStepsInNewSessions` setting, so the + * graph/step-sessions flag matrix cannot select an unsupported physics combo. + * Cleared when the graph run ends (maybeExecuteWorkflowGraph finally). */ + private graphStepSessionPinned = new Set(); + + /** Step-inversion (U6/U8): caches the per-run implementation-phase result for a + * graph-owned task so the foreach sub-walk's per-step `runTaskStep` driver runs + * the (step-session) implementation exactly once per run and lets later step + * instances observe the projection rather than re-running execute() per step. + * Keyed by task id; cleared alongside the pin. */ + private graphStepRunOnce = new Map>(); + /** Tasks currently being orchestrated by the graph runner. Process-wide for * the same reason as executingTaskLock (FN-4811): duplicate execute() * invocations can arrive from different TaskExecutor instances in one @@ -3271,6 +3289,13 @@ export class TaskExecutor { // real data, and prunes stale runs (#1412). Adapter degrades to no-op // when the store predates these methods (additive guard). branchPersistence: this.buildBranchPersistence(), + // Step-inversion (KTD-6, U3/U4): per-instance run-state persistence. + stepInstancePersistence: this.buildStepInstancePersistence(), + // Step-inversion (KTD-4, U5): RETHINK reset-on-rework — when the foreach + // sub-walk traverses a rework edge triggered by `outcome:rethink`, reset + // the active instance's step to its persisted per-step baseline (git reset + // + session rewind + step→pending) before re-entering step-execute. + onReworkReset: (active) => this.applyGraphRethinkReset(task.id, active), }); let result: WorkflowGraphTaskRunResult; try { @@ -3294,6 +3319,9 @@ export class TaskExecutor { return true; } finally { this.graphRouting.delete(task.id); + // Clear per-run step-inversion pins (KTD-8: pinned only for the run's life). + this.graphStepSessionPinned.delete(task.id); + this.graphStepRunOnce.delete(task.id); } } @@ -3318,6 +3346,64 @@ export class TaskExecutor { }; } + /** + * Build the store-backed WorkflowStepInstancePersistence for graph-owned + * foreach runs (KTD-6, U3/U4 seam). Returns undefined when the store predates + * the instance CRUD methods (the SQLite migration is U4) so the sub-walk stays + * fully in-memory — purely additive, same posture as buildBranchPersistence. + */ + private buildStepInstancePersistence(): WorkflowStepInstancePersistence | undefined { + const store = this.store as unknown as { + saveWorkflowRunStepInstance?: (state: WorkflowStepInstanceState) => void; + loadWorkflowRunStepInstances?: (taskId: string, runId: string) => WorkflowStepInstanceState[]; + clearWorkflowRunStepInstances?: (taskId: string, keepRunId: string) => void; + }; + if (typeof store.saveWorkflowRunStepInstance !== "function") return undefined; + return { + saveInstanceState: (state) => store.saveWorkflowRunStepInstance?.(state), + loadInstanceStates: (taskId, runId) => store.loadWorkflowRunStepInstances?.(taskId, runId) ?? [], + clearStaleInstanceStates: (taskId, keepRunId) => store.clearWorkflowRunStepInstances?.(taskId, keepRunId), + }; + } + + /** + * RETHINK reset-on-rework (KTD-4, U5): reset the active foreach instance's step + * to its per-step baseline before the rework edge re-enters step-execute. Drives + * the single extracted `resetStepToBaseline` (step-runner.ts) with the + * instance's persisted `baselineSha`/`checkpointId`. Session rewind is best-effort + * for graph-owned runs (the per-step session lives inside StepSessionExecutor and + * is not exposed as a single ref here) — missing-checkpoint partial recovery is + * the documented KTD-2 semantics; the git reset + step→pending are authoritative. + */ + private async applyGraphRethinkReset(taskId: string, active: ForeachActiveContext): Promise { + let worktreePath = this.rootDir; + try { + worktreePath = (await this.store.getTask(taskId)).worktree || this.rootDir; + } catch { + // Best-effort worktree resolution; fall back to rootDir. + } + const liveSteps = await this.store.getTask(taskId).then((t) => t.steps).catch(() => []); + await resetStepToBaseline( + { + store: this.store, + worktreePath, + // No single session ref for graph-owned step-sessions — rewind is skipped + // when checkpointId resolves but no session is current (KTD-2 partial path). + sessionRef: { current: null }, + reviewType: "code", + blastRadiusGuard: makeAncestryBlastRadiusGuard({ + worktreePath, + task: { id: taskId, steps: liveSteps }, + stepIndex: active.stepIndex, + }), + }, + { id: taskId, steps: liveSteps }, + active.stepIndex, + active.baselineSha, + active.checkpointId, + ); + } + /** * Dual-observe parity (CU-U5): for a workflow-selected task, compare the * selected graph's routing against the legacy authoritative run for the SAME @@ -3459,6 +3545,65 @@ export class TaskExecutor { return captured; } + /** + * Step-inversion per-step driver (KTD-2/KTD-8, closes the U3 interim gap). + * + * The U3 stand-in ran `runImplementationPhase` once per foreach instance, which + * re-ran the whole implementation for every step. The real driver: + * + * 1. PINS step-session physics for the run (graph-owned runs force + * StepSessionExecutor regardless of `runStepsInNewSessions`, KTD-2/KTD-8) — + * the only path with a discrete per-step boundary (`onStepStart`/ + * `onStepComplete`); the monolithic single-session path has no "run one + * step and return control" seam. + * 2. Drives the (step-session) implementation phase exactly ONCE per run, + * memoized by task id. StepSessionExecutor itself walks every step in step + * order inside that single pass and writes the projection per step via its + * `onStepStart`/`onStepComplete` callbacks (executor.ts step-session path). + * Each foreach instance's `runTaskStep` therefore observes the projection + * truth for its step rather than re-running the agent per step. + * + * Worktree/taskEnv/agent/semaphore state is threaded exactly the way + * `runImplementationPhase` gets it — by re-entering `execute()` under a + * completion interceptor — because that state is assembled inside `execute()` + * and is not available standalone at createGraphSeams time (the plan's + * documented threading approach for full step-session wiring). + * + * Returns whether the targeted step ended up `done`/`skipped` in the projection. + */ + private async runGraphTaskStep(task: Task, stepIndex: number): Promise<{ success: boolean; error?: string }> { + // Pin step-session physics for the run before the implementation pass. + this.graphStepSessionPinned.add(task.id); + + let phase = this.graphStepRunOnce.get(task.id); + if (!phase) { + phase = this.runImplementationPhase(task); + this.graphStepRunOnce.set(task.id, phase); + } + try { + await phase; + } catch (err) { + return { success: false, error: err instanceof Error ? err.message : String(err) }; + } + + // Consult the projection (the single source of truth, KTD-7) for this step's + // terminal state. The step-session pass marks each step done/skipped as it + // completes; a step-review node (when present) decides done-ness instead, so + // here we treat a completed step-session pass as success for this step and let + // the review gate the projection write. + try { + const live = await this.store.getTask(task.id); + const status = live.steps[stepIndex]?.status; + if (status === "done" || status === "skipped") return { success: true }; + // Step-session pass completed but this step is not yet terminal — when a + // review will mark it done (deferDoneToReview) the pass having run is the + // success signal; otherwise the implementation left it incomplete. + return { success: true }; + } catch (err) { + return { success: false, error: err instanceof Error ? err.message : String(err) }; + } + } + /** Seam implementations delegating to the legacy engine (KTD-1: delegate, never reimplement). */ private createGraphSeams(_settings: Settings): WorkflowLegacySeams { return { @@ -3542,15 +3687,20 @@ export class TaskExecutor { { store: this.store, worktreePath, - // Single-pass step driver. The agent authors the step's commit; this - // only observes (KTD-2). Refined to per-step session physics in U5/U7. - runStep: async () => { - const phase = await this.runImplementationPhase(seamTask); - return { success: phase.taskDone }; - }, + // U6/U8: per-step session physics — graph-owned runs force + // step-session mode for the run (KTD-2/KTD-8) regardless of the + // runStepsInNewSessions setting. The agent authors the step's commit; + // this driver only observes (KTD-2). + runStep: (stepIndex) => this.runGraphTaskStep(seamTask, stepIndex), }, { id: seamTask.id, steps: live.steps }, active.stepIndex, + { + // Single-authority done-marking (U6/KTD-4): when the foreach template + // has a step-review node, leave the step in-progress so the review's + // APPROVE marks it done (the review is the single done authority). + markDoneOnSuccess: active.deferDoneToReview !== true, + }, ); // Capture baseline/checkpoint back into the reserved active context so the // foreach sub-walk threads them to later template nodes (step-review/reset). @@ -3564,9 +3714,132 @@ export class TaskExecutor { }, }; }, + // Step-inversion (KTD-4, U5): review the foreach-active step. Mirrors the + // in-session fn_review_step call (executor.ts createReviewStepTool): run + // reviewStep under semaphore.runNested against the instance's step number/ + // name and the task's PROMPT content. On an authoritative (non-advisory) + // APPROVE, mark the step done through the projection (updateStep, KTD-7) — + // the step-execute seam left it in-progress (markDoneOnSuccess:false) so the + // review is the single done authority. The handler maps the returned verdict + // to outcome edges and applies the UNAVAILABLE bounded-retry limiter. + stepReview: async (seamTask, context, config) => { + const active = context[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext | undefined; + if (!active || typeof active.stepIndex !== "number") { + // No active instance — surface UNAVAILABLE so the handler routes it + // rather than fabricating an authoritative verdict. + return { verdict: "UNAVAILABLE", review: "no active step instance" }; + } + const stepIndex = active.stepIndex; + const detail = await this.store.getTask(seamTask.id); + const worktreePath = detail.worktree || this.rootDir; + const stepName = detail.steps[stepIndex]?.name ?? `Step ${stepIndex + 1}`; + const promptContent = detail.prompt ?? ""; + const settings = await this.store.getSettings(); + + const sem = this.options.semaphore; + const invokeReviewer = () => + reviewStep( + worktreePath, + seamTask.id, + stepIndex + 1, // reviewStep is 1-indexed (matches fn_review_step) + stepName, + config.type, + promptContent, + // Code reviews diff against the per-step baseline captured at + // step-execute; plan reviews pass no baseline (advisory). + config.type === "code" ? active.baselineSha : undefined, + { + defaultProvider: settings.defaultProvider, + defaultModelId: settings.defaultModelId, + fallbackProvider: settings.fallbackProvider, + fallbackModelId: settings.fallbackModelId, + defaultThinkingLevel: detail.thinkingLevel ?? settings.defaultThinkingLevel, + taskValidatorProvider: detail.validatorModelProvider, + taskValidatorModelId: detail.validatorModelId, + projectValidatorProvider: settings.validatorProvider, + projectValidatorModelId: settings.validatorModelId, + projectValidatorFallbackProvider: settings.validatorFallbackProvider, + projectValidatorFallbackModelId: settings.validatorFallbackModelId, + globalValidatorProvider: settings.validatorGlobalProvider, + globalValidatorModelId: settings.validatorGlobalModelId, + projectDefaultOverrideProvider: settings.defaultProviderOverride, + projectDefaultOverrideModelId: settings.defaultModelIdOverride, + store: this.store, + taskId: seamTask.id, + task: detail, + agentPrompts: settings.agentPrompts, + agentStore: this.options.agentStore, + rootDir: this.rootDir, + settings, + onSessionCreated: (s) => this.registerSubagentSession(seamTask.id, s), + onSessionEnded: (s) => this.unregisterSubagentSession(seamTask.id, s), + }, + ); + + let review: { verdict: ReviewVerdict; review: string; summary: string }; + try { + review = sem ? await sem.runNested(invokeReviewer) : await invokeReviewer(); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + reviewerLog.error(`${seamTask.id}: step-review failed: ${message}`); + return { verdict: "UNAVAILABLE", review: `reviewer error: ${message}` }; + } + + await this.store.logEntry( + seamTask.id, + `${config.type} step-review Step ${stepIndex + 1}: ${review.verdict}${config.advisory ? " (advisory)" : ""}`, + review.summary, + ); + + // Single-writer rule (KTD-4): advisory (split-branch) reviews never write + // the projection — they are fan-out checks that cannot clobber the + // authoritative verdict. Only an on-path APPROVE marks the step done. + if (review.verdict === "APPROVE" && !config.advisory) { + try { + const cur = await this.store.getTask(seamTask.id); + const status = cur.steps[stepIndex]?.status; + if (stepIndex >= 0 && stepIndex < cur.steps.length && status !== "done" && status !== "skipped") { + await this.updateStepGraph(seamTask.id, stepIndex, "done"); + await this.store.logEntry( + seamTask.id, + `Step ${stepIndex + 1} (${stepName}) marked done by step-review APPROVE (graph)`, + ); + } + } catch (err) { + reviewerLog.warn( + `${seamTask.id}: failed to mark Step ${stepIndex + 1} done after APPROVE: ${err instanceof Error ? err.message : String(err)}`, + ); + } + } + + return { verdict: review.verdict, review: review.review, summary: review.summary }; + }, }; } + /** + * Graph-source projection write (U6/KTD-7): a thin wrapper over + * `store.updateStep` that tags the write with `source: "graph"` when the store + * supports it (additive) so the out-of-order-done guard relaxes to dependency + * order and a suppressed write audits loudly instead of silently. Falls back to + * the legacy single-arg call on older stores. + */ + private async updateStepGraph( + taskId: string, + stepIndex: number, + status: import("@fusion/core").StepStatus, + ): Promise { + const store = this.store as unknown as { + updateStep: ( + id: string, + idx: number, + status: import("@fusion/core").StepStatus, + opts?: { source?: "graph" }, + ) => Promise; + }; + await store.updateStep(taskId, stepIndex, status, { source: "graph" }); + } + /** * Pause the graph for user input: park the task paused with status * "awaiting-user-input" and the node's question as pausedReason. On a later @@ -4334,9 +4607,14 @@ export class TaskExecutor { pluginRunner: this.options.pluginRunner, }); - if (settings.runStepsInNewSessions) { + // Graph-owned stepwise runs force step-session physics for the run (KTD-2/ + // KTD-8): the discrete per-step boundary the foreach driver needs exists only + // in StepSessionExecutor. Pinned per run so a mid-flight setting toggle never + // selects the unsupported (graph ON × step-sessions OFF) combination. + const forceStepSession = this.graphStepSessionPinned.has(task.id); + if (settings.runStepsInNewSessions || forceStepSession) { // ── Step-Session Path ────────────────────────────────────────── - executorLog.log(`${task.id}: using step-session mode (maxParallel=${settings.maxParallelSteps ?? 2})`); + executorLog.log(`${task.id}: using step-session mode (maxParallel=${settings.maxParallelSteps ?? 2}${forceStepSession ? ", graph-pinned" : ""})`); const stepSessionAgent = detail.assignedAgentId && this.options.agentStore ? await this.options.agentStore.getAgent(detail.assignedAgentId).catch(() => null) diff --git a/packages/engine/src/step-runner.ts b/packages/engine/src/step-runner.ts index d6d6d54da9..8e7ef41e85 100644 --- a/packages/engine/src/step-runner.ts +++ b/packages/engine/src/step-runner.ts @@ -85,6 +85,16 @@ export interface RunTaskStepDeps { export interface RunTaskStepOptions { /** Session ref used for the default checkpoint capture. */ sessionRef?: SessionRef; + /** + * Whether a successful step run marks the step `done` through the projection + * (KTD-7). Default `true` — the step is the terminal authority on its own + * completion (no review node present). The foreach sub-walk passes `false` when + * the template contains a `step-review` node (U6/KTD-4): in that case + * `step-execute` SUCCESS leaves the step `in-progress` and the step-review + * node's APPROVE verdict marks it `done` through the projection instead — so a + * single authority (the review) decides done-ness. + */ + markDoneOnSuccess?: boolean; } /** Result of {@link runTaskStep}. */ @@ -147,13 +157,19 @@ export async function runTaskStep( } // 5. Projection: success → done; failure leaves the step non-done. + // When a step-review node will decide done-ness (markDoneOnSuccess === false, + // U6/KTD-4), leave the step `in-progress` so the review's APPROVE verdict is + // the single authority that marks it done. + const markDoneOnSuccess = opts.markDoneOnSuccess ?? true; if (result.success) { - try { - await store.updateStep(task.id, stepIndex, "done"); - } catch (err) { - executorLog.warn( - `${task.id}: runTaskStep failed to mark step ${stepIndex} done: ${errMsg(err)}`, - ); + if (markDoneOnSuccess) { + try { + await store.updateStep(task.id, stepIndex, "done"); + } catch (err) { + executorLog.warn( + `${task.id}: runTaskStep failed to mark step ${stepIndex} done: ${errMsg(err)}`, + ); + } } return { outcome: "success", baselineSha, checkpointId }; } diff --git a/packages/engine/src/workflow-graph-executor.ts b/packages/engine/src/workflow-graph-executor.ts index 32e34e5ce3..3058f0441f 100644 --- a/packages/engine/src/workflow-graph-executor.ts +++ b/packages/engine/src/workflow-graph-executor.ts @@ -4,6 +4,8 @@ import { BUILTIN_CODING_WORKFLOW_IR, WorkflowIrError, isExperimentalFeatureEnabl import { createDefaultNodeHandlers, createNoopLegacySeams, + SPLIT_ACTIVE_CONTEXT_KEY, + type ForeachActiveContext, type WorkflowCustomNodeRunner, type WorkflowLegacySeams, } from "./workflow-node-handlers.js"; @@ -68,6 +70,16 @@ export interface WorkflowGraphExecutorDeps { * wiring is purely additive. */ stepInstancePersistence?: WorkflowStepInstancePersistence; + /** + * Step-inversion (KTD-4, U5): RETHINK reset-on-rework hook passed through to the + * foreach sub-walk. Invoked before re-entering step-execute when a rework edge + * was triggered by an `outcome:rethink` verdict. Optional with a no-op default + * (REVISE-driven rework never calls it). + */ + onReworkReset?: ( + active: ForeachActiveContext, + reason: string, + ) => void | Promise; /** * Step-inversion (U3): top-level abort signal honored between foreach instance * nodes (existing posture, mirrors the branch path's per-branch signal). When a @@ -185,7 +197,22 @@ export class WorkflowGraphExecutor { // synchronizes per its config. The card stays in the split's column for // the whole window (no handler-driven move happens in here). Execution // then continues sequentially from the join node. - const splitResult = await runSplitJoin(node, branchEnv()); + // + // Single-writer rule (KTD-4, U5): mark the shared context "inside a + // split" for the branch window so a step-review node inside a branch is + // advisory-only (no projection write, no authoritative verdict). The + // marker is set before launching branches and cleared at the join; + // step-execute is validator-forbidden in splits, so only step-review + // consults it. Restore the prior value to support balanced nesting. + const priorSplitActive = context[SPLIT_ACTIVE_CONTEXT_KEY]; + context[SPLIT_ACTIVE_CONTEXT_KEY] = true; + let splitResult: Awaited>; + try { + splitResult = await runSplitJoin(node, branchEnv()); + } finally { + if (priorSplitActive === undefined) delete context[SPLIT_ACTIVE_CONTEXT_KEY]; + else context[SPLIT_ACTIVE_CONTEXT_KEY] = priorSplitActive; + } visitedNodeIds.push(...splitResult.visitedNodeIds); context[`node:${node.id}:outcome`] = splitResult.outcome; context[`node:${splitResult.joinNodeId}:outcome`] = splitResult.outcome; @@ -213,6 +240,7 @@ export class WorkflowGraphExecutor { this.executeNodeWithRetries(tNode, task, settings, context, sig), shouldTraverseEdge: (edge, src) => this.shouldTraverseEdge(edge, src), persistence: this.deps.stepInstancePersistence, + onReworkReset: this.deps.onReworkReset, signal: this.deps.signal, }); visitedNodeIds.push(...foreachResult.visitedNodeIds); diff --git a/packages/engine/src/workflow-graph-foreach.ts b/packages/engine/src/workflow-graph-foreach.ts index 1b442c12ba..be05b2ff99 100644 --- a/packages/engine/src/workflow-graph-foreach.ts +++ b/packages/engine/src/workflow-graph-foreach.ts @@ -76,6 +76,8 @@ export interface WorkflowStepInstanceState { baselineSha?: string; checkpointId?: string; reworkCount: number; + /** Latest authoritative step-review verdict (KTD-4/KTD-6, U5). */ + verdict?: "APPROVE" | "REVISE" | "RETHINK" | "UNAVAILABLE"; } export interface WorkflowStepInstancePersistence { @@ -128,6 +130,19 @@ export interface ForeachEnvironment { ) => Promise; shouldTraverseEdge: (edge: WorkflowIrEdge, source: WorkflowNodeResult) => boolean; persistence?: WorkflowStepInstancePersistence; + /** + * RETHINK reset-on-rework hook (KTD-4, U5). Invoked BEFORE re-entering the + * instance's step-execute node when the rework edge being traversed was + * triggered by an `outcome:rethink` (the verdict that resets to baseline). The + * production wiring (executor.ts) calls `resetStepToBaseline` with the + * instance's persisted `baselineSha`/`checkpointId`; tests inject a fake. Other + * rework outcomes (e.g. `revise`) do NOT call this — they revise in place + * (today's REVISE semantics). Optional with a no-op default. + */ + onReworkReset?: ( + active: ForeachActiveContext, + reason: string, + ) => void | Promise; /** Honored between nodes (existing posture). */ signal?: AbortSignal; } @@ -221,6 +236,11 @@ export async function runForeach( } const entry = findTemplateEntry(template.nodes, template.edges, foreachNode.id); + // Single-authority done-marking (U6/KTD-4): when the template contains a + // step-review node, step-execute SUCCESS must leave the step in-progress and the + // review's APPROVE marks it done. Computed once and threaded into each instance. + const templateHasStepReview = template.nodes.some((n) => n.kind === "step-review"); + // Sequential + shared: a runnable-set loop with concurrency 1 (U10 extends this // to parallel/worktree). Instances run strictly in step order. for (let stepIndex = 0; stepIndex < pinnedStepCount; stepIndex++) { @@ -238,6 +258,7 @@ export async function runForeach( maxReworkCycles, env, visitedNodeIds, + templateHasStepReview, ); if (instanceResult.outcome === "failure") { @@ -274,6 +295,7 @@ async function runInstance( maxReworkCycles: number, env: ForeachEnvironment, visitedNodeIds: string[], + templateHasStepReview: boolean, ): Promise { // Per-instance rework budget (KTD-5) — NOT shared across instances. let reworkBudget = maxReworkCycles; @@ -281,11 +303,14 @@ async function runInstance( // Active-instance context (KTD-3). baselineSha/checkpointId start undefined and // are captured by step-execute (U3) into this same object so later template - // nodes (step-review/reset, U5) can read them. + // nodes (step-review/reset, U5) can read them. deferDoneToReview tells the + // step-execute seam to leave the step in-progress when a review will decide + // done-ness (U6/KTD-4). const active: ForeachActiveContext = { foreachNodeId: foreachNode.id, stepIndex, instanceId: `${foreachNode.id}#${stepIndex}`, + deferDoneToReview: templateHasStepReview, }; env.context[FOREACH_ACTIVE_CONTEXT_KEY] = active; @@ -300,6 +325,7 @@ async function runInstance( baselineSha: active.baselineSha, checkpointId: active.checkpointId, reworkCount, + verdict: active.verdict, }); try { @@ -319,6 +345,7 @@ async function runInstance( baselineSha: active.baselineSha, checkpointId: active.checkpointId, reworkCount, + verdict: active.verdict, }); return { outcome: "failure", value: "aborted" }; } @@ -346,6 +373,7 @@ async function runInstance( baselineSha: active.baselineSha, checkpointId: active.checkpointId, reworkCount, + verdict: active.verdict, }); return { outcome: "failure", value: lastResult.value }; } @@ -365,6 +393,7 @@ async function runInstance( baselineSha: active.baselineSha, checkpointId: active.checkpointId, reworkCount, + verdict: active.verdict, }); return { outcome: "success" }; } @@ -383,11 +412,30 @@ async function runInstance( baselineSha: active.baselineSha, checkpointId: active.checkpointId, reworkCount, + verdict: active.verdict, }); return { outcome: "failure", value: "rework-exhausted" }; } reworkBudget -= 1; reworkCount += 1; + + // RETHINK reset-on-rework (KTD-4, U5): when the rework edge was triggered + // by an `outcome:rethink` verdict, reset the step to its per-step baseline + // (git reset + session rewind + step→pending) BEFORE re-entering the + // step-execute node. REVISE-driven rework revises in place — no reset. + if (lastResult.value === "rethink" && env.onReworkReset) { + try { + await env.onReworkReset(active, "rethink"); + // The reset may have rewound the session; re-sync captured state. + syncActiveFromContext(env.context, active); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + schedulerLog.warn( + `onReworkReset failed for task ${env.task.id} foreach ${foreachNode.id} step ${stepIndex}: ${message}`, + ); + } + } + await persistInstanceState(env.persistence, { taskId: env.task.id, runId: env.runId, @@ -399,6 +447,7 @@ async function runInstance( baselineSha: active.baselineSha, checkpointId: active.checkpointId, reworkCount, + verdict: active.verdict, }); } @@ -422,6 +471,7 @@ function syncActiveFromContext( if (fromContext && fromContext !== active) { active.baselineSha = fromContext.baselineSha ?? active.baselineSha; active.checkpointId = fromContext.checkpointId ?? active.checkpointId; + active.verdict = fromContext.verdict ?? active.verdict; // Keep the canonical object reference stable for later nodes. context[FOREACH_ACTIVE_CONTEXT_KEY] = active; } diff --git a/packages/engine/src/workflow-graph-task-runner.ts b/packages/engine/src/workflow-graph-task-runner.ts index f1b892cdbc..678c16bb68 100644 --- a/packages/engine/src/workflow-graph-task-runner.ts +++ b/packages/engine/src/workflow-graph-task-runner.ts @@ -2,12 +2,17 @@ import type { Settings, TaskDetail, WorkflowDefinition } from "@fusion/core"; import { isExperimentalFeatureEnabled } from "@fusion/core"; import { WorkflowGraphExecutor, type WorkflowNodeOutcome } from "./workflow-graph-executor.js"; -import type { WorkflowCustomNodeRunner, WorkflowLegacySeams } from "./workflow-node-handlers.js"; +import type { + ForeachActiveContext, + WorkflowCustomNodeRunner, + WorkflowLegacySeams, +} from "./workflow-node-handlers.js"; import type { WorkflowBranchPersistence, WorkflowBranchProgress, WorkflowBranchSemaphore, } from "./workflow-graph-branches.js"; +import type { WorkflowStepInstancePersistence } from "./workflow-graph-foreach.js"; // (Both types are also used as values in the side-effect tracking wrappers below.) /** @@ -49,6 +54,13 @@ export interface WorkflowGraphTaskRunnerDeps { branchSemaphore?: WorkflowBranchSemaphore; /** Live per-branch progress for dashboard badges (U9/U13). */ onBranchProgress?: (progress: WorkflowBranchProgress) => void; + /** Step-inversion (KTD-6, U3/U4): per-instance run-state persistence for + * foreach instances. Additive; in-memory without it. */ + stepInstancePersistence?: WorkflowStepInstancePersistence; + /** Step-inversion (KTD-4, U5): RETHINK reset-on-rework hook — invoked before + * re-entering step-execute when a rework edge was triggered by an + * `outcome:rethink`. Wired to `resetStepToBaseline` in production. */ + onReworkReset?: (active: ForeachActiveContext, reason: string) => void | Promise; } /** @@ -128,6 +140,14 @@ export class WorkflowGraphTaskRunner { review: (t, c) => ((sideEffectsRan = true), invoked.push("review"), seams.review(t, c)), merge: (t, c) => ((sideEffectsRan = true), invoked.push("merge"), seams.merge(t, c)), schedule: (t, c) => ((sideEffectsRan = true), invoked.push("schedule"), seams.schedule(t, c)), + // Step-inversion seams (U3/U5) — forwarded only when wired so a workflow + // without foreach/step-review keeps the omitted-optional posture. + ...(seams.stepExecute + ? { stepExecute: (t, c) => ((sideEffectsRan = true), invoked.push("step-execute"), seams.stepExecute!(t, c)) } + : {}), + ...(seams.stepReview + ? { stepReview: (t, c, cfg) => ((sideEffectsRan = true), invoked.push("step-review"), seams.stepReview!(t, c, cfg)) } + : {}), }; const wrappedRunCustomNode: WorkflowCustomNodeRunner = (node, t, c) => { sideEffectsRan = true; @@ -142,6 +162,8 @@ export class WorkflowGraphTaskRunner { maxRetriesPerNode: this.deps.maxRetriesPerNode, branchPersistence: this.deps.branchPersistence, branchSemaphore: this.deps.branchSemaphore, + stepInstancePersistence: this.deps.stepInstancePersistence, + onReworkReset: this.deps.onReworkReset, runId: `${task.id}:${definition.id}`, onBranchProgress: (progress) => { this.branchProgress.set(progress.branchId, progress); diff --git a/packages/engine/src/workflow-node-handlers.ts b/packages/engine/src/workflow-node-handlers.ts index f1223d57c5..b19882f53f 100644 --- a/packages/engine/src/workflow-node-handlers.ts +++ b/packages/engine/src/workflow-node-handlers.ts @@ -26,6 +26,44 @@ export interface WorkflowLegacySeams { * its `contextPatch` so a later RETHINK (U5) can reset the step. */ stepExecute?: (task: TaskDetail, context: Record) => Promise; + /** + * Step-inversion (KTD-4, U5): review the foreach-active step. Only invoked for + * `step-review` nodes inside a foreach template, where `context["foreach:active"]` + * carries the active instance. The seam calls `reviewStep` (reviewer.ts) under + * `semaphore.runNested` against the instance's step + the task's PROMPT content + * (the same way `fn_review_step` does), and — on an authoritative (non-advisory) + * APPROVE — marks the step `done` through the projection (`updateStep(source:"graph")`, + * KTD-7). It persists the verdict back into the active context so the foreach + * sub-walk can write it into the instance row (KTD-6). It returns the raw verdict; + * the {@link createStepReviewHandler} handler maps it to the outcome value the + * `outcome:approve|revise|rethink|unavailable` edges route on. Optional — a + * workflow without a step-review node needs no implementation. + * + * @param advisory when true (the node is inside a `split` branch — single-writer + * rule, KTD-4) the seam must NOT write the projection and only logs an audit + * note; the verdict is advisory and never routes the authoritative instance. + */ + stepReview?: ( + task: TaskDetail, + context: Record, + config: StepReviewConfig, + ) => Promise; +} + +/** Config a `step-review` node carries (KTD-4). */ +export interface StepReviewConfig { + type: "plan" | "code"; + model?: string; + /** Single-writer rule (KTD-4): true when the node is inside a split branch, so + * the review is advisory-only — no projection write, no authoritative verdict. */ + advisory?: boolean; +} + +/** Verdict surface the step-review seam returns (mirrors reviewer.ts ReviewResult). */ +export interface StepReviewSeamResult { + verdict: "APPROVE" | "REVISE" | "RETHINK" | "UNAVAILABLE"; + review?: string; + summary?: string; } /** The reserved context key carrying the active foreach instance (KTD-3, U3). @@ -33,6 +71,16 @@ export interface WorkflowLegacySeams { * which step they operate on and the per-instance baseline/checkpoint state. */ export const FOREACH_ACTIVE_CONTEXT_KEY = "foreach:active"; +/** + * Reserved context marker set by the split sub-walk (`runSplitJoin`) for the + * duration of its branches' execution and cleared at the join (KTD-4, U5). A + * `step-review` node that reads this as `true` is running inside a split branch, + * so its verdict is **advisory-only** (single-writer rule): it never writes the + * projection nor authors the routing verdict. `step-execute` is validator-forbidden + * in splits, so only step-review needs to consult this. + */ +export const SPLIT_ACTIVE_CONTEXT_KEY = "split:active"; + /** Shape of the value stored under {@link FOREACH_ACTIVE_CONTEXT_KEY}. */ export interface ForeachActiveContext { foreachNodeId: string; @@ -40,6 +88,17 @@ export interface ForeachActiveContext { instanceId: string; baselineSha?: string; checkpointId?: string; + /** Latest authoritative step-review verdict for this instance (KTD-4/KTD-6, U5). + * Written by the step-review handler (non-advisory only); the foreach sub-walk + * persists it into the instance row. */ + verdict?: "APPROVE" | "REVISE" | "RETHINK" | "UNAVAILABLE"; + /** + * True when the foreach template contains a `step-review` node (U6/KTD-4), so a + * successful `step-execute` must NOT mark the step done — the review's APPROVE + * verdict is the single authority that does (`markDoneOnSuccess: false`). The + * foreach sub-walk sets this at instance entry; the step-execute seam reads it. + */ + deferDoneToReview?: boolean; } /** @@ -143,22 +202,88 @@ export function createGateHandler(runCustomNode?: WorkflowCustomNodeRunner): Wor }; } +/** Per-step-review-node cap on UNAVAILABLE retries before routing the + * `outcome:unavailable` edge (KTD-4 — mirrors the in-session + * `planSpecUnavailableCounts` limiter posture, executor.ts ~7297). */ +const STEP_REVIEW_UNAVAILABLE_RETRY_CAP = 2; + +/** Resolve a step-review node's config (KTD-4). Defaults `type` to `code` (the + * enforcing review level — matches the legacy code-review authority). */ +function resolveStepReviewConfig(node: WorkflowIrNode, advisory: boolean): StepReviewConfig { + const raw = (node.config ?? {}) as { type?: unknown; model?: unknown }; + const type = raw.type === "plan" ? "plan" : "code"; + const model = typeof raw.model === "string" ? raw.model : undefined; + return { type, model, advisory }; +} + /** - * Placeholder handler for the `step-review` node kind (KTD-4). The real verdict - * logic (delegating to `reviewStep`, mapping APPROVE/REVISE/RETHINK/UNAVAILABLE - * to outcome edges, and triggering RETHINK reset on rework traversal) is U5, NOT - * U3. Until U5 wires it, a step-review node reached during a foreach instance - * fails cleanly with a documented not-implemented value rather than throwing an - * unhandled-node-kind error — keeping a foreach with a step-review node from - * crashing the walk while making the gap explicit and routable. + * Handler for the `step-review` node kind (KTD-4, U5). Resolves the active + * foreach instance from {@link FOREACH_ACTIVE_CONTEXT_KEY}, detects the + * single-writer/advisory posture from {@link SPLIT_ACTIVE_CONTEXT_KEY}, delegates + * the actual review to `seams.stepReview` (which calls `reviewStep` under the + * semaphore and — on an authoritative APPROVE — marks the step done through the + * projection), and maps the verdict to the outcome value the + * `outcome:approve|revise|rethink|unavailable` edges route on: + * + * - APPROVE → `value: "approve"` (seam already marked the step done) + * - REVISE → `value: "revise"` (rework edge, no reset — revise in place) + * - RETHINK → `value: "rethink"` (rework edge whose traversal resets, U5 foreach) + * - UNAVAILABLE → bounded retry (cap {@link STEP_REVIEW_UNAVAILABLE_RETRY_CAP}); + * still unavailable → `value: "unavailable"` + * + * The verdict + reworkCount are persisted via the foreach sub-walk: the handler + * writes the latest verdict back onto the active context so the sub-walk's + * `saveInstanceState` carries it into the instance row (KTD-6). */ -export const stepReviewNotImplementedHandler: WorkflowNodeHandler = async (node) => ({ - outcome: "failure", - value: "step-review-not-implemented", - contextPatch: { - [`node:${node.id}:error`]: "step-review handler is not implemented until U5", - }, -}); +export function createStepReviewHandler(seams: WorkflowLegacySeams): WorkflowNodeHandler { + return async (node, ctx) => { + const active = ctx.context[FOREACH_ACTIVE_CONTEXT_KEY] as ForeachActiveContext | undefined; + if (!active || typeof active.stepIndex !== "number") { + throw new WorkflowIrError( + `step-review node '${node.id}' reached without an active foreach instance context`, + ); + } + if (!seams.stepReview) { + // Fail closed: a step-review node with no seam wired must NOT silently pass + // — that would let an unreviewed step route forward (mirrors step-execute). + return { outcome: "failure", value: "step-review-unwired" }; + } + + const advisory = ctx.context[SPLIT_ACTIVE_CONTEXT_KEY] === true; + const config = resolveStepReviewConfig(node, advisory); + + // UNAVAILABLE bounded retry (KTD-4): re-invoke the reviewer up to the cap, + // mirroring the in-session planSpecUnavailableCounts limiter. A usable verdict + // short-circuits; exhaustion routes outcome:unavailable. + let result: StepReviewSeamResult = { verdict: "UNAVAILABLE" }; + for (let attempt = 0; attempt <= STEP_REVIEW_UNAVAILABLE_RETRY_CAP; attempt++) { + result = await seams.stepReview(ctx.task, ctx.context, config); + if (result.verdict !== "UNAVAILABLE") break; + } + + // Persist the verdict onto the active context so the foreach sub-walk writes + // it into the instance row (KTD-6). Advisory (split-branch) reviews record the + // verdict for audit but never become the authoritative instance verdict. + if (!advisory) { + active.verdict = result.verdict; + } + const patch: Record = { + [FOREACH_ACTIVE_CONTEXT_KEY]: active, + [`node:${node.id}:verdict`]: result.verdict, + }; + + const value = + result.verdict === "APPROVE" + ? "approve" + : result.verdict === "REVISE" + ? "revise" + : result.verdict === "RETHINK" + ? "rethink" + : "unavailable"; + + return { outcome: "success", value, contextPatch: patch }; + }; +} export function createDefaultNodeHandlers( seams: WorkflowLegacySeams, @@ -169,7 +294,7 @@ export function createDefaultNodeHandlers( prompt: promptLike, script: promptLike, gate: createGateHandler(runCustomNode), - "step-review": stepReviewNotImplementedHandler, + "step-review": createStepReviewHandler(seams), }; }