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"]);
|
||||
});
|
||||
|
||||
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.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);
|
||||
}
|
||||
}
|
||||
// Stop the shadow walk at the live terminal seam. The graph walker visits a
|
||||
// node *before* 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 on healthy runs.
|
||||
// The built-in coding IR now enters an interpreter-owned merge-policy
|
||||
// primitive region after review; graph execution collapses that region to a
|
||||
// synthetic legacy merge seam recorded as `merge` until merge-policy cutover.
|
||||
stageByNodeId.set("merge", "merge");
|
||||
// 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
|
||||
// reached: merged → merge, reachedReview → review, else → 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 {
|
||||
@@ -164,6 +173,27 @@ const TERMINAL_FAILURE: WorkflowGraphExecutorResult = {
|
||||
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 {
|
||||
if (typeof value === "string") {
|
||||
const trimmed = value.trim().replaceAll("\\", "/").replace(/^\.\//, "");
|
||||
@@ -273,6 +303,12 @@ export class WorkflowGraphExecutor {
|
||||
};
|
||||
const visitedNodeIds: 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
|
||||
// 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 (
|
||||
node: WorkflowIrNode,
|
||||
sourceResult: WorkflowNodeResult,
|
||||
@@ -532,6 +588,11 @@ export class WorkflowGraphExecutor {
|
||||
aggregate = sourceResult;
|
||||
continue;
|
||||
}
|
||||
if (target && isMergeRegionKind(target.kind)) {
|
||||
aggregate = await runLegacyMergeSeam();
|
||||
if (aggregate.outcome === "failure") break;
|
||||
continue;
|
||||
}
|
||||
const child = await walk(edge.to);
|
||||
// A ReworkSignal propagated from deeper: bubble it further up unchanged.
|
||||
if (isReworkSignal(child)) return child;
|
||||
|
||||
Reference in New Issue
Block a user