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):
|
||||
"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();
|
||||
});
|
||||
|
||||
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 () => {
|
||||
const ordinary = task("FN-ORDINARY", { worktree: undefined, workflowIrPinNodeId: undefined });
|
||||
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);
|
||||
});
|
||||
|
||||
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 () => {
|
||||
const zeroStore = createMockStore();
|
||||
const zeroTask = task({ postReviewFixCount: 0, column: "in-progress" });
|
||||
|
||||
@@ -1662,11 +1662,12 @@ export class TaskExecutor {
|
||||
private workflowRerunPending = new Set<string>();
|
||||
/**
|
||||
* Task ids whose current `task:moved` event is being emitted by this
|
||||
* executor's workflow column-boundary hook. The store emits synchronously,
|
||||
* so this narrowly distinguishes a graph's own transition from an external
|
||||
* engine/user move that must still hard-cancel the active run.
|
||||
* executor's workflow lifecycle handling (column boundaries or Plan Review
|
||||
* replans). The store emits synchronously, so this narrowly distinguishes a
|
||||
* 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
|
||||
* task:moved (away from in-progress) and task:deleted listeners populate
|
||||
* 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") {
|
||||
if (this.workflowBoundaryMovesInFlight.has(task.id) && this.graphRouting.has(task.id)) {
|
||||
if (this.workflowLifecycleMovesInFlight.has(task.id) && this.graphRouting.has(task.id)) {
|
||||
executorLog.log(
|
||||
`[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);
|
||||
const originColumn = task.column;
|
||||
const promotedFromTodo = originColumn === "todo";
|
||||
if (promotedFromTodo) {
|
||||
const promotedFromPlannerColumn = originColumn === "todo" || originColumn === "triage";
|
||||
let completionTask = task;
|
||||
if (promotedFromPlannerColumn) {
|
||||
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");
|
||||
if (promotedFromTodo) {
|
||||
await this.handoffTaskToReview(completionTask, "completed-task-recovered");
|
||||
if (promotedFromPlannerColumn) {
|
||||
this.recoveringCompleted.delete(task.id);
|
||||
}
|
||||
this.clearCompletedTaskWatchdog(task.id);
|
||||
@@ -4757,7 +4777,12 @@ export class TaskExecutor {
|
||||
optionalStepRevisionLogOutcome(feedback, revisionKey),
|
||||
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, {
|
||||
status: "needs-replan",
|
||||
error: null,
|
||||
@@ -5836,7 +5861,7 @@ export class TaskExecutor {
|
||||
// row fields so an ordinary requeue re-resolves the CURRENT IR fresh.
|
||||
clearPin: pinPersistence.clearPin,
|
||||
moveTask: async (toColumn, ctx) => {
|
||||
this.workflowBoundaryMovesInFlight.add(task.id);
|
||||
this.workflowLifecycleMovesInFlight.add(task.id);
|
||||
try {
|
||||
await this.store.moveTask(task.id, toColumn, {
|
||||
moveSource: "engine",
|
||||
@@ -5846,7 +5871,7 @@ export class TaskExecutor {
|
||||
workflowMoveMetadata: { fromColumn: ctx.fromColumn, nodeId: ctx.nodeId },
|
||||
});
|
||||
} finally {
|
||||
this.workflowBoundaryMovesInFlight.delete(task.id);
|
||||
this.workflowLifecycleMovesInFlight.delete(task.id);
|
||||
}
|
||||
},
|
||||
emitAudit: async (event) => {
|
||||
|
||||
@@ -3028,6 +3028,11 @@ export class SelfHealingManager {
|
||||
const tasks = await this.store.listTasks({ column: "triage", slim: true });
|
||||
const executingIds = this.options.getExecutingTaskIds?.() ?? 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) =>
|
||||
task.column === "triage"
|
||||
&& task.status == null
|
||||
@@ -3038,7 +3043,8 @@ export class SelfHealingManager {
|
||||
&& !executingIds.has(task.id)
|
||||
&& !planningIds.has(task.id)
|
||||
&& this.options.isTaskActive?.(task.id) !== true
|
||||
&& !(task.worktree && activeSessionRegistry.isPathActive(task.worktree)),
|
||||
&& this.options.getActiveMergeTaskId?.() !== task.id
|
||||
&& !hasForeignPathOwner(task),
|
||||
);
|
||||
|
||||
let recovered = 0;
|
||||
@@ -3056,11 +3062,21 @@ export class SelfHealingManager {
|
||||
|| !live.workflowIrPinNodeId
|
||||
|| (this.options.getExecutingTaskIds?.() ?? new Set<string>()).has(live.id)
|
||||
|| this.options.isTaskActive?.(live.id) === true
|
||||
|| activeSessionRegistry.isPathActive(live.worktree)
|
||||
|| this.options.getActiveMergeTaskId?.() === live.id
|
||||
|| hasForeignPathOwner(live)
|
||||
) {
|
||||
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 complete = steps.length > 0
|
||||
&& 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.getPlanningTaskIds?.() ?? new Set<string>()).has(current.id)
|
||||
&& this.options.isTaskActive?.(current.id) !== true
|
||||
&& !(current.worktree && activeSessionRegistry.isPathActive(current.worktree)),
|
||||
&& this.options.getActiveMergeTaskId?.() !== current.id
|
||||
&& !hasForeignPathOwner(current),
|
||||
{
|
||||
moveSource: "engine",
|
||||
workflowMoveSource: "self-healing-advanced-triage",
|
||||
|
||||
Reference in New Issue
Block a user