diff --git a/packages/engine/src/__tests__/reliability-interactions/completion-handoff-limbo.test.ts b/packages/engine/src/__tests__/reliability-interactions/completion-handoff-limbo.test.ts index 7353eae60c..1c91fcc2f9 100644 --- a/packages/engine/src/__tests__/reliability-interactions/completion-handoff-limbo.test.ts +++ b/packages/engine/src/__tests__/reliability-interactions/completion-handoff-limbo.test.ts @@ -18,6 +18,10 @@ function makeTask(overrides: Partial = {}): Task { } as Task; } +function doneMarker(minutesAgo = 6) { + return { action: "Task marked done by agent", timestamp: new Date(Date.now() - minutesAgo * 60_000).toISOString() } as any; +} + function createStore(task: Task) { let current = { ...task } as Task; return { @@ -35,18 +39,22 @@ function createStore(task: Task) { } as any; } +function limboTask(overrides: Partial = {}): Task { + return makeTask({ + worktree: "/tmp/wt", + status: undefined, + review: undefined, + reviewState: undefined, + mergeDetails: undefined, + log: [doneMarker()], + ...overrides, + }); +} + describe("FN-4999 reliability interactions: completion-handoff-limbo", () => { - it("recovers exact signature by requeueing auto-merge", async () => { - const task = makeTask({ - worktree: "/tmp/wt", - status: undefined, - review: undefined, - reviewState: undefined, - mergeDetails: undefined, - log: [{ action: "Task marked done by agent", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any], - }); - const store = createStore(task); - const requeueForAutoMerge = vi.fn(); + it("recovers exact signature by requeueing auto-merge after accepted handoff", async () => { + const store = createStore(limboTask()); + const requeueForAutoMerge = vi.fn(() => true); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); await manager.recoverCompletionHandoffLimbo(); @@ -59,15 +67,19 @@ describe("FN-4999 reliability interactions: completion-handoff-limbo", () => { expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:auto-recover-completion-handoff-limbo", target: "FN-4999-T", - metadata: expect.objectContaining({ ageMs: expect.any(Number), source: "self-healing-in-review-sweep" }), + metadata: expect.objectContaining({ + ageMs: expect.any(Number), + source: "self-healing-in-review-sweep", + attempts: 1, + }), })); const event = store.recordRunAuditEvent.mock.calls.find((call: any[]) => call[0].mutationType === "task:auto-recover-completion-handoff-limbo")?.[0]; expect(event.metadata.ageMs).toBeGreaterThanOrEqual(COMPLETION_HANDOFF_LIMBO_GRACE_MS); }); it("is no-op before grace period elapses", async () => { - const store = createStore(makeTask({ status: undefined, review: undefined, reviewState: undefined, mergeDetails: undefined, log: [{ action: "Task marked done by agent", timestamp: new Date(Date.now() - 30_000).toISOString() } as any] })); - const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge: vi.fn() }); + const store = createStore(limboTask({ log: [doneMarker(0.5)] })); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge: vi.fn(() => true) }); await manager.recoverCompletionHandoffLimbo(); expect(store.updateTask).not.toHaveBeenCalled(); expect(store.logEntry).not.toHaveBeenCalled(); @@ -75,32 +87,32 @@ describe("FN-4999 reliability interactions: completion-handoff-limbo", () => { }); it("skips active tasks", async () => { - const requeueForAutoMerge = vi.fn(); - const store = createStore(makeTask({ status: undefined, review: undefined, reviewState: undefined, mergeDetails: undefined, log: [{ action: "Task marked done by agent", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any] })); + const requeueForAutoMerge = vi.fn(() => true); + const store = createStore(limboTask()); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge, isTaskActive: () => true }); await manager.recoverCompletionHandoffLimbo(); expect(requeueForAutoMerge).not.toHaveBeenCalled(); }); it("honors legitimate merge blockers", async () => { - const requeueForAutoMerge = vi.fn(); - const store = createStore(makeTask({ status: "failed", review: undefined, reviewState: undefined, mergeDetails: undefined, log: [{ action: "Task marked done by agent", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any] })); + const requeueForAutoMerge = vi.fn(() => true); + const store = createStore(limboTask({ status: "failed" })); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); await manager.recoverCompletionHandoffLimbo(); expect(requeueForAutoMerge).not.toHaveBeenCalled(); }); it("is no-op when marker is absent", async () => { - const requeueForAutoMerge = vi.fn(); - const store = createStore(makeTask({ status: undefined, review: undefined, reviewState: undefined, mergeDetails: undefined, log: [{ action: "workflow step", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any] })); + const requeueForAutoMerge = vi.fn(() => true); + const store = createStore(limboTask({ log: [{ action: "workflow step", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any] })); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); await manager.recoverCompletionHandoffLimbo(); expect(requeueForAutoMerge).not.toHaveBeenCalled(); }); it("emits exhausted event and fails task at cap", async () => { - const requeueForAutoMerge = vi.fn(); - const store = createStore(makeTask({ completionHandoffLimboRecoveryCount: MAX_COMPLETION_HANDOFF_LIMBO_RECOVERIES, status: undefined, review: undefined, reviewState: undefined, mergeDetails: undefined, log: [{ action: "Task marked done by agent", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any] })); + const requeueForAutoMerge = vi.fn(() => true); + const store = createStore(limboTask({ completionHandoffLimboRecoveryCount: MAX_COMPLETION_HANDOFF_LIMBO_RECOVERIES })); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); await manager.recoverCompletionHandoffLimbo(); expect(requeueForAutoMerge).not.toHaveBeenCalled(); @@ -109,10 +121,9 @@ describe("FN-4999 reliability interactions: completion-handoff-limbo", () => { expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:auto-recover-completion-handoff-limbo-exhausted" })); }); - it("increments completionHandoffLimboRecoveryCount on each successful recovery", async () => { - const task = makeTask({ status: undefined, review: undefined, reviewState: undefined, mergeDetails: undefined, log: [{ action: "Task marked done by agent", timestamp: new Date(Date.now() - 6 * 60_000).toISOString() } as any] }); - const store = createStore(task); - const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge: vi.fn() }); + it("increments completionHandoffLimboRecoveryCount on each accepted recovery", async () => { + const store = createStore(limboTask()); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge: vi.fn(() => true) }); await manager.recoverCompletionHandoffLimbo(); await manager.recoverCompletionHandoffLimbo(); @@ -123,4 +134,67 @@ describe("FN-4999 reliability interactions: completion-handoff-limbo", () => { .filter((value: unknown) => typeof value === "number"); expect(increments).toEqual([1, 2, 3]); }); + + it("does not increment completionHandoffLimboRecoveryCount when merge requeue is rejected", async () => { + const store = createStore(limboTask({ completionHandoffLimboRecoveryCount: 2 })); + const requeueForAutoMerge = vi.fn(() => false); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); + + await manager.recoverCompletionHandoffLimbo(); + + expect(requeueForAutoMerge).toHaveBeenCalledWith("FN-4999-T"); + expect(store._get().completionHandoffLimboRecoveryCount).toBe(2); + expect(store._get().status).toBeUndefined(); + expect(store.logEntry).not.toHaveBeenCalled(); + expect(store.recordRunAuditEvent).not.toHaveBeenCalledWith(expect.objectContaining({ + mutationType: "task:auto-recover-completion-handoff-limbo", + })); + }); + + it("increments completionHandoffLimboRecoveryCount when merge requeue is accepted", async () => { + const store = createStore(limboTask({ completionHandoffLimboRecoveryCount: 1 })); + const requeueForAutoMerge = vi.fn(() => true); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); + + await manager.recoverCompletionHandoffLimbo(); + + expect(store._get().completionHandoffLimboRecoveryCount).toBe(2); + expect(store.logEntry).toHaveBeenCalledWith("FN-4999-T", expect.stringMatching(/Auto-recovered \(FN-4999\)/)); + expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ + mutationType: "task:auto-recover-completion-handoff-limbo", + metadata: expect.objectContaining({ attempts: 2 }), + })); + }); + + it("skips tasks already held by the merge queue without incrementing the count", async () => { + const store = createStore(limboTask({ completionHandoffLimboRecoveryCount: 1 })); + const requeueForAutoMerge = vi.fn(() => false); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); + + await manager.recoverCompletionHandoffLimbo(); + await manager.recoverCompletionHandoffLimbo(); + + expect(requeueForAutoMerge).toHaveBeenCalledTimes(2); + expect(store._get().completionHandoffLimboRecoveryCount).toBe(1); + expect(store._get().status).toBeUndefined(); + expect(store.logEntry).not.toHaveBeenCalled(); + }); + + it("exhausts only after three accepted limbo recoveries", async () => { + const store = createStore(limboTask()); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge: vi.fn(() => true) }); + + await manager.recoverCompletionHandoffLimbo(); + await manager.recoverCompletionHandoffLimbo(); + await manager.recoverCompletionHandoffLimbo(); + expect(store._get().completionHandoffLimboRecoveryCount).toBe(MAX_COMPLETION_HANDOFF_LIMBO_RECOVERIES); + expect(store._get().status).toBeUndefined(); + + await manager.recoverCompletionHandoffLimbo(); + + expect(store.updateTask).toHaveBeenLastCalledWith("FN-4999-T", expect.objectContaining({ + status: "failed", + error: "Completion handoff limbo recovery exhausted", + })); + }); }); diff --git a/packages/engine/src/__tests__/reliability-interactions/in-review-stall-deadlock-disposition.test.ts b/packages/engine/src/__tests__/reliability-interactions/in-review-stall-deadlock-disposition.test.ts index da8407d468..1fe1cccce0 100644 --- a/packages/engine/src/__tests__/reliability-interactions/in-review-stall-deadlock-disposition.test.ts +++ b/packages/engine/src/__tests__/reliability-interactions/in-review-stall-deadlock-disposition.test.ts @@ -28,6 +28,7 @@ function createStore(task: Task, settings: Record = {}): TaskSt auditEvents.push(event); }); (emitter as any).moveTask = vi.fn().mockResolvedValue(undefined); + (emitter as any).enqueueMergeQueue = vi.fn().mockResolvedValue(undefined); return emitter; } @@ -114,6 +115,48 @@ describe("reliability interactions: in-review stall deadlock disposition", () => manager.stop(); }); + it("FN-6070: rejected limbo requeues do not increment into deadlock disposition", async () => { + const task = { + id: "FN-6070-REJECTED", + column: "in-review", + paused: false, + userPaused: false, + status: undefined, + error: undefined, + branch: "fusion/fn-6070-rejected", + worktree: "/tmp/fn-6070-rejected", + mergeDetails: undefined, + mergeRetries: 0, + completionHandoffLimboRecoveryCount: 2, + steps: [{ name: "merge", status: "done" }], + workflowStepResults: [], + updatedAt: "2026-01-01T00:00:00.000Z", + log: [{ action: "Task marked done by agent", timestamp: "2026-01-01T00:00:00.000Z" }], + } as any satisfies Task; + + const store = createStore(task); + const manager = new SelfHealingManager(store, { + rootDir: "/tmp/repo", + requeueForAutoMerge: vi.fn(() => false), + }); + + vi.setSystemTime(new Date("2026-01-01T00:10:00.000Z")); + await manager.recoverCompletionHandoffLimbo(); + vi.setSystemTime(new Date("2026-01-01T00:25:00.000Z")); + await manager.recoverCompletionHandoffLimbo(); + vi.setSystemTime(new Date("2026-01-01T00:40:00.000Z")); + await manager.recoverCompletionHandoffLimbo(); + + expect(task.completionHandoffLimboRecoveryCount).toBe(2); + expect(task.status).toBeUndefined(); + expect(task.error).toBeUndefined(); + expect(await manager.surfaceInReviewStalls()).toBe(0); + expect(task.pausedReason).not.toBe("in-review-stall-deadlock"); + expect(task.log.some((entry: { action: string }) => entry.action.includes("Completion handoff limbo recovery exhausted"))).toBe(false); + + manager.stop(); + }); + it("does not auto-dispose userPaused tasks with repeated identical stalls", async () => { const task = { id: "FN-4860-PAUSED", diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index 33c2804279..40ed13d67f 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -788,7 +788,7 @@ export class InProcessRuntime getPlanningTaskIds: () => this.triageProcessor?.getProcessingTaskIds() ?? new Set(), evictStaleTriageProcessing: () => this.triageProcessor?.evictStaleProcessing() ?? new Set(), enqueueMerge: this.mergeEnqueuer ? (taskId: string) => this.mergeEnqueuer?.(taskId) ?? false : undefined, - requeueForAutoMerge: this.mergeEnqueuer ? (taskId: string) => { this.mergeEnqueuer?.(taskId); } : undefined, + requeueForAutoMerge: this.mergeEnqueuer ? (taskId: string) => this.mergeEnqueuer?.(taskId) ?? false : undefined, isTaskActive: (taskId: string) => this.executor.isTaskActive(taskId), clearMergeActive: this.clearMergeActive ? (taskId: string) => this.clearMergeActive?.(taskId) : undefined, getActiveMergeTaskId: () => this.activeMergeTaskIdProvider?.() ?? null, diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index 5d4483e6be..20c6902adb 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -269,7 +269,7 @@ export interface SelfHealingOptions { * the polling sweep's enqueue to silently no-op). */ enqueueMerge?: (taskId: string) => boolean; - requeueForAutoMerge?: (taskId: string) => void | Promise; + requeueForAutoMerge?: (taskId: string) => boolean | void | Promise; isTaskActive?: (taskId: string) => boolean; clearMergeActive?: (taskId: string) => void; /** @@ -7002,6 +7002,30 @@ export class SelfHealingManager { continue; } + const requeueForAutoMerge = this.options.requeueForAutoMerge ?? this.options.enqueueMerge; + if (!requeueForAutoMerge) { + log.warn(`recoverCompletionHandoffLimbo: requeueForAutoMerge callback missing for ${task.id}`); + continue; + } + + try { + // FN-5353: strict targetTaskId leasing in reuse handoff requires an + // explicit queue row before re-emitting auto-merge. + await this.store.enqueueMergeQueue(task.id); + } catch (err) { + const errorMessage = err instanceof Error ? err.message : String(err); + log.warn(`recoverCompletionHandoffLimbo: enqueue failed for ${task.id}: ${errorMessage}`); + continue; + } + + const accepted = await requeueForAutoMerge(task.id); + if (accepted !== true) { + log.log( + `recoverCompletionHandoffLimbo: skipped recovery count for ${task.id} because merge requeue was not accepted`, + ); + continue; + } + await this.store.updateTask(task.id, { completionHandoffLimboRecoveryCount: currentCount + 1, }); @@ -7016,24 +7040,10 @@ export class SelfHealingManager { await audit.database({ type: "task:auto-recover-completion-handoff-limbo", target: task.id, - metadata: { ageMs, source: "self-healing-in-review-sweep" }, + metadata: { ageMs, source: "self-healing-in-review-sweep", attempts: currentCount + 1 }, }); await this.store.logEntry(task.id, "Auto-recovered (FN-4999): task in 'in-review' past handoff grace with no merge fan-out — re-emitting auto-merge handoff"); - if (this.options.requeueForAutoMerge) { - try { - // FN-5353: strict targetTaskId leasing in reuse handoff requires an - // explicit queue row before re-emitting auto-merge. - await this.store.enqueueMergeQueue(task.id); - } catch (err) { - const errorMessage = err instanceof Error ? err.message : String(err); - log.warn(`recoverCompletionHandoffLimbo: enqueue failed for ${task.id}: ${errorMessage}`); - continue; - } - await this.options.requeueForAutoMerge(task.id); - } else { - log.warn(`recoverCompletionHandoffLimbo: requeueForAutoMerge callback missing for ${task.id}`); - } } }