fix(engine): recover completed triage tasks
Preserve workflow ownership across Plan Review replan moves and route advanced completed triage rows through legal lifecycle transitions before review. Clear only stale same-task session claims after live executor, planner, and merger ownership checks.
This commit is contained in:
7
.changeset/quiet-triage-recovery.md
Normal file
7
.changeset/quiet-triage-recovery.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Prevent Plan Review replans from stranding completed tasks in Triage and recover affected tasks automatically.
|
||||||
|
category: fix
|
||||||
|
dev: Preserves graph ownership during executor-authored replan moves and clears stale same-task session claims during recovery.
|
||||||
@@ -1022,6 +1022,44 @@ describe("In-progress task resume after restart", () => {
|
|||||||
}));
|
}));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("recoverCompletedTask() legally re-homes a completed triage zombie before review handoff", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = makeTask("FN-TRIAGE-ZOMBIE", "triage", {
|
||||||
|
worktree: "/tmp/wt/FN-TRIAGE-ZOMBIE",
|
||||||
|
steps: makeSteps("done"),
|
||||||
|
enabledWorkflowSteps: ["plan-review", "code-review"],
|
||||||
|
workflowIrPinNodeId: "merge",
|
||||||
|
workflowStepResults: [
|
||||||
|
{ workflowStepId: "plan-review", workflowStepName: "Plan Review", phase: "pre-merge", status: "passed" },
|
||||||
|
{ workflowStepId: "code-review", workflowStepName: "Code Review", phase: "pre-merge", status: "passed" },
|
||||||
|
],
|
||||||
|
});
|
||||||
|
store.getTask.mockResolvedValue(makeTaskDetail("FN-TRIAGE-ZOMBIE", "triage", {
|
||||||
|
worktree: task.worktree,
|
||||||
|
steps: makeSteps("done"),
|
||||||
|
enabledWorkflowSteps: task.enabledWorkflowSteps,
|
||||||
|
workflowIrPinNodeId: "merge",
|
||||||
|
workflowStepResults: task.workflowStepResults,
|
||||||
|
}));
|
||||||
|
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/test");
|
||||||
|
vi.spyOn(executor as any, "captureModifiedFiles").mockResolvedValue([]);
|
||||||
|
|
||||||
|
const recovered = await executor.recoverCompletedTask(task);
|
||||||
|
|
||||||
|
expect(recovered).toBe(true);
|
||||||
|
expect(store.moveTask).toHaveBeenNthCalledWith(
|
||||||
|
1,
|
||||||
|
task.id,
|
||||||
|
"todo",
|
||||||
|
expect.objectContaining({ recoveryRehome: true, preserveProgress: true, preserveWorktree: true }),
|
||||||
|
);
|
||||||
|
expect(store.moveTask).toHaveBeenNthCalledWith(2, task.id, "in-progress");
|
||||||
|
expect(store.handoffToReview).toHaveBeenCalledWith(task.id, expect.objectContaining({
|
||||||
|
evidence: expect.objectContaining({ reason: "completed-task-recovered" }),
|
||||||
|
}));
|
||||||
|
});
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:EngineTests 2026-07-19-18:55 (U10b):
|
FNXC:EngineTests 2026-07-19-18:55 (U10b):
|
||||||
"Enabled steps are ABSENT" is the condition under test, so it must be stated rather than
|
"Enabled steps are ABSENT" is the condition under test, so it must be stated rather than
|
||||||
|
|||||||
@@ -95,6 +95,32 @@ describe("advanced workflow tasks stranded in triage", () => {
|
|||||||
expect(store.moveTaskIf).not.toHaveBeenCalled();
|
expect(store.moveTaskIf).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("clears a stale same-task session-path claim before promoting completed pinned work", async () => {
|
||||||
|
const stranded = task("FN-STALE-REGISTRY", {
|
||||||
|
workflowIrPinNodeId: "merge",
|
||||||
|
workflowIrPinColumnId: undefined,
|
||||||
|
steps: [{ name: "Implement", status: "done" }],
|
||||||
|
});
|
||||||
|
const store = storeFor([stranded]);
|
||||||
|
const recoverCompletedTask = vi.fn(async () => true);
|
||||||
|
activeSessionRegistry.registerPath(stranded.worktree!, {
|
||||||
|
taskId: stranded.id,
|
||||||
|
kind: "workflow-step",
|
||||||
|
ownerKey: `${stranded.id}#stale`,
|
||||||
|
});
|
||||||
|
const manager = new SelfHealingManager(store, {
|
||||||
|
rootDir: "/repo",
|
||||||
|
recoverCompletedTask,
|
||||||
|
getExecutingTaskIds: () => new Set<string>(),
|
||||||
|
getPlanningTaskIds: () => new Set<string>(),
|
||||||
|
isTaskActive: () => false,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(await manager.recoverAdvancedTriageTasks()).toBe(1);
|
||||||
|
expect(recoverCompletedTask).toHaveBeenCalledWith(stranded);
|
||||||
|
expect(activeSessionRegistry.isPathActive(stranded.worktree!)).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
it("leaves ordinary planning rows and actively-owned graph runs untouched", async () => {
|
it("leaves ordinary planning rows and actively-owned graph runs untouched", async () => {
|
||||||
const ordinary = task("FN-ORDINARY", { worktree: undefined, workflowIrPinNodeId: undefined });
|
const ordinary = task("FN-ORDINARY", { worktree: undefined, workflowIrPinNodeId: undefined });
|
||||||
const active = task("FN-ACTIVE");
|
const active = task("FN-ACTIVE");
|
||||||
|
|||||||
@@ -178,6 +178,44 @@ describe("TaskExecutor pre-merge optional-step fix seam", () => {
|
|||||||
expect((executor as any).pausedAborted.has("FN-7066")).toBe(false);
|
expect((executor as any).pausedAborted.has("FN-7066")).toBe(false);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("does not hard-cancel the graph that performs its own Plan Review replan move", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const liveTask = task({ postReviewFixCount: 0, column: "in-progress", status: null });
|
||||||
|
store.getTask.mockResolvedValue(liveTask);
|
||||||
|
store.getSettings.mockResolvedValue({ maxPostReviewFixes: 3 });
|
||||||
|
const executor = new TaskExecutor(store, "/tmp/test");
|
||||||
|
const abortSpy = vi
|
||||||
|
.spyOn(executor as any, "awaitAbortInFlightTaskWork")
|
||||||
|
.mockResolvedValue(undefined);
|
||||||
|
store.moveTask.mockImplementation(async (_taskId: string, column: string) => {
|
||||||
|
await (store as any)._triggerAsync("task:moved", {
|
||||||
|
task: { ...liveTask, column },
|
||||||
|
from: "in-progress",
|
||||||
|
to: column,
|
||||||
|
source: "engine",
|
||||||
|
});
|
||||||
|
return { ...liveTask, column };
|
||||||
|
});
|
||||||
|
|
||||||
|
(executor as any).graphRouting.add(liveTask.id);
|
||||||
|
try {
|
||||||
|
await (executor as any).requestPreMergeOptionalStepFix(liveTask.id, liveTask, {
|
||||||
|
stepName: "Plan Review",
|
||||||
|
feedback: "PROMPT.md needs a revision",
|
||||||
|
phase: "pre-merge" as const,
|
||||||
|
status: "failed" as const,
|
||||||
|
verdict: "REVISE",
|
||||||
|
nodeId: "plan-review",
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(store.moveTask).toHaveBeenCalledWith(liveTask.id, "triage");
|
||||||
|
expect(abortSpy).not.toHaveBeenCalled();
|
||||||
|
expect((executor as any).pausedAborted.has(liveTask.id)).toBe(false);
|
||||||
|
} finally {
|
||||||
|
(executor as any).graphRouting.delete(liveTask.id);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
it("honors Plan Review workflow-setting caps before automatic replan", async () => {
|
it("honors Plan Review workflow-setting caps before automatic replan", async () => {
|
||||||
const zeroStore = createMockStore();
|
const zeroStore = createMockStore();
|
||||||
const zeroTask = task({ postReviewFixCount: 0, column: "in-progress" });
|
const zeroTask = task({ postReviewFixCount: 0, column: "in-progress" });
|
||||||
|
|||||||
@@ -1662,11 +1662,12 @@ export class TaskExecutor {
|
|||||||
private workflowRerunPending = new Set<string>();
|
private workflowRerunPending = new Set<string>();
|
||||||
/**
|
/**
|
||||||
* Task ids whose current `task:moved` event is being emitted by this
|
* Task ids whose current `task:moved` event is being emitted by this
|
||||||
* executor's workflow column-boundary hook. The store emits synchronously,
|
* executor's workflow lifecycle handling (column boundaries or Plan Review
|
||||||
* so this narrowly distinguishes a graph's own transition from an external
|
* replans). The store emits synchronously, so this narrowly distinguishes a
|
||||||
* engine/user move that must still hard-cancel the active run.
|
* graph's own transition from an external engine/user move that must still
|
||||||
|
* hard-cancel the active run.
|
||||||
*/
|
*/
|
||||||
private workflowBoundaryMovesInFlight = new Set<string>();
|
private workflowLifecycleMovesInFlight = new Set<string>();
|
||||||
/** FN-5256: in-flight session-disposal promises keyed by taskId. The
|
/** FN-5256: in-flight session-disposal promises keyed by taskId. The
|
||||||
* task:moved (away from in-progress) and task:deleted listeners populate
|
* task:moved (away from in-progress) and task:deleted listeners populate
|
||||||
* this so a fast re-dispatch (task:moved → in-progress) awaits the prior
|
* this so a fast re-dispatch (task:moved → in-progress) awaits the prior
|
||||||
@@ -3045,7 +3046,7 @@ export class TaskExecutor {
|
|||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
} else if (from === "in-progress") {
|
} else if (from === "in-progress") {
|
||||||
if (this.workflowBoundaryMovesInFlight.has(task.id) && this.graphRouting.has(task.id)) {
|
if (this.workflowLifecycleMovesInFlight.has(task.id) && this.graphRouting.has(task.id)) {
|
||||||
executorLog.log(
|
executorLog.log(
|
||||||
`[event:task:moved] Preserving graph run for ${task.id} across its own ${from} → ${to} boundary`,
|
`[event:task:moved] Preserving graph run for ${task.id} across its own ${from} → ${to} boundary`,
|
||||||
);
|
);
|
||||||
@@ -4570,13 +4571,32 @@ export class TaskExecutor {
|
|||||||
}
|
}
|
||||||
await this.persistTokenUsage(task.id);
|
await this.persistTokenUsage(task.id);
|
||||||
const originColumn = task.column;
|
const originColumn = task.column;
|
||||||
const promotedFromTodo = originColumn === "todo";
|
const promotedFromPlannerColumn = originColumn === "todo" || originColumn === "triage";
|
||||||
if (promotedFromTodo) {
|
let completionTask = task;
|
||||||
|
if (promotedFromPlannerColumn) {
|
||||||
this.recoveringCompleted.add(task.id);
|
this.recoveringCompleted.add(task.id);
|
||||||
await this.store.moveTask(task.id, "in-progress");
|
/*
|
||||||
|
FNXC:WorkflowLifecycle 2026-07-20-08:42:
|
||||||
|
Advanced-triage recovery reaches this shared seam with completed work, a
|
||||||
|
preserved worktree, and a durable merge pin. The workflow transition map
|
||||||
|
deliberately rejects triage -> in-review, so re-home through the legal
|
||||||
|
triage -> todo -> in-progress path while the recovery ownership set prevents
|
||||||
|
scheduler/executor dispatch. Todo callers retain their existing single hop.
|
||||||
|
*/
|
||||||
|
if (originColumn === "triage") {
|
||||||
|
completionTask = await this.store.moveTask(task.id, "todo", {
|
||||||
|
moveSource: "engine",
|
||||||
|
recoveryRehome: true,
|
||||||
|
bypassGuards: true,
|
||||||
|
preserveProgress: true,
|
||||||
|
preserveWorktree: true,
|
||||||
|
preserveResumeState: true,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
completionTask = await this.store.moveTask(task.id, "in-progress");
|
||||||
}
|
}
|
||||||
await this.handoffTaskToReview(task, "completed-task-recovered");
|
await this.handoffTaskToReview(completionTask, "completed-task-recovered");
|
||||||
if (promotedFromTodo) {
|
if (promotedFromPlannerColumn) {
|
||||||
this.recoveringCompleted.delete(task.id);
|
this.recoveringCompleted.delete(task.id);
|
||||||
}
|
}
|
||||||
this.clearCompletedTaskWatchdog(task.id);
|
this.clearCompletedTaskWatchdog(task.id);
|
||||||
@@ -4757,7 +4777,12 @@ export class TaskExecutor {
|
|||||||
optionalStepRevisionLogOutcome(feedback, revisionKey),
|
optionalStepRevisionLogOutcome(feedback, revisionKey),
|
||||||
this.getRunContextFor(taskId),
|
this.getRunContextFor(taskId),
|
||||||
);
|
);
|
||||||
await moveTaskToReplanColumn(this.store, { id: taskId, column: liveTask.column }, replanColumn);
|
this.workflowLifecycleMovesInFlight.add(taskId);
|
||||||
|
try {
|
||||||
|
await moveTaskToReplanColumn(this.store, { id: taskId, column: liveTask.column }, replanColumn);
|
||||||
|
} finally {
|
||||||
|
this.workflowLifecycleMovesInFlight.delete(taskId);
|
||||||
|
}
|
||||||
await this.store.updateTask(taskId, {
|
await this.store.updateTask(taskId, {
|
||||||
status: "needs-replan",
|
status: "needs-replan",
|
||||||
error: null,
|
error: null,
|
||||||
@@ -5836,7 +5861,7 @@ export class TaskExecutor {
|
|||||||
// row fields so an ordinary requeue re-resolves the CURRENT IR fresh.
|
// row fields so an ordinary requeue re-resolves the CURRENT IR fresh.
|
||||||
clearPin: pinPersistence.clearPin,
|
clearPin: pinPersistence.clearPin,
|
||||||
moveTask: async (toColumn, ctx) => {
|
moveTask: async (toColumn, ctx) => {
|
||||||
this.workflowBoundaryMovesInFlight.add(task.id);
|
this.workflowLifecycleMovesInFlight.add(task.id);
|
||||||
try {
|
try {
|
||||||
await this.store.moveTask(task.id, toColumn, {
|
await this.store.moveTask(task.id, toColumn, {
|
||||||
moveSource: "engine",
|
moveSource: "engine",
|
||||||
@@ -5846,7 +5871,7 @@ export class TaskExecutor {
|
|||||||
workflowMoveMetadata: { fromColumn: ctx.fromColumn, nodeId: ctx.nodeId },
|
workflowMoveMetadata: { fromColumn: ctx.fromColumn, nodeId: ctx.nodeId },
|
||||||
});
|
});
|
||||||
} finally {
|
} finally {
|
||||||
this.workflowBoundaryMovesInFlight.delete(task.id);
|
this.workflowLifecycleMovesInFlight.delete(task.id);
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
emitAudit: async (event) => {
|
emitAudit: async (event) => {
|
||||||
|
|||||||
@@ -3028,6 +3028,11 @@ export class SelfHealingManager {
|
|||||||
const tasks = await this.store.listTasks({ column: "triage", slim: true });
|
const tasks = await this.store.listTasks({ column: "triage", slim: true });
|
||||||
const executingIds = this.options.getExecutingTaskIds?.() ?? new Set<string>();
|
const executingIds = this.options.getExecutingTaskIds?.() ?? new Set<string>();
|
||||||
const planningIds = this.options.getPlanningTaskIds?.() ?? new Set<string>();
|
const planningIds = this.options.getPlanningTaskIds?.() ?? new Set<string>();
|
||||||
|
const hasForeignPathOwner = (task: Task) => {
|
||||||
|
if (!task.worktree) return false;
|
||||||
|
const owner = activeSessionRegistry.lookupByPath(task.worktree);
|
||||||
|
return owner != null && owner.taskId !== task.id;
|
||||||
|
};
|
||||||
const candidates = tasks.filter((task) =>
|
const candidates = tasks.filter((task) =>
|
||||||
task.column === "triage"
|
task.column === "triage"
|
||||||
&& task.status == null
|
&& task.status == null
|
||||||
@@ -3038,7 +3043,8 @@ export class SelfHealingManager {
|
|||||||
&& !executingIds.has(task.id)
|
&& !executingIds.has(task.id)
|
||||||
&& !planningIds.has(task.id)
|
&& !planningIds.has(task.id)
|
||||||
&& this.options.isTaskActive?.(task.id) !== true
|
&& this.options.isTaskActive?.(task.id) !== true
|
||||||
&& !(task.worktree && activeSessionRegistry.isPathActive(task.worktree)),
|
&& this.options.getActiveMergeTaskId?.() !== task.id
|
||||||
|
&& !hasForeignPathOwner(task),
|
||||||
);
|
);
|
||||||
|
|
||||||
let recovered = 0;
|
let recovered = 0;
|
||||||
@@ -3056,11 +3062,21 @@ export class SelfHealingManager {
|
|||||||
|| !live.workflowIrPinNodeId
|
|| !live.workflowIrPinNodeId
|
||||||
|| (this.options.getExecutingTaskIds?.() ?? new Set<string>()).has(live.id)
|
|| (this.options.getExecutingTaskIds?.() ?? new Set<string>()).has(live.id)
|
||||||
|| this.options.isTaskActive?.(live.id) === true
|
|| this.options.isTaskActive?.(live.id) === true
|
||||||
|| activeSessionRegistry.isPathActive(live.worktree)
|
|| this.options.getActiveMergeTaskId?.() === live.id
|
||||||
|
|| hasForeignPathOwner(live)
|
||||||
) {
|
) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A prior aborted graph can leave its own registry claim behind after every
|
||||||
|
// executable/planning/merge owner has gone away. The durable task row and the
|
||||||
|
// liveness callbacks above prove that claim is stale; clear it so completed
|
||||||
|
// recovery can reuse the preserved worktree instead of rejecting its own path.
|
||||||
|
const pathOwner = activeSessionRegistry.lookupByPath(live.worktree);
|
||||||
|
if (pathOwner?.taskId === live.id) {
|
||||||
|
activeSessionRegistry.unregisterPath(live.worktree);
|
||||||
|
}
|
||||||
|
|
||||||
const steps = live.steps ?? [];
|
const steps = live.steps ?? [];
|
||||||
const complete = steps.length > 0
|
const complete = steps.length > 0
|
||||||
&& steps.every((step) => step.status === "done" || step.status === "skipped");
|
&& steps.every((step) => step.status === "done" || step.status === "skipped");
|
||||||
@@ -3083,7 +3099,8 @@ export class SelfHealingManager {
|
|||||||
&& !(this.options.getExecutingTaskIds?.() ?? new Set<string>()).has(current.id)
|
&& !(this.options.getExecutingTaskIds?.() ?? new Set<string>()).has(current.id)
|
||||||
&& !(this.options.getPlanningTaskIds?.() ?? new Set<string>()).has(current.id)
|
&& !(this.options.getPlanningTaskIds?.() ?? new Set<string>()).has(current.id)
|
||||||
&& this.options.isTaskActive?.(current.id) !== true
|
&& this.options.isTaskActive?.(current.id) !== true
|
||||||
&& !(current.worktree && activeSessionRegistry.isPathActive(current.worktree)),
|
&& this.options.getActiveMergeTaskId?.() !== current.id
|
||||||
|
&& !hasForeignPathOwner(current),
|
||||||
{
|
{
|
||||||
moveSource: "engine",
|
moveSource: "engine",
|
||||||
workflowMoveSource: "self-healing-advanced-triage",
|
workflowMoveSource: "self-healing-advanced-triage",
|
||||||
|
|||||||
Reference in New Issue
Block a user