FN-6294: collapse workflow merge primitives
Collapse built-in workflow merge-policy primitives back to the legacy merge seam until graph execution owns merge policy end-to-end. - Treat all merge-policy primitive node kinds as a synthetic legacy merge node during graph traversal. - Preserve shadow-stage parity by mapping the synthetic merge node into legacy stage transitions. - Add regression coverage for success, failure, pre-merge review failure, and alternate merge-region entry points. - Record a patch changeset for the published Fusion package. Files changed: .changeset/fn-6294-merge-region-collapse.md | 5 + .../workflow-graph-executor-handlers.test.ts | 13 ++- .../workflow-graph-merge-region-collapse.test.ts | 116 +++++++++++++++++++++ packages/engine/src/executor.ts | 14 ++- packages/engine/src/workflow-graph-executor.ts | 63 ++++++++++- 5 files changed, 203 insertions(+), 8 deletions(-) Fusion-Task-Id: FN-6294 Fusion-Task-Lineage: fe86b5e2-861f-448d-8db0-b0332e6ad76f
This commit is contained in:
5
.changeset/fn-6294-merge-region-collapse.md
Normal file
5
.changeset/fn-6294-merge-region-collapse.md
Normal file
@@ -0,0 +1,5 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
Fix workflow graph execution for the built-in coding workflow's merge-policy primitive region by collapsing any merge-region entry back to the legacy `merge` seam until the workflow interpreter owns merge policy execution.
|
||||||
@@ -391,9 +391,18 @@ describe("WorkflowGraphExecutor traversal", () => {
|
|||||||
expect(order).toEqual(["b", "c"]);
|
expect(order).toEqual(["b", "c"]);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("builtin coding workflow ir exposes expected lifecycle nodes", () => {
|
it("builtin coding workflow ir exposes expected lifecycle and merge-policy nodes", () => {
|
||||||
expect(BUILTIN_CODING_WORKFLOW_IR.nodes.map((node) => node.id)).toEqual(
|
expect(BUILTIN_CODING_WORKFLOW_IR.nodes.map((node) => node.id)).toEqual(
|
||||||
expect.arrayContaining(["start", "execute", "review", "merge", "end"]),
|
expect.arrayContaining([
|
||||||
|
"start",
|
||||||
|
"execute",
|
||||||
|
"review",
|
||||||
|
"merge-gate",
|
||||||
|
"branch-group-member-integration",
|
||||||
|
"branch-group-promotion",
|
||||||
|
"merge-attempt",
|
||||||
|
"end",
|
||||||
|
]),
|
||||||
);
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,116 @@
|
|||||||
|
import { describe, expect, it, vi } from "vitest";
|
||||||
|
import type { TaskDetail, WorkflowIr, WorkflowIrNodeKind } from "@fusion/core";
|
||||||
|
import { BUILTIN_CODING_WORKFLOW_IR } from "@fusion/core";
|
||||||
|
|
||||||
|
import { WorkflowGraphExecutor } from "../workflow-graph-executor.js";
|
||||||
|
import type { WorkflowLegacySeams } from "../workflow-node-handlers.js";
|
||||||
|
|
||||||
|
const task = { id: "FN-6294" } as TaskDetail;
|
||||||
|
const settings = { experimentalFeatures: { workflowGraphExecutor: true } };
|
||||||
|
|
||||||
|
const mergeRegionEntries: Array<{ id: string; kind: WorkflowIrNodeKind }> = [
|
||||||
|
{ id: "merge-gate", kind: "merge-gate" },
|
||||||
|
{ id: "merge-attempt", kind: "merge-attempt" },
|
||||||
|
{ id: "merge-manual-hold", kind: "manual-merge-hold" },
|
||||||
|
{ id: "merge-retry", kind: "retry-backoff" },
|
||||||
|
{ id: "recovery-router", kind: "recovery-router" },
|
||||||
|
{ id: "branch-group-member-integration", kind: "branch-group-member-integration" },
|
||||||
|
{ id: "branch-group-promotion", kind: "branch-group-promotion" },
|
||||||
|
];
|
||||||
|
const rawMergeRegionNodeIds = mergeRegionEntries.map((entry) => entry.id);
|
||||||
|
|
||||||
|
function createSeams(overrides: Partial<WorkflowLegacySeams> = {}): WorkflowLegacySeams {
|
||||||
|
return {
|
||||||
|
planning: async () => ({ outcome: "success" }),
|
||||||
|
execute: async () => ({ outcome: "success" }),
|
||||||
|
workflowStep: async () => ({ outcome: "success" }),
|
||||||
|
review: async () => ({ outcome: "success" }),
|
||||||
|
merge: async () => ({ outcome: "success" }),
|
||||||
|
schedule: async () => ({ outcome: "success" }),
|
||||||
|
...overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function expectNoRawMergeRegionVisits(visitedNodeIds: string[]) {
|
||||||
|
for (const rawNodeId of rawMergeRegionNodeIds) {
|
||||||
|
expect(visitedNodeIds).not.toContain(rawNodeId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function irEnteringMergeRegionAt(entryId: string): WorkflowIr {
|
||||||
|
return {
|
||||||
|
...BUILTIN_CODING_WORKFLOW_IR,
|
||||||
|
edges: BUILTIN_CODING_WORKFLOW_IR.edges.map((edge) =>
|
||||||
|
edge.from === "review" && edge.to === "merge-gate" && edge.condition === "success"
|
||||||
|
? { ...edge, to: entryId }
|
||||||
|
: edge,
|
||||||
|
),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("WorkflowGraphExecutor merge-region collapse", () => {
|
||||||
|
it("collapses the built-in merge-policy region to one legacy merge seam", async () => {
|
||||||
|
const calls: string[] = [];
|
||||||
|
const merge = vi.fn(async () => {
|
||||||
|
calls.push("merge");
|
||||||
|
return { outcome: "success" as const };
|
||||||
|
});
|
||||||
|
const executor = new WorkflowGraphExecutor({ seams: createSeams({ merge }) });
|
||||||
|
|
||||||
|
const result = await executor.run(task, settings, BUILTIN_CODING_WORKFLOW_IR);
|
||||||
|
|
||||||
|
expect(result.outcome).toBe("success");
|
||||||
|
expect(merge).toHaveBeenCalledOnce();
|
||||||
|
expect(calls).toEqual(["merge"]);
|
||||||
|
expect(result.visitedNodeIds).toEqual(["start", "planning", "execute", "workflow-step", "review", "merge"]);
|
||||||
|
expect(result.context["node:merge:outcome"]).toBe("success");
|
||||||
|
expectNoRawMergeRegionVisits(result.visitedNodeIds);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("routes legacy merge seam failures to a failure terminal without visiting raw merge primitives", async () => {
|
||||||
|
const merge = vi.fn(async () => ({ outcome: "failure" as const, value: "FileScopeViolationError" }));
|
||||||
|
const executor = new WorkflowGraphExecutor({ seams: createSeams({ merge }) });
|
||||||
|
|
||||||
|
const result = await executor.run(task, settings, BUILTIN_CODING_WORKFLOW_IR);
|
||||||
|
|
||||||
|
expect(result.outcome).toBe("failure");
|
||||||
|
expect(merge).toHaveBeenCalledOnce();
|
||||||
|
expect(result.visitedNodeIds).toEqual(["start", "planning", "execute", "workflow-step", "review", "merge"]);
|
||||||
|
expect(result.context["node:merge:outcome"]).toBe("failure");
|
||||||
|
expect(result.context["node:merge:value"]).toBe("FileScopeViolationError");
|
||||||
|
expectNoRawMergeRegionVisits(result.visitedNodeIds);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not collapse to merge when review fails before the merge-policy region", async () => {
|
||||||
|
const merge = vi.fn(async () => ({ outcome: "success" as const }));
|
||||||
|
const executor = new WorkflowGraphExecutor({
|
||||||
|
seams: createSeams({
|
||||||
|
review: async () => ({ outcome: "failure", value: "manual-merge-required" }),
|
||||||
|
merge,
|
||||||
|
}),
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await executor.run(task, settings, BUILTIN_CODING_WORKFLOW_IR);
|
||||||
|
|
||||||
|
expect(result.outcome).toBe("failure");
|
||||||
|
expect(merge).not.toHaveBeenCalled();
|
||||||
|
expect(result.visitedNodeIds).toEqual(["start", "planning", "execute", "workflow-step", "review"]);
|
||||||
|
expect(result.visitedNodeIds).not.toContain("merge");
|
||||||
|
expectNoRawMergeRegionVisits(result.visitedNodeIds);
|
||||||
|
});
|
||||||
|
|
||||||
|
it.each(mergeRegionEntries)(
|
||||||
|
"treats $kind as a merge-region boundary when entered directly",
|
||||||
|
async ({ id }) => {
|
||||||
|
const merge = vi.fn(async () => ({ outcome: "success" as const }));
|
||||||
|
const executor = new WorkflowGraphExecutor({ seams: createSeams({ merge }) });
|
||||||
|
|
||||||
|
const result = await executor.run(task, settings, irEnteringMergeRegionAt(id));
|
||||||
|
|
||||||
|
expect(result.outcome).toBe("success");
|
||||||
|
expect(merge).toHaveBeenCalledOnce();
|
||||||
|
expect(result.visitedNodeIds).toEqual(["start", "planning", "execute", "workflow-step", "review", "merge"]);
|
||||||
|
expectNoRawMergeRegionVisits(result.visitedNodeIds);
|
||||||
|
},
|
||||||
|
);
|
||||||
|
});
|
||||||
@@ -4457,11 +4457,15 @@ export class TaskExecutor {
|
|||||||
stageByNodeId.set(node.id, seam);
|
stageByNodeId.set(node.id, seam);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Stop the shadow walk at the live terminal seam. The graph walker visits a
|
// The built-in coding IR now enters an interpreter-owned merge-policy
|
||||||
// node *before* invoking its seam, so even a failing merge seam (the case
|
// primitive region after review; graph execution collapses that region to a
|
||||||
// when the live task is parked in-review with autoMerge off) still records a
|
// synthetic legacy merge seam recorded as `merge` until merge-policy cutover.
|
||||||
// "merge" stage. The legacy side never reports merge for an in-review task,
|
stageByNodeId.set("merge", "merge");
|
||||||
// so that phantom stage manufactures stageTransitions drift on healthy runs.
|
// Stop the shadow walk at the live terminal seam. The graph walker records a
|
||||||
|
// merge stage before/while invoking its seam, so even a failing merge seam
|
||||||
|
// (the case when the live task is parked in-review with autoMerge off) still
|
||||||
|
// records a "merge" stage. The legacy side never reports merge for an
|
||||||
|
// in-review task, so that phantom stage manufactures stageTransitions drift.
|
||||||
// Truncate the visited-stage sequence at the stage the live task actually
|
// Truncate the visited-stage sequence at the stage the live task actually
|
||||||
// reached: merged → merge, reachedReview → review, else → execute.
|
// reached: merged → merge, reachedReview → review, else → execute.
|
||||||
const terminalStage: WorkflowStage = merged ? "merge" : reachedReview ? "review" : "execute";
|
const terminalStage: WorkflowStage = merged ? "merge" : reachedReview ? "review" : "execute";
|
||||||
|
|||||||
@@ -1,4 +1,13 @@
|
|||||||
import type { Settings, TaskDetail, TaskStep, WorkflowIr, WorkflowIrEdge, WorkflowIrNode, WorkflowNodeExtensionResult } from "@fusion/core";
|
import type {
|
||||||
|
Settings,
|
||||||
|
TaskDetail,
|
||||||
|
TaskStep,
|
||||||
|
WorkflowIr,
|
||||||
|
WorkflowIrEdge,
|
||||||
|
WorkflowIrNode,
|
||||||
|
WorkflowIrNodeKind,
|
||||||
|
WorkflowNodeExtensionResult,
|
||||||
|
} from "@fusion/core";
|
||||||
import { BUILTIN_CODING_WORKFLOW_IR, WorkflowIrError, getWorkflowExtensionRegistry, isExperimentalFeatureEnabled, resolveMaxReworkCycles } from "@fusion/core";
|
import { BUILTIN_CODING_WORKFLOW_IR, WorkflowIrError, getWorkflowExtensionRegistry, isExperimentalFeatureEnabled, resolveMaxReworkCycles } from "@fusion/core";
|
||||||
|
|
||||||
import {
|
import {
|
||||||
@@ -164,6 +173,27 @@ const TERMINAL_FAILURE: WorkflowGraphExecutorResult = {
|
|||||||
visitedNodeIds: [],
|
visitedNodeIds: [],
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Engine-local mirror of core's workflow-owned merge/retry/recovery primitive
|
||||||
|
* region. Until the workflow interpreter owns merge policy end-to-end, graph
|
||||||
|
* execution treats any entry into this region as the terminal legacy `merge`
|
||||||
|
* seam so observable lifecycle behavior stays byte-identical with the legacy
|
||||||
|
* executor. Consolidate with a core export when one exists.
|
||||||
|
*/
|
||||||
|
const MERGE_REGION_KINDS = new Set<WorkflowIrNodeKind>([
|
||||||
|
"merge-gate",
|
||||||
|
"merge-attempt",
|
||||||
|
"manual-merge-hold",
|
||||||
|
"retry-backoff",
|
||||||
|
"recovery-router",
|
||||||
|
"branch-group-member-integration",
|
||||||
|
"branch-group-promotion",
|
||||||
|
]);
|
||||||
|
|
||||||
|
function isMergeRegionKind(kind: WorkflowIrNodeKind): boolean {
|
||||||
|
return MERGE_REGION_KINDS.has(kind);
|
||||||
|
}
|
||||||
|
|
||||||
function normalizeTouchedFile(value: unknown): string | undefined {
|
function normalizeTouchedFile(value: unknown): string | undefined {
|
||||||
if (typeof value === "string") {
|
if (typeof value === "string") {
|
||||||
const trimmed = value.trim().replaceAll("\\", "/").replace(/^\.\//, "");
|
const trimmed = value.trim().replaceAll("\\", "/").replace(/^\.\//, "");
|
||||||
@@ -273,6 +303,12 @@ export class WorkflowGraphExecutor {
|
|||||||
};
|
};
|
||||||
const visitedNodeIds: string[] = [];
|
const visitedNodeIds: string[] = [];
|
||||||
const inStack = new Set<string>();
|
const inStack = new Set<string>();
|
||||||
|
const syntheticMergeNode: WorkflowIrNode = {
|
||||||
|
id: "merge",
|
||||||
|
kind: "prompt",
|
||||||
|
column: "in-review",
|
||||||
|
config: { seam: "merge" },
|
||||||
|
};
|
||||||
|
|
||||||
// Bounded-rework generalization (U6). A `kind: "rework"` edge is the only
|
// Bounded-rework generalization (U6). A `kind: "rework"` edge is the only
|
||||||
// legal cycle: it loops back to a "rework region head" (the edge's `to` node).
|
// legal cycle: it loops back to a "rework region head" (the edge's `to` node).
|
||||||
@@ -499,6 +535,26 @@ export class WorkflowGraphExecutor {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
const runLegacyMergeSeam = async (): Promise<WorkflowNodeResult> => {
|
||||||
|
// The merge-policy primitive region is interpreter-owned policy. While the
|
||||||
|
// legacy lifecycle remains authoritative, reaching any of its node kinds is
|
||||||
|
// the terminal merge boundary: dispatch the same prompt/seam handler a
|
||||||
|
// legacy `config.seam: "merge"` node used, but record it under the stable
|
||||||
|
// legacy node id `merge` and never expose raw merge-region primitive ids.
|
||||||
|
visitedNodeIds.push(syntheticMergeNode.id);
|
||||||
|
const result = await this.executeNodeWithRetries(
|
||||||
|
syntheticMergeNode,
|
||||||
|
task,
|
||||||
|
settings,
|
||||||
|
context,
|
||||||
|
ir,
|
||||||
|
);
|
||||||
|
if (result.contextPatch) Object.assign(context, result.contextPatch);
|
||||||
|
context[`node:${syntheticMergeNode.id}:outcome`] = result.outcome;
|
||||||
|
if (result.value !== undefined) context[`node:${syntheticMergeNode.id}:value`] = result.value;
|
||||||
|
return result;
|
||||||
|
};
|
||||||
|
|
||||||
const traverseChildren = async (
|
const traverseChildren = async (
|
||||||
node: WorkflowIrNode,
|
node: WorkflowIrNode,
|
||||||
sourceResult: WorkflowNodeResult,
|
sourceResult: WorkflowNodeResult,
|
||||||
@@ -532,6 +588,11 @@ export class WorkflowGraphExecutor {
|
|||||||
aggregate = sourceResult;
|
aggregate = sourceResult;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
if (target && isMergeRegionKind(target.kind)) {
|
||||||
|
aggregate = await runLegacyMergeSeam();
|
||||||
|
if (aggregate.outcome === "failure") break;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
const child = await walk(edge.to);
|
const child = await walk(edge.to);
|
||||||
// A ReworkSignal propagated from deeper: bubble it further up unchanged.
|
// A ReworkSignal propagated from deeper: bubble it further up unchanged.
|
||||||
if (isReworkSignal(child)) return child;
|
if (isReworkSignal(child)) return child;
|
||||||
|
|||||||
Reference in New Issue
Block a user