FN-8601: enforce foreach merge proof
Require complete foreach execution evidence before workflow merge review. - Add reusable foreach instance coverage proof evaluation. - Block checklist projection and merge admission on incomplete or failed node results. - Cover core proof logic and PostgreSQL merge-boundary behavior. - Add a patch changeset for the merge safeguard. Files changed: .changeset/fn-8601-foreach-merge-proof.md | 7 ++ .../src/__tests__/workflow-merge-proof.test.ts | 43 ++++++++ packages/core/src/index.gate.ts | 2 + packages/core/src/index.ts | 2 + packages/core/src/workflow-merge-proof.ts | 74 +++++++++++++ ...xecutor-merge-boundary-foreach-proof.pg.test.ts | 111 +++++++++++++++++++ packages/engine/src/executor.ts | 117 +++++++++++++-------- 7 files changed, 314 insertions(+), 42 deletions(-) Fusion-Task-Id: FN-8601 Fusion-Task-Lineage: 40578171-0b13-4538-8f38-3948ed1e92c0 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8601-foreach-merge-proof.md
Normal file
7
.changeset/fn-8601-foreach-merge-proof.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Block incomplete foreach workflow steps from entering merge review.
|
||||
category: fix
|
||||
dev: Merge proof now correlates foreach step-execute results with every expanded instance.
|
||||
43
packages/core/src/__tests__/workflow-merge-proof.test.ts
Normal file
43
packages/core/src/__tests__/workflow-merge-proof.test.ts
Normal file
@@ -0,0 +1,43 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { instanceNodeId } from "../column-agent-resolver.js";
|
||||
import { evaluateForeachMergeProof } from "../workflow-merge-proof.js";
|
||||
import type { TaskStep, WorkflowIr, WorkflowStepResult } from "../types.js";
|
||||
|
||||
const steps = (statuses: Array<TaskStep["status"]>): TaskStep[] => statuses.map((status, index) => ({ id: `s${index}`, title: `Step ${index}`, status }));
|
||||
const result = (id: string, status: WorkflowStepResult["status"] = "passed"): WorkflowStepResult => ({ workflowStepId: id, source: "node", phase: "pre-merge", status, completedAt: new Date().toISOString() });
|
||||
const foreachIr = (seam = true): WorkflowIr => ({ version: 2, columns: [], nodes: [{ id: "steps", kind: "foreach", config: { source: "task-steps", template: { nodes: seam ? [{ id: "step-execute", kind: "prompt", config: { seam: "step-execute" } }] : [{ id: "review", kind: "prompt" }], edges: [] } } }], edges: [] });
|
||||
const plainIr: WorkflowIr = { version: 2, columns: [], nodes: [{ id: "execute", kind: "prompt", config: { seam: "execute" } }], edges: [] };
|
||||
|
||||
describe("evaluateForeachMergeProof", () => {
|
||||
it("is vacuously complete without a foreach seam, including empty results", () => {
|
||||
for (const ir of [plainIr, foreachIr(false)]) {
|
||||
const proof = evaluateForeachMergeProof({ ir, steps: steps(["pending"]), workflowStepResults: [] });
|
||||
expect(proof).toMatchObject({ hasForeachStepExecute: false, expectedInstanceIds: [], missingInstanceIds: [] });
|
||||
}
|
||||
});
|
||||
|
||||
it("covers complete, partial, zero-step, already-terminal, failed, post-merge, and duplicate result shapes", () => {
|
||||
const ids = [0, 1, 2].map((index) => instanceNodeId("steps", index, "step-execute"));
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending", "pending", "pending"]), workflowStepResults: ids.map((id) => result(id)) })).toMatchObject({ missingInstanceIds: [] });
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending", "pending", "pending"]), workflowStepResults: [result(ids[0])] })).toMatchObject({ missingInstanceIds: [ids[1], ids[2]] });
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: [], workflowStepResults: [] })).toMatchObject({ hasForeachStepExecute: true, missingInstanceIds: [] });
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["done", "skipped", "pending"]), workflowStepResults: [result(ids[2])] })).toMatchObject({ missingInstanceIds: [] });
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending"]), workflowStepResults: [result(ids[0], "failed")] }).missingInstanceIds).toEqual([ids[0]]);
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending"]), workflowStepResults: [{ ...result(ids[0]), phase: "post-merge" }] }).missingInstanceIds).toEqual([ids[0]]);
|
||||
expect(evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending"]), workflowStepResults: [result(ids[0]), result(ids[0])] }).missingInstanceIds).toEqual([]);
|
||||
});
|
||||
|
||||
it("uses persisted rows only to widen expected identities, never as satisfaction", () => {
|
||||
const id0 = instanceNodeId("steps", 0, "step-execute");
|
||||
const id1 = instanceNodeId("steps", 1, "step-execute");
|
||||
const proof = evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending"]), workflowStepResults: [result(id0)], persistedInstances: [{ foreachNodeId: "steps", stepIndex: 1, pinnedStepCount: 2 }] });
|
||||
expect(proof.expectedInstanceIds).toEqual([id0, id1]);
|
||||
expect(proof.missingInstanceIds).toEqual([id1]);
|
||||
});
|
||||
|
||||
it("keeps coverage independent from unrelated failed node results", () => {
|
||||
const ids = [0, 1].map((index) => instanceNodeId("steps", index, "step-execute"));
|
||||
const proof = evaluateForeachMergeProof({ ir: foreachIr(), steps: steps(["pending", "pending"]), workflowStepResults: [...ids.map((id) => result(id)), result("unrelated", "failed")] });
|
||||
expect(proof.missingInstanceIds).toEqual([]);
|
||||
});
|
||||
});
|
||||
@@ -248,6 +248,8 @@ export type {
|
||||
EffectiveAgentInput,
|
||||
EffectiveAgentResult,
|
||||
} from "./column-agent-resolver.js";
|
||||
export { evaluateForeachMergeProof } from "./workflow-merge-proof.js";
|
||||
export type { ForeachMergeProof, ForeachMergeProofInput } from "./workflow-merge-proof.js";
|
||||
export { BUILTIN_CODING_WORKFLOW_IR } from "./builtin-coding-workflow-ir.js";
|
||||
export { PLAN_REVIEW_GROUP_ID } from "./builtin-plan-review-group.js";
|
||||
export { BUILTIN_MARKETING_WORKFLOW_IR } from "./builtin-marketing-workflow-ir.js";
|
||||
|
||||
@@ -277,6 +277,8 @@ export type {
|
||||
EffectiveAgentInput,
|
||||
EffectiveAgentResult,
|
||||
} from "./column-agent-resolver.js";
|
||||
export { evaluateForeachMergeProof } from "./workflow-merge-proof.js";
|
||||
export type { ForeachMergeProof, ForeachMergeProofInput } from "./workflow-merge-proof.js";
|
||||
export { BUILTIN_CODING_WORKFLOW_IR } from "./builtin-coding-workflow-ir.js";
|
||||
export { BUILTIN_CODING_IDEAS_WORKFLOW_IR } from "./builtin-coding-ideas-workflow-ir.js";
|
||||
export { PLAN_REVIEW_GROUP_ID } from "./builtin-plan-review-group.js";
|
||||
|
||||
74
packages/core/src/workflow-merge-proof.ts
Normal file
74
packages/core/src/workflow-merge-proof.ts
Normal file
@@ -0,0 +1,74 @@
|
||||
import { instanceNodeId } from "./column-agent-resolver.js";
|
||||
import type { TaskStep, WorkflowRunStepInstance, WorkflowStepResult } from "./types.js";
|
||||
import type { WorkflowIr } from "./workflow-ir-types.js";
|
||||
|
||||
export interface ForeachMergeProofInput {
|
||||
ir: WorkflowIr;
|
||||
steps: readonly Pick<TaskStep, "status">[] | undefined;
|
||||
workflowStepResults: readonly WorkflowStepResult[] | undefined;
|
||||
persistedInstances?: readonly Pick<WorkflowRunStepInstance, "foreachNodeId" | "stepIndex" | "pinnedStepCount">[];
|
||||
}
|
||||
|
||||
export interface ForeachMergeProof {
|
||||
expectedInstanceIds: string[];
|
||||
satisfiedInstanceIds: string[];
|
||||
missingInstanceIds: string[];
|
||||
hasForeachStepExecute: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* FNXC:WorkflowMerge 2026-07-27-12:00:
|
||||
* FN-8601 requires every expanded step-execute identity to have a terminal-success
|
||||
* workflowStepResult, unless its live task step is already done/skipped. Callers
|
||||
* must independently require a relevant pre-merge node result and that every such
|
||||
* result is passed/skipped: coverage is vacuously complete on non-foreach/no-seam
|
||||
* shapes and is never sole proof. Persisted rows only widen expectation; partial
|
||||
* projections cannot complete a checklist, prove implementation, or enter merge.
|
||||
*/
|
||||
export function evaluateForeachMergeProof(input: ForeachMergeProofInput): ForeachMergeProof {
|
||||
const steps = input.steps ?? [];
|
||||
const foreachTemplates = input.ir.nodes
|
||||
.filter((node) => node.kind === "foreach")
|
||||
.flatMap((node) => {
|
||||
const template = (node.config as { template?: { nodes?: WorkflowIr["nodes"] } } | undefined)?.template;
|
||||
return (template?.nodes ?? [])
|
||||
.filter((templateNode) => templateNode.kind === "prompt" && templateNode.config?.seam === "step-execute")
|
||||
.map((templateNode) => ({ foreachNodeId: node.id, templateNodeId: templateNode.id }));
|
||||
});
|
||||
const hasForeachStepExecute = foreachTemplates.length > 0;
|
||||
const expected = new Set<string>();
|
||||
const terminalLiveSteps = new Set<string>();
|
||||
|
||||
for (const { foreachNodeId, templateNodeId } of foreachTemplates) {
|
||||
for (let stepIndex = 0; stepIndex < steps.length; stepIndex += 1) {
|
||||
const id = instanceNodeId(foreachNodeId, stepIndex, templateNodeId);
|
||||
expected.add(id);
|
||||
if (steps[stepIndex]?.status === "done" || steps[stepIndex]?.status === "skipped") terminalLiveSteps.add(id);
|
||||
}
|
||||
}
|
||||
|
||||
// Rows may survive an IR/live-step discrepancy. Their pinned count can only widen
|
||||
// the IR/live expectation; a completed row never satisfies an identity by itself.
|
||||
for (const template of foreachTemplates) {
|
||||
const rows = (input.persistedInstances ?? []).filter((row) => row.foreachNodeId === template.foreachNodeId);
|
||||
const persistedCount = rows.reduce((maximum, row) => Math.max(maximum, row.pinnedStepCount, row.stepIndex + 1), 0);
|
||||
for (let stepIndex = 0; stepIndex < Math.max(steps.length, persistedCount); stepIndex += 1) {
|
||||
expected.add(instanceNodeId(template.foreachNodeId, stepIndex, template.templateNodeId));
|
||||
}
|
||||
}
|
||||
|
||||
const successfulResults = new Set(
|
||||
(input.workflowStepResults ?? [])
|
||||
.filter((result) => result.source === "node" && (result.phase ?? "pre-merge") === "pre-merge")
|
||||
.filter((result) => result.status === "passed" || result.status === "skipped")
|
||||
.map((result) => result.workflowStepId),
|
||||
);
|
||||
const expectedInstanceIds = [...expected].sort();
|
||||
const satisfiedInstanceIds = expectedInstanceIds.filter((id) => terminalLiveSteps.has(id) || successfulResults.has(id));
|
||||
return {
|
||||
expectedInstanceIds,
|
||||
satisfiedInstanceIds,
|
||||
missingInstanceIds: expectedInstanceIds.filter((id) => !terminalLiveSteps.has(id) && !successfulResults.has(id)),
|
||||
hasForeachStepExecute,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,111 @@
|
||||
import { afterEach, expect, it, vi } from "vitest";
|
||||
import { instanceNodeId, type TaskStep, type TaskStore, type WorkflowIr } from "@fusion/core";
|
||||
import { createTaskStoreForTest, pgDescribe, type PgTestHarness } from "../../../core/src/__test-utils__/pg-test-harness.js";
|
||||
import { TaskExecutor } from "../executor.js";
|
||||
|
||||
const runId = "foreach-proof-run";
|
||||
const columns = [{ id: "todo", name: "Todo", traits: [] }, { id: "in-review", name: "Review", traits: [{ trait: "merge" }] }];
|
||||
const foreachIr: WorkflowIr = {
|
||||
version: "v2", name: "foreach proof", columns,
|
||||
nodes: [
|
||||
{ id: "parse", kind: "parse-steps", column: "todo" },
|
||||
{ id: "steps", kind: "foreach", column: "todo", config: { source: "task-steps", template: { nodes: [{ id: "step-execute", kind: "prompt", config: { seam: "step-execute" } }], edges: [] } } },
|
||||
{ id: "merge", kind: "merge-gate", column: "in-review" }, { id: "end", kind: "end" },
|
||||
], edges: [],
|
||||
};
|
||||
const foreachWithoutSeamIr: WorkflowIr = {
|
||||
...foreachIr,
|
||||
nodes: foreachIr.nodes.map((node) => node.id === "steps"
|
||||
? { ...node, config: { source: "task-steps", template: { nodes: [{ id: "review", kind: "prompt" }], edges: [] } } }
|
||||
: node),
|
||||
};
|
||||
const executeIr: WorkflowIr = {
|
||||
version: "v2", name: "execute proof", columns,
|
||||
nodes: [{ id: "execute", kind: "prompt", column: "todo", config: { seam: "execute" } }, { id: "merge", kind: "merge-gate", column: "in-review" }], edges: [],
|
||||
};
|
||||
const nodeResult = (id: string, status: "passed" | "failed" = "passed") => ({ workflowStepId: id, source: "node" as const, phase: "pre-merge" as const, status, completedAt: new Date().toISOString() });
|
||||
const pendingSteps = (): TaskStep[] => ["one", "two", "three"].map((title, index) => ({ id: `${index}`, title, status: "pending" }));
|
||||
|
||||
let harness: PgTestHarness | undefined;
|
||||
afterEach(async () => { await harness?.teardown(); harness = undefined; vi.restoreAllMocks(); });
|
||||
|
||||
async function setup({ ir = foreachIr, results = [], steps = pendingSteps() }: { ir?: WorkflowIr; results?: ReturnType<typeof nodeResult>[]; steps?: TaskStep[] } = {}) {
|
||||
harness = await createTaskStoreForTest();
|
||||
const store = harness.store;
|
||||
const task = await store.createTask({ description: "foreach merge proof" });
|
||||
await store.moveTask(task.id, "todo");
|
||||
await store.moveTask(task.id, "in-progress");
|
||||
await store.updateTask(task.id, { steps, workflowStepResults: results });
|
||||
vi.spyOn(store as unknown as { getTaskWorkflowSelectionAsync: () => Promise<unknown> }, "getTaskWorkflowSelectionAsync").mockResolvedValue({ workflowId: "wf-proof", stepIds: [] });
|
||||
vi.spyOn(store as unknown as { getTaskWorkflowSelection: () => unknown }, "getTaskWorkflowSelection").mockReturnValue({ workflowId: "wf-proof", stepIds: [] });
|
||||
vi.spyOn(store as unknown as { getWorkflowDefinition: () => Promise<unknown> }, "getWorkflowDefinition").mockResolvedValue({ id: "wf-proof", ir });
|
||||
return { store, task, executor: new TaskExecutor(store as TaskStore, process.cwd()) };
|
||||
}
|
||||
|
||||
async function merge(executor: TaskExecutor, task: unknown) {
|
||||
return (executor as unknown as { ensureWorkflowMergeBoundaryTask: (t: unknown, m: unknown) => Promise<unknown> })
|
||||
.ensureWorkflowMergeBoundaryTask(task, { reason: "test", nodeId: "merge", workflowId: "wf-proof", runId });
|
||||
}
|
||||
|
||||
async function proofFailure(executor: TaskExecutor, task: unknown) {
|
||||
return (executor as unknown as { getWorkflowMergeImplementationProofFailure: (t: unknown) => Promise<string | undefined> })
|
||||
.getWorkflowMergeImplementationProofFailure(task);
|
||||
}
|
||||
|
||||
pgDescribe("FN-8601 — PostgreSQL foreach merge proof", () => {
|
||||
it("blocks partial results even with completed persisted rows, then advances complete coverage", async () => {
|
||||
const ids = [0, 1, 2].map((index) => instanceNodeId("steps", index, "step-execute"));
|
||||
const { store, task, executor } = await setup({ results: [nodeResult(ids[0])] });
|
||||
for (let index = 0; index < 3; index += 1) await store.saveWorkflowRunStepInstanceAsync({ taskId: task.id, runId, foreachNodeId: "steps", stepIndex: index, pinnedStepCount: 3, currentNodeId: "step-execute", status: "completed", reworkCount: 0 });
|
||||
await merge(executor, task);
|
||||
let live = await store.getTask(task.id);
|
||||
expect(live?.column).not.toBe("in-review");
|
||||
expect(live?.steps.map((step) => step.status)).toEqual(["pending", "pending", "pending"]);
|
||||
expect(await proofFailure(executor, live!)).toContain("foreach step instances are incomplete");
|
||||
const evaluation = await (executor as unknown as { evaluateWorkflowMergeBoundary: (t: unknown, run?: string) => Promise<{ missingInstanceIds: string[] }> }).evaluateWorkflowMergeBoundary(live!, runId);
|
||||
expect(evaluation.missingInstanceIds).toEqual([ids[1], ids[2]]);
|
||||
await store.updateTask(task.id, { workflowStepResults: ids.map((id) => nodeResult(id)) });
|
||||
live = await store.getTask(task.id);
|
||||
await merge(executor, live!);
|
||||
live = await store.getTask(task.id);
|
||||
expect(live?.column).toBe("in-review");
|
||||
expect(live?.steps.every((step) => step.status === "done")).toBe(true);
|
||||
});
|
||||
|
||||
it("blocks complete coverage when an unrelated pre-merge node result failed", async () => {
|
||||
const ids = [0, 1, 2].map((index) => instanceNodeId("steps", index, "step-execute"));
|
||||
const { store, task, executor } = await setup({ results: [...ids.map((id) => nodeResult(id)), nodeResult("unrelated", "failed")] });
|
||||
await merge(executor, task);
|
||||
const live = await store.getTask(task.id);
|
||||
expect(live?.column).not.toBe("in-review");
|
||||
expect(live?.steps.map((step) => step.status)).toEqual(["pending", "pending", "pending"]);
|
||||
});
|
||||
|
||||
it("keeps empty results as missing implementation proof for execute and foreach-without-seam workflows", async () => {
|
||||
for (const ir of [executeIr, foreachWithoutSeamIr]) {
|
||||
const { store, task, executor } = await setup({ ir, results: [] });
|
||||
const live = await store.getTask(task.id);
|
||||
expect(await proofFailure(executor, live!)).toContain("implementation did not run");
|
||||
await merge(executor, live!);
|
||||
const afterMerge = await store.getTask(task.id);
|
||||
expect(afterMerge?.steps.map((step) => step.status)).toEqual(["pending", "pending", "pending"]);
|
||||
}
|
||||
});
|
||||
|
||||
it("preserves non-foreach and resumed foreach merge behavior", async () => {
|
||||
const { store: executeStore, task: executeTask, executor: executeExecutor } = await setup({ ir: executeIr, results: [nodeResult("execute")] });
|
||||
await merge(executeExecutor, executeTask);
|
||||
let live = await executeStore.getTask(executeTask.id);
|
||||
expect(live?.column).toBe("in-review");
|
||||
expect(live?.steps.every((step) => step.status === "done")).toBe(true);
|
||||
await harness?.teardown(); harness = undefined;
|
||||
|
||||
const id = instanceNodeId("steps", 2, "step-execute");
|
||||
const resumed = pendingSteps(); resumed[0] = { ...resumed[0], status: "done" }; resumed[1] = { ...resumed[1], status: "skipped" };
|
||||
const { store, task, executor } = await setup({ results: [nodeResult(id)], steps: resumed });
|
||||
await merge(executor, task);
|
||||
live = await store.getTask(task.id);
|
||||
expect(live?.column).toBe("in-review");
|
||||
expect(live?.steps.every((step) => step.status === "done" || step.status === "skipped")).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -14,7 +14,7 @@ import { existsSync, lstatSync, realpathSync } from "node:fs";
|
||||
import { readFile, rm, writeFile } from "node:fs/promises";
|
||||
import type { TaskStore, Task, TaskDetail, TaskTokenUsage, StepStatus, Settings, WorkflowStep, MissionStore, AsyncMissionStore, Slice, AgentState, AgentCapability, RunMutationContext, AgentHeartbeatConfig, Agent, AgentMemoryInclusionMode, ProjectSettings, MergeResult, WorkflowIrNode, WorkflowIrNodeKind, WorkflowStepResult as CoreWorkflowStepResult, ThinkingLevel } from "@fusion/core";
|
||||
import { getUnmetSchedulingDependencies } from "./scheduler.js";
|
||||
import { RetryStormError, serializeRetryStormError, evaluateCompletedPromotionFailureProvenance, evaluateSkipBypassTaint, resolveWorkflowIrForTask, resolveCompleteColumn, resolveMergeOrchestrationColumn, resolveReboundTarget, resolveColumnAgentBinding, resolveEffectiveAgent, instanceNodeId, getWorkflowExtensionRegistry, getBuiltinWorkflow, parseNoOpCompletionMarker, allowsAutoMergeProcessing, resolveEffectiveAutoMerge, isLiveSharedBranchGroupMemberIntegration, resolveMaxAutoMergeRetries, resolveMaxConsecutiveToolFailureRetries, resolveConsecutiveToolFailureRetryBackoffMs, resolveConsecutiveToolFailureThreshold, resolveExecutorEscalationTarget, resolveOptionalStepRevisionBudget, resolveOptionalReviewRevisionBudget, COMPLETION_SUMMARY_NODE_ID, upsertWorkflowStepResult, AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, ACTIVE_WORKFLOW_WORK_ITEM_STATES, AgentStore, resolveExecutorFallbackModel } from "@fusion/core";
|
||||
import { RetryStormError, serializeRetryStormError, evaluateCompletedPromotionFailureProvenance, evaluateSkipBypassTaint, resolveWorkflowIrForTask, evaluateForeachMergeProof, resolveCompleteColumn, resolveMergeOrchestrationColumn, resolveReboundTarget, resolveColumnAgentBinding, resolveEffectiveAgent, instanceNodeId, getWorkflowExtensionRegistry, getBuiltinWorkflow, parseNoOpCompletionMarker, allowsAutoMergeProcessing, resolveEffectiveAutoMerge, isLiveSharedBranchGroupMemberIntegration, resolveMaxAutoMergeRetries, resolveMaxConsecutiveToolFailureRetries, resolveConsecutiveToolFailureRetryBackoffMs, resolveConsecutiveToolFailureThreshold, resolveExecutorEscalationTarget, resolveOptionalStepRevisionBudget, resolveOptionalReviewRevisionBudget, COMPLETION_SUMMARY_NODE_ID, upsertWorkflowStepResult, AWAITING_APPROVAL_PAUSE_REASON, THINKING_LEVELS, ACTIVE_WORKFLOW_WORK_ITEM_STATES, AgentStore, resolveExecutorFallbackModel } from "@fusion/core";
|
||||
import { finalizeProvenAutoMergeTask } from "./auto-merge-finalization.js";
|
||||
import { mergeEffectiveSettings } from "./effective-settings.js";
|
||||
import { generateFeatureVideo, type GenerateFeatureVideoOptions } from "./review-artifacts/feature-video.js";
|
||||
@@ -7623,7 +7623,18 @@ export class TaskExecutor {
|
||||
FNXC:WorkflowMerge 2026-06-29-15:28:
|
||||
Compound Engineering and similar graph-native workflows execute skill nodes instead of legacy parsed task steps. The graph records those nodes as `workflowStepResults.source = "node"`; at the merge boundary, project a successful graph-native run onto the legacy checklist so `task has incomplete steps` cannot block a workflow that already completed its authoritative nodes.
|
||||
*/
|
||||
if (this.shouldCompleteChecklistAtWorkflowMerge(live)) {
|
||||
const mergeProof = await this.evaluateWorkflowMergeBoundary(live, metadata.runId);
|
||||
if (mergeProof.hasForeachStepExecute && !mergeProof.complete) {
|
||||
const reason = !mergeProof.hasRelevantNodeResult
|
||||
? "no pre-merge node result recorded"
|
||||
: !mergeProof.allResultsTerminal
|
||||
? `non-terminal pre-merge node result ${mergeProof.nonTerminalResult?.workflowStepId ?? "unknown"} (${mergeProof.nonTerminalResult?.status ?? "unknown"})`
|
||||
: `foreach step instances incomplete at merge boundary: missing ${mergeProof.missingInstanceIds.join(", ")}`;
|
||||
await this.store.logEntry(live.id, `Workflow merge boundary blocked: ${reason}`, undefined, this.getRunContextFor(live.id));
|
||||
return live;
|
||||
}
|
||||
|
||||
if (this.shouldCompleteChecklistAtWorkflowMerge(live, mergeProof)) {
|
||||
const completedSteps = live.steps.map((step) =>
|
||||
step.status === "done" || step.status === "skipped"
|
||||
? step
|
||||
@@ -7664,6 +7675,49 @@ export class TaskExecutor {
|
||||
return { ...live, column: targetColumn };
|
||||
}
|
||||
|
||||
private async evaluateWorkflowMergeBoundary(task: TaskDetail, runId?: string): Promise<{
|
||||
resolved: boolean;
|
||||
hasRelevantNodeResult: boolean;
|
||||
allResultsTerminal: boolean;
|
||||
coverageComplete: boolean;
|
||||
hasForeachStepExecute: boolean;
|
||||
missingInstanceIds: string[];
|
||||
nonTerminalResult?: CoreWorkflowStepResult;
|
||||
complete: boolean;
|
||||
}> {
|
||||
const relevant = (task.workflowStepResults ?? []).filter((result) =>
|
||||
result.source === "node" && (result.phase ?? "pre-merge") === "pre-merge",
|
||||
);
|
||||
// FNXC:WorkflowMerge 2026-07-27-12:30: FN-8601 keeps required presence
|
||||
// independent from terminality: a failed node result proves execution occurred,
|
||||
// while allResultsTerminal separately rejects it at the merge boundary.
|
||||
const hasRelevantNodeResult = relevant.length > 0;
|
||||
const nonTerminalResult = relevant.find((result) => result.status !== "passed" && result.status !== "skipped");
|
||||
const allResultsTerminal = nonTerminalResult === undefined;
|
||||
let ir: WorkflowIr | undefined;
|
||||
try { ir = await resolveWorkflowIrForTask(this.store, task.id); } catch { /* preserve legacy behavior for unresolved IRs */ }
|
||||
if (!ir) return { resolved: false, hasRelevantNodeResult, allResultsTerminal, coverageComplete: true, hasForeachStepExecute: false, missingInstanceIds: [], nonTerminalResult, complete: false };
|
||||
|
||||
let persistedInstances: Awaited<ReturnType<typeof this.loadMergeBoundaryInstances>> = [];
|
||||
try { persistedInstances = await this.loadMergeBoundaryInstances(task.id, runId); } catch { /* persistence is additive */ }
|
||||
const coverage = evaluateForeachMergeProof({ ir, steps: task.steps, workflowStepResults: task.workflowStepResults, persistedInstances });
|
||||
const complete = hasRelevantNodeResult && allResultsTerminal && coverage.missingInstanceIds.length === 0;
|
||||
return { resolved: true, hasRelevantNodeResult, allResultsTerminal, coverageComplete: coverage.missingInstanceIds.length === 0, hasForeachStepExecute: coverage.hasForeachStepExecute, missingInstanceIds: coverage.missingInstanceIds, nonTerminalResult, complete };
|
||||
}
|
||||
|
||||
private async loadMergeBoundaryInstances(taskId: string, runId?: string): Promise<Array<{ foreachNodeId: string; stepIndex: number; pinnedStepCount: number }>> {
|
||||
if (!runId) return [];
|
||||
const store = this.store as typeof this.store & {
|
||||
loadWorkflowRunStepInstancesAsync?: (id: string, idRun: string) => Promise<Array<{ foreachNodeId: string; stepIndex: number; pinnedStepCount: number }>>;
|
||||
loadWorkflowRunStepInstances?: (id: string, idRun: string) => Array<{ foreachNodeId: string; stepIndex: number; pinnedStepCount: number }>;
|
||||
};
|
||||
try {
|
||||
return await store.loadWorkflowRunStepInstancesAsync?.(taskId, runId)
|
||||
?? store.loadWorkflowRunStepInstances?.(taskId, runId)
|
||||
?? [];
|
||||
} catch { return []; }
|
||||
}
|
||||
|
||||
private async getWorkflowMergeImplementationProofFailure(task: TaskDetail): Promise<string | undefined> {
|
||||
/*
|
||||
FNXC:Lifecycle 2026-07-16-21:40:
|
||||
@@ -7674,62 +7728,41 @@ export class TaskExecutor {
|
||||
Runs before the noCommitsExpected exemption so a tainted task cannot slip past it.
|
||||
*/
|
||||
const taint = evaluateSkipBypassTaint(task);
|
||||
if (taint.blocked) {
|
||||
return "implementation did not run: steps were skipped after a bulk-step-completion refusal without an accepted fn_task_done";
|
||||
}
|
||||
if (taint.blocked) return "implementation did not run: steps were skipped after a bulk-step-completion refusal without an accepted fn_task_done";
|
||||
if (task.noCommitsExpected === true) return undefined;
|
||||
|
||||
let ir: WorkflowIr | undefined;
|
||||
try {
|
||||
ir = await resolveWorkflowIrForTask(this.store, task.id);
|
||||
} catch {
|
||||
ir = undefined;
|
||||
}
|
||||
try { ir = await resolveWorkflowIrForTask(this.store, task.id); } catch { ir = undefined; }
|
||||
if (!ir) return undefined;
|
||||
|
||||
const usesParsedSteps = ir.nodes.some((node) => node.kind === "parse-steps");
|
||||
const usesExecuteSeam = ir.nodes.some((node) => node.kind === "prompt" && node.config?.seam === "execute");
|
||||
if (!usesParsedSteps && !usesExecuteSeam) return undefined;
|
||||
|
||||
const steps = Array.isArray(task.steps) ? task.steps : [];
|
||||
const hasTerminalParsedSteps =
|
||||
steps.length > 0 && steps.every((step) => step.status === "done" || step.status === "skipped");
|
||||
const hasTerminalParsedSteps = steps.length > 0 && steps.every((step) => step.status === "done" || step.status === "skipped");
|
||||
const hasModifiedFiles = (task.modifiedFiles?.length ?? 0) > 0;
|
||||
const hasGraphNativeImplementationProof = (task.workflowStepResults ?? []).some((result) =>
|
||||
result.source === "node"
|
||||
&& (result.phase ?? "pre-merge") === "pre-merge"
|
||||
&& (result.status === "passed" || result.status === "skipped")
|
||||
);
|
||||
|
||||
/*
|
||||
FNXC:WorkflowMerge 2026-06-30-00:38:
|
||||
Stepwise Coding proves implementation through parsed task steps/foreach projection. Legacy monolithic Coding may prove through modified files or explicit no-op completion. Do not accept an empty step list as success for parse-step workflows; a valid PROMPT.md with unparsed steps must resume execution, not no-op merge.
|
||||
*/
|
||||
const proof = await this.evaluateWorkflowMergeBoundary(task);
|
||||
const hasGraphNativeImplementationProof = proof.hasRelevantNodeResult && proof.allResultsTerminal && proof.coverageComplete;
|
||||
if (usesParsedSteps) {
|
||||
return hasTerminalParsedSteps || hasGraphNativeImplementationProof
|
||||
? undefined
|
||||
if (hasTerminalParsedSteps || hasGraphNativeImplementationProof) return undefined;
|
||||
return proof.hasForeachStepExecute && !proof.coverageComplete
|
||||
? `implementation did not run: foreach step instances are incomplete (missing ${proof.missingInstanceIds.join(", ")})`
|
||||
: "implementation did not run: parsed coding steps are missing or incomplete";
|
||||
}
|
||||
|
||||
if (usesExecuteSeam) {
|
||||
return hasTerminalParsedSteps || hasModifiedFiles || hasGraphNativeImplementationProof
|
||||
? undefined
|
||||
: "implementation did not run: execute seam has no completion proof";
|
||||
}
|
||||
|
||||
if (usesExecuteSeam) return hasTerminalParsedSteps || hasModifiedFiles || hasGraphNativeImplementationProof ? undefined : "implementation did not run: execute seam has no completion proof";
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private shouldCompleteChecklistAtWorkflowMerge(task: TaskDetail): boolean {
|
||||
/*
|
||||
FNXC:WorkflowMerge 2026-07-27-12:00:
|
||||
FN-8601 gates checklist projection and foreach merge admission on required node-result
|
||||
presence, terminal status for every present result, and expanded-instance coverage.
|
||||
Non-foreach/no-seam coverage is vacuous and does not change legacy move behavior.
|
||||
*/
|
||||
private shouldCompleteChecklistAtWorkflowMerge(task: TaskDetail, proof?: { complete: boolean }): boolean {
|
||||
if (!Array.isArray(task.steps) || task.steps.length === 0) return false;
|
||||
if (task.steps.every((step) => step.status === "done" || step.status === "skipped")) return false;
|
||||
|
||||
const graphNodeResults = (task.workflowStepResults ?? []).filter((result) =>
|
||||
result.source === "node" && (result.phase ?? "pre-merge") === "pre-merge"
|
||||
);
|
||||
if (graphNodeResults.length === 0) return false;
|
||||
|
||||
return graphNodeResults.every((result) => result.status === "passed" || result.status === "skipped");
|
||||
if (proof) return proof.complete;
|
||||
const graphNodeResults = (task.workflowStepResults ?? []).filter((result) => result.source === "node" && (result.phase ?? "pre-merge") === "pre-merge");
|
||||
return graphNodeResults.length > 0 && graphNodeResults.every((result) => result.status === "passed" || result.status === "skipped");
|
||||
}
|
||||
|
||||
public createAuthoritativeWorkflowSeams(_settings: Settings): WorkflowLegacySeams {
|
||||
|
||||
Reference in New Issue
Block a user