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:
gsxdsm
2026-07-20 08:46:34 -07:00
parent 8166c80bad
commit 71c0d0a970
6 changed files with 167 additions and 16 deletions

View 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.

View File

@@ -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

View File

@@ -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");

View File

@@ -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" });

View File

@@ -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) => {

View File

@@ -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",