From f430f5db39971260d5e6949d074946a8f07a7376 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Mon, 29 Jun 2026 03:23:01 -0700 Subject: [PATCH] fix(FN-7228): respect workflow step dependencies --- .../fn-7228-step-session-dependencies.md | 7 ++ .../__tests__/store-update-step-order.test.ts | 22 +++- packages/core/src/store.ts | 19 +-- .../__tests__/step-session-executor.test.ts | 110 ++++++++---------- packages/engine/src/step-session-executor.ts | 29 +++-- 5 files changed, 107 insertions(+), 80 deletions(-) create mode 100644 .changeset/fn-7228-step-session-dependencies.md diff --git a/.changeset/fn-7228-step-session-dependencies.md b/.changeset/fn-7228-step-session-dependencies.md new file mode 100644 index 0000000000..9b95eb536a --- /dev/null +++ b/.changeset/fn-7228-step-session-dependencies.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Keep verification steps from starting before earlier workflow steps finish. +category: fix +dev: Step-session wave planning now respects graph step dependencies; unannotated steps remain sequential by default. 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 de5b0deb4d..8362137353 100644 --- a/packages/core/src/__tests__/store-update-step-order.test.ts +++ b/packages/core/src/__tests__/store-update-step-order.test.ts @@ -57,13 +57,13 @@ describe("TaskStore.updateStep step-order guard", () => { // ── U6: graph-source projection discipline (KTD-7/KTD-11) ────────────────── - it("graph source: done is legal for independent steps even when an earlier step is in-progress", async () => { - // Graph-owned step sessions can complete independent steps out of index order. - // Missing dependsOn means the graph/wave planner already decided the step was - // runnable; TaskStore must not invent a hidden previous-step dependency. + it("graph source: done is legal for explicitly independent steps even when an earlier step is in-progress", async () => { const store = harness.store(); const task = await harness.createTaskWithSteps(); await store.updateStep(task.id, 0, "pending"); + const primed = await store.getTask(task.id); + const steps = primed.steps.map((s, i) => (i === 2 ? { ...s, dependsOn: [] } : { ...s })); + await store.updateTask(task.id, { steps }); await store.updateStep(task.id, 1, "in-progress", { source: "graph" }); const updated = await store.updateStep(task.id, 2, "done", { source: "graph" }); @@ -73,6 +73,20 @@ describe("TaskStore.updateStep step-order guard", () => { expect(updated.log.some((e) => e.action.includes("Ignored out-of-order done for step 2"))).toBe(false); }); + it("graph source: missing dependsOn defaults to previous step and blocks early verification", async () => { + const store = harness.store(); + const task = await harness.createTaskWithSteps(); + await store.updateStep(task.id, 0, "pending"); + + await store.updateStep(task.id, 1, "in-progress", { source: "graph" }); + const updated = await store.updateStep(task.id, 2, "done", { source: "graph" }); + + expect(updated.steps[2].status).toBe("pending"); + expect( + updated.log.some((e) => e.action.includes("Ignored dependency-order done for step 2")), + ).toBe(true); + }); + it("graph source: explicit dependsOn still suppresses completion until dependencies finish", async () => { const store = harness.store(); const task = await harness.createTaskWithSteps(); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index a039507119..837e93dca5 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -9241,9 +9241,8 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} // (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 explicitly declared - // dependsOn step is done/skipped; absent dependsOn means the graph has - // already decided the step can run independently); + // 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 @@ -9294,20 +9293,22 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} if (status === "done") { // The set of predecessor steps that must be done/skipped before this step // may go done. Legacy: strict index order (every earlier step). Graph: - // only the step's explicit dependsOn list. When a graph-owned step session - // finishes before an earlier independent step, the task card must record - // that real completion instead of emitting a false "out-of-order" pause. + // the step's dependsOn list, with absent dependsOn defaulting to the + // immediately-preceding step. A deliberately empty dependsOn array is the + // opt-in for an independent graph step. /* - FNXC:WorkflowStepControl 2026-06-29-10:12: - Graph-owned execution can complete independent steps out of index order. Do not infer a hidden previous-step dependency in TaskStore projection; the graph scheduler/step-session wave planner is the authority that decided the step was runnable. Legacy fn_task_update keeps strict index order for non-graph sessions. + FNXC:WorkflowStepControl 2026-06-29-10:51: + Graph-owned execution may complete explicitly independent steps out of index order, but unannotated task plans are sequential by default. FN-7228 showed Testing & Verification starting while Preflight/implementation were still active because step-session planning treated missing dependencies as independent. Keep TaskStore projection consistent with the graph scheduler: absent dependsOn means previous-step dependency; explicit dependsOn: [] means independent. */ 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 + Array.isArray(deps) ? deps + : stepIndex > 0 + ? [stepIndex - 1] : []; for (const i of depIndices) { const priorStatus = task.steps[i]?.status; diff --git a/packages/engine/src/__tests__/step-session-executor.test.ts b/packages/engine/src/__tests__/step-session-executor.test.ts index e53750247e..c87a170e4f 100644 --- a/packages/engine/src/__tests__/step-session-executor.test.ts +++ b/packages/engine/src/__tests__/step-session-executor.test.ts @@ -60,6 +60,14 @@ function makeTaskDetail(overrides: Partial = {}): TaskDetail { }; } +function makeIndependentSteps(count: number): TaskDetail["steps"] { + return Array.from({ length: count }, (_, i) => ({ + name: `Step ${i}`, + status: "pending" as const, + dependsOn: [], + })); +} + // ── parseStepFileScopes tests ────────────────────────────────────────── describe("parseStepFileScopes", () => { @@ -289,40 +297,61 @@ describe("buildConflictMatrix", () => { // ── determineParallelWaves tests ─────────────────────────────────────── describe("determineParallelWaves", () => { + const independentSteps = (count: number) => + Array.from({ length: count }, () => ({ dependsOn: [] as number[] })); + it("maxParallel=1 → all steps sequential (one per wave)", () => { const scopes = new Map([ [0, ["src/a.ts"]], [1, ["src/b.ts"]], [2, ["src/c.ts"]], ]); - const waves = determineParallelWaves(scopes, 1); + const waves = determineParallelWaves(scopes, 1, independentSteps(3)); expect(waves).toHaveLength(3); expect(waves[0]).toEqual({ indices: [0], waveNumber: 0 }); expect(waves[1]).toEqual({ indices: [1], waveNumber: 1 }); expect(waves[2]).toEqual({ indices: [2], waveNumber: 2 }); }); - it("no conflicts → all steps in wave 0 (capped by maxParallel)", () => { + it("explicitly independent non-conflicting steps share waves capped by maxParallel", () => { const scopes = new Map([ [0, ["src/a.ts"]], [1, ["src/b.ts"]], [2, ["src/c.ts"]], [3, ["src/d.ts"]], ]); - const waves = determineParallelWaves(scopes, 2); + const waves = determineParallelWaves(scopes, 2, independentSteps(4)); // All non-conflicting, capped at 2 per wave → 2 waves expect(waves).toHaveLength(2); expect(waves[0].indices).toEqual([0, 1]); expect(waves[1].indices).toEqual([2, 3]); }); + it("unannotated steps are sequential so verification cannot start before implementation", () => { + const scopes = new Map([ + [0, ["preflight.md"]], + [1, ["src/implementation.ts"]], + [2, ["src/implementation.test.ts"]], + [3, ["verification.log"]], + ]); + + const waves = determineParallelWaves(scopes, 4, [ + {}, + {}, + {}, + {}, + ]); + + expect(waves.map((wave) => wave.indices)).toEqual([[0], [1], [2], [3]]); + }); + it("all steps conflict → each step in its own wave", () => { const scopes = new Map([ [0, ["src/a.ts"]], [1, ["src/a.ts"]], [2, ["src/a.ts"]], ]); - const waves = determineParallelWaves(scopes, 4); + const waves = determineParallelWaves(scopes, 4, independentSteps(3)); // All conflict with each other → each in own wave expect(waves).toHaveLength(3); expect(waves[0].indices).toEqual([0]); @@ -338,7 +367,7 @@ describe("determineParallelWaves", () => { [2, ["packages/engine"]], [3, ["packages/engine/src/executor.ts"]], ]); - const waves = determineParallelWaves(scopes, 2); + const waves = determineParallelWaves(scopes, 2, independentSteps(4)); // 0 and 1 conflict (0 is prefix of 1), 2 and 3 conflict (2 is prefix of 3) // 0 and 2 don't conflict → wave 0: [0, 2] // 1 and 3 don't conflict → wave 1: [1, 3] @@ -352,7 +381,7 @@ describe("determineParallelWaves", () => { [0, []], [1, []], ]); - const waves = determineParallelWaves(scopes, 4); + const waves = determineParallelWaves(scopes, 4, independentSteps(2)); expect(waves).toHaveLength(1); expect(waves[0].indices).toEqual([0, 1]); }); @@ -364,7 +393,7 @@ describe("determineParallelWaves", () => { [2, ["c.ts"]], [3, ["d.ts"]], ]); - const waves = determineParallelWaves(scopes, 2); + const waves = determineParallelWaves(scopes, 2, independentSteps(4)); expect(waves).toHaveLength(2); expect(waves[0].indices).toHaveLength(2); expect(waves[1].indices).toHaveLength(2); @@ -375,7 +404,7 @@ describe("determineParallelWaves", () => { [0, ["a.ts"]], [1, ["b.ts"]], ]); - const waves = determineParallelWaves(scopes, 1); + const waves = determineParallelWaves(scopes, 1, independentSteps(2)); for (let i = 0; i < waves.length; i++) { expect(waves[i].waveNumber).toBe(i); } @@ -387,7 +416,7 @@ describe("determineParallelWaves", () => { [1, ["b.ts"]], [2, ["c.ts"]], ]); - const waves = determineParallelWaves(scopes, 2); + const waves = determineParallelWaves(scopes, 2, independentSteps(3)); // 3 non-conflicting steps, max 2 → wave 0: [0,1], wave 1: [2] expect(waves).toHaveLength(2); expect(waves[0].indices).toEqual([0, 1]); @@ -1578,10 +1607,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -1623,10 +1649,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -1680,10 +1703,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -1732,10 +1752,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -1825,10 +1842,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); const semaphore = new AgentSemaphore(4); @@ -1868,10 +1882,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -1956,11 +1967,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - { name: "Step 2", status: "pending" }, - ], + steps: makeIndependentSteps(3), }); const settings = makeSettings({ maxParallelSteps: 3 }); @@ -2019,11 +2026,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - { name: "Step 2", status: "pending" }, - ], + steps: makeIndependentSteps(3), }); const settings = makeSettings({ maxParallelSteps: 3 }); @@ -2129,10 +2132,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -2181,11 +2181,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - { name: "Step 2", status: "pending" }, - ], + steps: makeIndependentSteps(3), }); const settings = makeSettings({ maxParallelSteps: 3 }); @@ -2331,10 +2327,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); @@ -2398,10 +2391,7 @@ describe("StepSessionExecutor", () => { const task = makeTaskDetail({ prompt, - steps: [ - { name: "Step 0", status: "pending" }, - { name: "Step 1", status: "pending" }, - ], + steps: makeIndependentSteps(2), }); const settings = makeSettings({ maxParallelSteps: 2 }); diff --git a/packages/engine/src/step-session-executor.ts b/packages/engine/src/step-session-executor.ts index 3a30695a25..fece1d8c60 100644 --- a/packages/engine/src/step-session-executor.ts +++ b/packages/engine/src/step-session-executor.ts @@ -294,11 +294,12 @@ function pathsOverlap(a: string[], b: string[]): boolean { } /** - * Group non-conflicting steps into parallel execution waves. + * Group dependency-ready, non-conflicting steps into parallel execution waves. * - * Uses a greedy left-to-right scan: for each wave, add steps that don't - * conflict with any step already in the wave. Each wave's `indices` array - * is capped at `maxParallel`. Remaining steps spill into subsequent waves. + * Uses a greedy left-to-right scan: for each wave, add steps whose dependencies + * have all been assigned to earlier waves and that don't conflict with any step + * already in the wave. Each wave's `indices` array is capped at `maxParallel`. + * Remaining steps spill into subsequent waves. * * @param stepScopes - Map from step index to array of file paths. * @param maxParallel - Maximum number of steps per wave (range 1–4). @@ -307,6 +308,7 @@ function pathsOverlap(a: string[], b: string[]): boolean { export function determineParallelWaves( stepScopes: Map, maxParallel: number, + steps: Array<{ dependsOn?: number[] }> = [], ): ParallelWave[] { const indices = [...stepScopes.keys()].sort((a, b) => a - b); if (indices.length === 0) return []; @@ -326,6 +328,8 @@ export function determineParallelWaves( for (const idx of indices) { if (assigned.has(idx)) continue; if (waveIndices.length >= clampedMax) break; + const deps = resolveStepSessionDependsOn(steps, idx); + if (!deps.every((dep) => assigned.has(dep))) continue; const pos = posMap.get(idx)!; @@ -341,8 +345,9 @@ export function determineParallelWaves( } if (waveIndices.length === 0) { - // Safety: if no steps could be added (shouldn't happen with diagonal=true), - // force-add the next unassigned step + // Safety: malformed/cyclic dependencies should not infinite-loop the + // executor. Force the next unassigned step and let downstream execution/ + // validation report the real problem. const next = indices.find((idx) => !assigned.has(idx)); if (next !== undefined) { waveIndices.push(next); @@ -362,6 +367,12 @@ export function determineParallelWaves( return waves; } +function resolveStepSessionDependsOn(steps: Array<{ dependsOn?: number[] }>, stepIndex: number): number[] { + const deps = steps[stepIndex]?.dependsOn; + if (Array.isArray(deps)) return deps; + return stepIndex > 0 ? [stepIndex - 1] : []; +} + // ── Step Prompt Builder ─────────────────────────────────────────────── /** @@ -765,7 +776,11 @@ export class StepSessionExecutor { } } - const waves = determineParallelWaves(stepScopes, this.maxParallel); + /* + * FNXC:WorkflowStepControl 2026-06-29-10:45: + * Step-session parallelism must respect the workflow graph's dependency model. FN-7228 showed Testing & Verification running while Preflight/implementation were still in progress because this planner considered only file conflicts. Unannotated steps are sequential by default; only explicit dependsOn metadata may unlock out-of-index parallelism. + */ + const waves = determineParallelWaves(stepScopes, this.maxParallel, taskDetail.steps ?? []); stepExecLog.log( `Executing ${stepCount} steps in ${waves.length} wave(s) for task ${taskDetail.id} ` +