fix(core): allow dependency-ready steps to finalize
Honor explicit workflow dependencies across completion writers, keep the progress cursor aligned with unfinished work, and fail closed when dependency metadata is malformed.
This commit is contained in:
7
.changeset/fix-parallel-step-completion-order.md
Normal file
7
.changeset/fix-parallel-step-completion-order.md
Normal file
@@ -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.
|
||||||
@@ -35,6 +35,104 @@ pgTest("TaskStore.updateStep step-order guard (PostgreSQL)", () => {
|
|||||||
expect(updated.steps[2].status).toBe("pending");
|
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 () => {
|
it("allows done when prior steps are skipped", async () => {
|
||||||
const store = h.store();
|
const store = h.store();
|
||||||
const task = await h.createTaskWithSteps();
|
const task = await h.createTaskWithSteps();
|
||||||
|
|||||||
@@ -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<Task> {
|
export async function updateStepImpl(store: TaskStore, id: string, stepIndex: number, status: import("../types.js").StepStatus, options?: { source?: "graph" },): Promise<Task> {
|
||||||
|
// FNXC:WorkflowStepOrdering 2026-07-20-20:05:
|
||||||
// Step-inversion projection discipline (U6/KTD-7). A `source: "graph"` write
|
// Step-inversion projection discipline (U6/KTD-7). A `source: "graph"` write
|
||||||
// is the workflow-graph executor projecting a foreach instance's lifecycle
|
// is the workflow-graph executor projecting a foreach instance's lifecycle
|
||||||
// (in-progress / done / pending) onto Task.steps[] with EXPLICIT indices. Three
|
// (in-progress / done / pending) onto Task.steps[] with EXPLICIT indices. Three
|
||||||
// behaviors diverge from the legacy (default) write:
|
// behaviors diverge from a legacy write with no explicit dependency metadata:
|
||||||
// (a) the out-of-order-done guard relaxes from strict index order to
|
// (a) the out-of-order-done guard uses DEPENDENCY order (a done write is
|
||||||
// DEPENDENCY order (a done write is legal when every dependsOn step —
|
// legal when every dependsOn step — default: the immediately-preceding
|
||||||
// default: the immediately-preceding step — is done/skipped, KTD-11);
|
// 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
|
// (b) a guard that DOES suppress a graph write logs an audit warning loudly
|
||||||
// (legacy stays silent — a graph suppression is a projection bug);
|
// (legacy stays silent — a graph suppression is a projection bug);
|
||||||
// (c) the auto-reinit-from-PROMPT.md path is bypassed (the graph pinned the
|
// (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") {
|
if (status === "done") {
|
||||||
// The set of predecessor steps that must be done/skipped before this step
|
// 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
|
// may go done. Explicit dependency metadata is authoritative regardless
|
||||||
// step's dependsOn list (default = the immediately-preceding step when the
|
// of which execution surface performs the write: parallel step sessions
|
||||||
// annotation is absent — preserving sequential behavior, KTD-11).
|
// 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 blockingIndex = -1;
|
||||||
let blockingStatus: import("../types.js").StepStatus | undefined;
|
let blockingStatus: import("../types.js").StepStatus | undefined;
|
||||||
if (graphSource) {
|
if (dependencyOrdered) {
|
||||||
const deps = task.steps[stepIndex]?.dependsOn;
|
|
||||||
const depIndices =
|
const depIndices =
|
||||||
Array.isArray(deps) && deps.length > 0
|
validExplicitDependencies
|
||||||
? deps
|
? explicitDependencies!
|
||||||
: stepIndex > 0
|
: stepIndex > 0
|
||||||
? [stepIndex - 1]
|
? [stepIndex - 1]
|
||||||
: [];
|
: [];
|
||||||
@@ -164,12 +189,12 @@ export async function updateStepImpl(store: TaskStore, id: string, stepIndex: nu
|
|||||||
if (blockingIndex !== -1) {
|
if (blockingIndex !== -1) {
|
||||||
const ts = new Date().toISOString();
|
const ts = new Date().toISOString();
|
||||||
task.updatedAt = ts;
|
task.updatedAt = ts;
|
||||||
const kind = graphSource ? "dependency-order" : "out-of-order";
|
const kind = dependencyOrdered ? "dependency-order" : "out-of-order";
|
||||||
task.log.push({
|
task.log.push({
|
||||||
timestamp: ts,
|
timestamp: ts,
|
||||||
action:
|
action:
|
||||||
`Ignored ${kind} ${status} for step ${stepIndex} (${task.steps[stepIndex].name}) — ` +
|
`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
|
// Graph-source suppression is a projection bug — surface it loudly in
|
||||||
// the activity log (U6) rather than the legacy silent ignore.
|
// 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.steps[stepIndex].status = status;
|
||||||
task.updatedAt = new Date().toISOString();
|
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") {
|
if (status === "done") {
|
||||||
while (
|
const firstUnfinished = task.steps.findIndex(
|
||||||
task.currentStep < task.steps.length &&
|
(step) => step.status !== "done" && step.status !== "skipped",
|
||||||
(task.steps[task.currentStep].status === "done" || task.steps[task.currentStep].status === "skipped")
|
);
|
||||||
) {
|
task.currentStep = firstUnfinished === -1 ? task.steps.length : firstUnfinished;
|
||||||
task.currentStep++;
|
|
||||||
}
|
|
||||||
} else if (status === "in-progress") {
|
} else if (status === "in-progress") {
|
||||||
task.currentStep = stepIndex;
|
task.currentStep = stepIndex;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user