fix(FN-7228): respect workflow step dependencies
This commit is contained in:
7
.changeset/fn-7228-step-session-dependencies.md
Normal file
7
.changeset/fn-7228-step-session-dependencies.md
Normal file
@@ -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.
|
||||||
@@ -57,13 +57,13 @@ describe("TaskStore.updateStep step-order guard", () => {
|
|||||||
|
|
||||||
// ── U6: graph-source projection discipline (KTD-7/KTD-11) ──────────────────
|
// ── 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 () => {
|
it("graph source: done is legal for explicitly 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.
|
|
||||||
const store = harness.store();
|
const store = harness.store();
|
||||||
const task = await harness.createTaskWithSteps();
|
const task = await harness.createTaskWithSteps();
|
||||||
await store.updateStep(task.id, 0, "pending");
|
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" });
|
await store.updateStep(task.id, 1, "in-progress", { source: "graph" });
|
||||||
const updated = await store.updateStep(task.id, 2, "done", { 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);
|
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 () => {
|
it("graph source: explicit dependsOn still suppresses completion until dependencies finish", async () => {
|
||||||
const store = harness.store();
|
const store = harness.store();
|
||||||
const task = await harness.createTaskWithSteps();
|
const task = await harness.createTaskWithSteps();
|
||||||
|
|||||||
@@ -9241,9 +9241,8 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS}
|
|||||||
// (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 the legacy (default) write:
|
||||||
// (a) the out-of-order-done guard relaxes from strict index order to
|
// (a) the out-of-order-done guard relaxes from strict index order to
|
||||||
// DEPENDENCY order (a done write is legal when every explicitly declared
|
// DEPENDENCY order (a done write is legal when every dependsOn step —
|
||||||
// dependsOn step is done/skipped; absent dependsOn means the graph has
|
// default: the immediately-preceding step — is done/skipped, KTD-11);
|
||||||
// already decided the step can run independently);
|
|
||||||
// (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
|
||||||
@@ -9294,20 +9293,22 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS}
|
|||||||
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:
|
// may go done. Legacy: strict index order (every earlier step). Graph:
|
||||||
// only the step's explicit dependsOn list. When a graph-owned step session
|
// the step's dependsOn list, with absent dependsOn defaulting to the
|
||||||
// finishes before an earlier independent step, the task card must record
|
// immediately-preceding step. A deliberately empty dependsOn array is the
|
||||||
// that real completion instead of emitting a false "out-of-order" pause.
|
// opt-in for an independent graph step.
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowStepControl 2026-06-29-10:12:
|
FNXC:WorkflowStepControl 2026-06-29-10:51:
|
||||||
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.
|
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 blockingIndex = -1;
|
||||||
let blockingStatus: import("./types.js").StepStatus | undefined;
|
let blockingStatus: import("./types.js").StepStatus | undefined;
|
||||||
if (graphSource) {
|
if (graphSource) {
|
||||||
const deps = task.steps[stepIndex]?.dependsOn;
|
const deps = task.steps[stepIndex]?.dependsOn;
|
||||||
const depIndices =
|
const depIndices =
|
||||||
Array.isArray(deps) && deps.length > 0
|
Array.isArray(deps)
|
||||||
? deps
|
? deps
|
||||||
|
: stepIndex > 0
|
||||||
|
? [stepIndex - 1]
|
||||||
: [];
|
: [];
|
||||||
for (const i of depIndices) {
|
for (const i of depIndices) {
|
||||||
const priorStatus = task.steps[i]?.status;
|
const priorStatus = task.steps[i]?.status;
|
||||||
|
|||||||
@@ -60,6 +60,14 @@ function makeTaskDetail(overrides: Partial<TaskDetail> = {}): TaskDetail {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function makeIndependentSteps(count: number): TaskDetail["steps"] {
|
||||||
|
return Array.from({ length: count }, (_, i) => ({
|
||||||
|
name: `Step ${i}`,
|
||||||
|
status: "pending" as const,
|
||||||
|
dependsOn: [],
|
||||||
|
}));
|
||||||
|
}
|
||||||
|
|
||||||
// ── parseStepFileScopes tests ──────────────────────────────────────────
|
// ── parseStepFileScopes tests ──────────────────────────────────────────
|
||||||
|
|
||||||
describe("parseStepFileScopes", () => {
|
describe("parseStepFileScopes", () => {
|
||||||
@@ -289,40 +297,61 @@ describe("buildConflictMatrix", () => {
|
|||||||
// ── determineParallelWaves tests ───────────────────────────────────────
|
// ── determineParallelWaves tests ───────────────────────────────────────
|
||||||
|
|
||||||
describe("determineParallelWaves", () => {
|
describe("determineParallelWaves", () => {
|
||||||
|
const independentSteps = (count: number) =>
|
||||||
|
Array.from({ length: count }, () => ({ dependsOn: [] as number[] }));
|
||||||
|
|
||||||
it("maxParallel=1 → all steps sequential (one per wave)", () => {
|
it("maxParallel=1 → all steps sequential (one per wave)", () => {
|
||||||
const scopes = new Map<number, string[]>([
|
const scopes = new Map<number, string[]>([
|
||||||
[0, ["src/a.ts"]],
|
[0, ["src/a.ts"]],
|
||||||
[1, ["src/b.ts"]],
|
[1, ["src/b.ts"]],
|
||||||
[2, ["src/c.ts"]],
|
[2, ["src/c.ts"]],
|
||||||
]);
|
]);
|
||||||
const waves = determineParallelWaves(scopes, 1);
|
const waves = determineParallelWaves(scopes, 1, independentSteps(3));
|
||||||
expect(waves).toHaveLength(3);
|
expect(waves).toHaveLength(3);
|
||||||
expect(waves[0]).toEqual({ indices: [0], waveNumber: 0 });
|
expect(waves[0]).toEqual({ indices: [0], waveNumber: 0 });
|
||||||
expect(waves[1]).toEqual({ indices: [1], waveNumber: 1 });
|
expect(waves[1]).toEqual({ indices: [1], waveNumber: 1 });
|
||||||
expect(waves[2]).toEqual({ indices: [2], waveNumber: 2 });
|
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<number, string[]>([
|
const scopes = new Map<number, string[]>([
|
||||||
[0, ["src/a.ts"]],
|
[0, ["src/a.ts"]],
|
||||||
[1, ["src/b.ts"]],
|
[1, ["src/b.ts"]],
|
||||||
[2, ["src/c.ts"]],
|
[2, ["src/c.ts"]],
|
||||||
[3, ["src/d.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
|
// All non-conflicting, capped at 2 per wave → 2 waves
|
||||||
expect(waves).toHaveLength(2);
|
expect(waves).toHaveLength(2);
|
||||||
expect(waves[0].indices).toEqual([0, 1]);
|
expect(waves[0].indices).toEqual([0, 1]);
|
||||||
expect(waves[1].indices).toEqual([2, 3]);
|
expect(waves[1].indices).toEqual([2, 3]);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("unannotated steps are sequential so verification cannot start before implementation", () => {
|
||||||
|
const scopes = new Map<number, string[]>([
|
||||||
|
[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", () => {
|
it("all steps conflict → each step in its own wave", () => {
|
||||||
const scopes = new Map<number, string[]>([
|
const scopes = new Map<number, string[]>([
|
||||||
[0, ["src/a.ts"]],
|
[0, ["src/a.ts"]],
|
||||||
[1, ["src/a.ts"]],
|
[1, ["src/a.ts"]],
|
||||||
[2, ["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
|
// All conflict with each other → each in own wave
|
||||||
expect(waves).toHaveLength(3);
|
expect(waves).toHaveLength(3);
|
||||||
expect(waves[0].indices).toEqual([0]);
|
expect(waves[0].indices).toEqual([0]);
|
||||||
@@ -338,7 +367,7 @@ describe("determineParallelWaves", () => {
|
|||||||
[2, ["packages/engine"]],
|
[2, ["packages/engine"]],
|
||||||
[3, ["packages/engine/src/executor.ts"]],
|
[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 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]
|
// 0 and 2 don't conflict → wave 0: [0, 2]
|
||||||
// 1 and 3 don't conflict → wave 1: [1, 3]
|
// 1 and 3 don't conflict → wave 1: [1, 3]
|
||||||
@@ -352,7 +381,7 @@ describe("determineParallelWaves", () => {
|
|||||||
[0, []],
|
[0, []],
|
||||||
[1, []],
|
[1, []],
|
||||||
]);
|
]);
|
||||||
const waves = determineParallelWaves(scopes, 4);
|
const waves = determineParallelWaves(scopes, 4, independentSteps(2));
|
||||||
expect(waves).toHaveLength(1);
|
expect(waves).toHaveLength(1);
|
||||||
expect(waves[0].indices).toEqual([0, 1]);
|
expect(waves[0].indices).toEqual([0, 1]);
|
||||||
});
|
});
|
||||||
@@ -364,7 +393,7 @@ describe("determineParallelWaves", () => {
|
|||||||
[2, ["c.ts"]],
|
[2, ["c.ts"]],
|
||||||
[3, ["d.ts"]],
|
[3, ["d.ts"]],
|
||||||
]);
|
]);
|
||||||
const waves = determineParallelWaves(scopes, 2);
|
const waves = determineParallelWaves(scopes, 2, independentSteps(4));
|
||||||
expect(waves).toHaveLength(2);
|
expect(waves).toHaveLength(2);
|
||||||
expect(waves[0].indices).toHaveLength(2);
|
expect(waves[0].indices).toHaveLength(2);
|
||||||
expect(waves[1].indices).toHaveLength(2);
|
expect(waves[1].indices).toHaveLength(2);
|
||||||
@@ -375,7 +404,7 @@ describe("determineParallelWaves", () => {
|
|||||||
[0, ["a.ts"]],
|
[0, ["a.ts"]],
|
||||||
[1, ["b.ts"]],
|
[1, ["b.ts"]],
|
||||||
]);
|
]);
|
||||||
const waves = determineParallelWaves(scopes, 1);
|
const waves = determineParallelWaves(scopes, 1, independentSteps(2));
|
||||||
for (let i = 0; i < waves.length; i++) {
|
for (let i = 0; i < waves.length; i++) {
|
||||||
expect(waves[i].waveNumber).toBe(i);
|
expect(waves[i].waveNumber).toBe(i);
|
||||||
}
|
}
|
||||||
@@ -387,7 +416,7 @@ describe("determineParallelWaves", () => {
|
|||||||
[1, ["b.ts"]],
|
[1, ["b.ts"]],
|
||||||
[2, ["c.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]
|
// 3 non-conflicting steps, max 2 → wave 0: [0,1], wave 1: [2]
|
||||||
expect(waves).toHaveLength(2);
|
expect(waves).toHaveLength(2);
|
||||||
expect(waves[0].indices).toEqual([0, 1]);
|
expect(waves[0].indices).toEqual([0, 1]);
|
||||||
@@ -1578,10 +1607,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -1623,10 +1649,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -1680,10 +1703,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -1732,10 +1752,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -1825,10 +1842,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
const semaphore = new AgentSemaphore(4);
|
const semaphore = new AgentSemaphore(4);
|
||||||
@@ -1868,10 +1882,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -1956,11 +1967,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(3),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
{ name: "Step 2", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 3 });
|
const settings = makeSettings({ maxParallelSteps: 3 });
|
||||||
|
|
||||||
@@ -2019,11 +2026,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(3),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
{ name: "Step 2", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 3 });
|
const settings = makeSettings({ maxParallelSteps: 3 });
|
||||||
|
|
||||||
@@ -2129,10 +2132,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -2181,11 +2181,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(3),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
{ name: "Step 2", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 3 });
|
const settings = makeSettings({ maxParallelSteps: 3 });
|
||||||
|
|
||||||
@@ -2331,10 +2327,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
@@ -2398,10 +2391,7 @@ describe("StepSessionExecutor", () => {
|
|||||||
|
|
||||||
const task = makeTaskDetail({
|
const task = makeTaskDetail({
|
||||||
prompt,
|
prompt,
|
||||||
steps: [
|
steps: makeIndependentSteps(2),
|
||||||
{ name: "Step 0", status: "pending" },
|
|
||||||
{ name: "Step 1", status: "pending" },
|
|
||||||
],
|
|
||||||
});
|
});
|
||||||
const settings = makeSettings({ maxParallelSteps: 2 });
|
const settings = makeSettings({ maxParallelSteps: 2 });
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
* Uses a greedy left-to-right scan: for each wave, add steps whose dependencies
|
||||||
* conflict with any step already in the wave. Each wave's `indices` array
|
* have all been assigned to earlier waves and that don't conflict with any step
|
||||||
* is capped at `maxParallel`. Remaining steps spill into subsequent waves.
|
* 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 stepScopes - Map from step index to array of file paths.
|
||||||
* @param maxParallel - Maximum number of steps per wave (range 1–4).
|
* @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(
|
export function determineParallelWaves(
|
||||||
stepScopes: Map<number, string[]>,
|
stepScopes: Map<number, string[]>,
|
||||||
maxParallel: number,
|
maxParallel: number,
|
||||||
|
steps: Array<{ dependsOn?: number[] }> = [],
|
||||||
): ParallelWave[] {
|
): ParallelWave[] {
|
||||||
const indices = [...stepScopes.keys()].sort((a, b) => a - b);
|
const indices = [...stepScopes.keys()].sort((a, b) => a - b);
|
||||||
if (indices.length === 0) return [];
|
if (indices.length === 0) return [];
|
||||||
@@ -326,6 +328,8 @@ export function determineParallelWaves(
|
|||||||
for (const idx of indices) {
|
for (const idx of indices) {
|
||||||
if (assigned.has(idx)) continue;
|
if (assigned.has(idx)) continue;
|
||||||
if (waveIndices.length >= clampedMax) break;
|
if (waveIndices.length >= clampedMax) break;
|
||||||
|
const deps = resolveStepSessionDependsOn(steps, idx);
|
||||||
|
if (!deps.every((dep) => assigned.has(dep))) continue;
|
||||||
|
|
||||||
const pos = posMap.get(idx)!;
|
const pos = posMap.get(idx)!;
|
||||||
|
|
||||||
@@ -341,8 +345,9 @@ export function determineParallelWaves(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (waveIndices.length === 0) {
|
if (waveIndices.length === 0) {
|
||||||
// Safety: if no steps could be added (shouldn't happen with diagonal=true),
|
// Safety: malformed/cyclic dependencies should not infinite-loop the
|
||||||
// force-add the next unassigned step
|
// executor. Force the next unassigned step and let downstream execution/
|
||||||
|
// validation report the real problem.
|
||||||
const next = indices.find((idx) => !assigned.has(idx));
|
const next = indices.find((idx) => !assigned.has(idx));
|
||||||
if (next !== undefined) {
|
if (next !== undefined) {
|
||||||
waveIndices.push(next);
|
waveIndices.push(next);
|
||||||
@@ -362,6 +367,12 @@ export function determineParallelWaves(
|
|||||||
return waves;
|
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 ───────────────────────────────────────────────
|
// ── 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(
|
stepExecLog.log(
|
||||||
`Executing ${stepCount} steps in ${waves.length} wave(s) for task ${taskDetail.id} ` +
|
`Executing ${stepCount} steps in ${waves.length} wave(s) for task ${taskDetail.id} ` +
|
||||||
|
|||||||
Reference in New Issue
Block a user