diff --git a/.changeset/fix-parallel-step-completion-order.md b/.changeset/fix-parallel-step-completion-order.md new file mode 100644 index 0000000000..c1336542c7 --- /dev/null +++ b/.changeset/fix-parallel-step-completion-order.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Allow dependency-ready workflow steps to finalize when earlier independent steps are still running. +category: fix +dev: Makes explicit step dependency metadata authoritative for every step-completion writer. diff --git a/packages/core/src/__tests__/postgres/store-update-step-order.pg.test.ts b/packages/core/src/__tests__/postgres/store-update-step-order.pg.test.ts index fe70e883ce..7a27776968 100644 --- a/packages/core/src/__tests__/postgres/store-update-step-order.pg.test.ts +++ b/packages/core/src/__tests__/postgres/store-update-step-order.pg.test.ts @@ -35,6 +35,104 @@ pgTest("TaskStore.updateStep step-order guard (PostgreSQL)", () => { expect(updated.steps[2].status).toBe("pending"); }); + it("allows an explicitly independent step to finish out of index order", async () => { + const store = h.store(); + const task = await h.createTaskWithSteps(); + await store.updateTask(task.id, { + steps: task.steps.map((step, index) => + index === 2 ? { ...step, dependsOn: [] } : step, + ), + }); + + const active = await store.updateStep(task.id, 2, "in-progress"); + expect(active.steps[2].status).toBe("in-progress"); + + const updated = await store.updateStep(task.id, 2, "done"); + + expect(updated.steps[0].status).toBe("pending"); + expect(updated.steps[1].status).toBe("pending"); + expect(updated.steps[2].status).toBe("done"); + expect(updated.currentStep).toBe(0); + + const remaining = await store.updateStep(task.id, 0, "done"); + expect(remaining.currentStep).toBe(1); + const finalized = await store.updateStep(task.id, 1, "done"); + expect(finalized.steps.every((step) => step.status === "done")).toBe(true); + expect(finalized.currentStep).toBe(3); + }); + + it("still blocks out-of-index completion when an explicit dependency is pending", async () => { + const store = h.store(); + const task = await h.createTaskWithSteps(); + await store.updateTask(task.id, { + steps: task.steps.map((step, index) => + index === 2 ? { ...step, dependsOn: [1] } : step, + ), + }); + + await store.updateStep(task.id, 1, "in-progress"); + const updated = await store.updateStep(task.id, 2, "done"); + + expect(updated.steps[2].status).toBe("pending"); + expect(updated.log.at(-1)?.action).toContain("dependency step 1"); + }); + + it("allows a satisfied explicit dependency while an unrelated earlier step is pending", async () => { + const store = h.store(); + const task = await h.createTaskWithSteps(); + await store.updateTask(task.id, { + steps: task.steps.map((step, index) => + index === 1 + ? { ...step, dependsOn: [] } + : index === 2 + ? { ...step, dependsOn: [1] } + : step, + ), + }); + + await store.updateStep(task.id, 1, "done"); + const updated = await store.updateStep(task.id, 2, "done"); + + expect(updated.steps[0].status).toBe("pending"); + expect(updated.steps[1].status).toBe("done"); + expect(updated.steps[2].status).toBe("done"); + }); + + it("honors an explicit empty dependency list for graph-source completion", async () => { + const store = h.store(); + const task = await h.createTaskWithSteps(); + await store.updateTask(task.id, { + steps: task.steps.map((step, index) => + index === 2 ? { ...step, dependsOn: [] } : step, + ), + }); + + const updated = await store.updateStep(task.id, 2, "done", { source: "graph" }); + + expect(updated.steps[2].status).toBe("done"); + }); + + it.each([{ dependsOn: [2] }, { dependsOn: [99] }, { dependsOn: [-1] }, { dependsOn: [1.5] }])( + "falls back to strict ordering for malformed dependsOn %j", + async ({ dependsOn }) => { + const store = h.store(); + const task = await h.createTaskWithSteps(); + await store.updateTask(task.id, { + steps: task.steps.map((step, index) => + index === 2 ? { ...step, dependsOn } : step, + ), + }); + + const updated = await store.updateStep(task.id, 2, "done"); + + expect(updated.steps[2].status).toBe("pending"); + expect(updated.log.map((entry) => entry.action)).toContainEqual( + expect.stringContaining("[integrity-warning] invalid dependsOn"), + ); + expect(updated.log.at(-1)?.action).toContain("earlier step 0"); + }, + ); + it("allows done when prior steps are skipped", async () => { const store = h.store(); const task = await h.createTaskWithSteps(); diff --git a/packages/core/src/task-store/merge-queue-ops.ts b/packages/core/src/task-store/merge-queue-ops.ts index 8afde2a355..b5cde90e11 100644 --- a/packages/core/src/task-store/merge-queue-ops.ts +++ b/packages/core/src/task-store/merge-queue-ops.ts @@ -54,13 +54,15 @@ async function appendProactiveStepStatus(store: TaskStore, taskId: string, messa } export async function updateStepImpl(store: TaskStore, id: string, stepIndex: number, status: import("../types.js").StepStatus, options?: { source?: "graph" },): Promise { + // FNXC:WorkflowStepOrdering 2026-07-20-20:05: // 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); + // behaviors diverge from a legacy write with no explicit dependency metadata: + // (a) the out-of-order-done guard uses DEPENDENCY order (a done write is + // legal when every dependsOn step — default: the immediately-preceding + // step — is done/skipped, KTD-11). Explicit dependsOn metadata is + // authoritative for every writer, not only graph-tagged writes; // (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 @@ -130,16 +132,39 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu 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: the - // step's dependsOn list (default = the immediately-preceding step when the - // annotation is absent — preserving sequential behavior, KTD-11). + // may go done. Explicit dependency metadata is authoritative regardless + // of which execution surface performs the write: parallel step sessions + // and graph foreach instances share the same Task.steps[] contract. When + // metadata is absent, legacy callers retain strict index order while a + // graph-source write defaults to the immediately preceding step (KTD-11). + const explicitDependencies = task.steps[stepIndex]?.dependsOn; + const hasExplicitDependencies = Array.isArray(explicitDependencies); + const validExplicitDependencies = + hasExplicitDependencies && + explicitDependencies.every( + (dependency) => + Number.isInteger(dependency) && + dependency >= 0 && + dependency < stepIndex, + ); + const dependencyOrdered = + validExplicitDependencies || (!hasExplicitDependencies && graphSource); + + if (hasExplicitDependencies && !validExplicitDependencies) { + const ts = new Date().toISOString(); + task.log.push({ + timestamp: ts, + action: + `[integrity-warning] invalid dependsOn metadata for step ${stepIndex} ` + + `(${task.steps[stepIndex].name}); using strict index-order completion guard`, + }); + } let blockingIndex = -1; let blockingStatus: import("../types.js").StepStatus | undefined; - if (graphSource) { - const deps = task.steps[stepIndex]?.dependsOn; + if (dependencyOrdered) { const depIndices = - Array.isArray(deps) && deps.length > 0 - ? deps + validExplicitDependencies + ? explicitDependencies! : stepIndex > 0 ? [stepIndex - 1] : []; @@ -164,12 +189,12 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu if (blockingIndex !== -1) { const ts = new Date().toISOString(); task.updatedAt = ts; - const kind = graphSource ? "dependency-order" : "out-of-order"; + const kind = dependencyOrdered ? "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}`, + `${dependencyOrdered ? "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. @@ -192,14 +217,13 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu task.steps[stepIndex].status = status; task.updatedAt = new Date().toISOString(); - // Advance currentStep to first non-done/non-skipped step + // Recompute from the full list: an out-of-index parallel step may have + // moved currentStep ahead of an earlier unfinished dependency branch. if (status === "done") { - while ( - task.currentStep < task.steps.length && - (task.steps[task.currentStep].status === "done" || task.steps[task.currentStep].status === "skipped") - ) { - task.currentStep++; - } + const firstUnfinished = task.steps.findIndex( + (step) => step.status !== "done" && step.status !== "skipped", + ); + task.currentStep = firstUnfinished === -1 ? task.steps.length : firstUnfinished; } else if (status === "in-progress") { task.currentStep = stepIndex; }