diff --git a/.changeset/merge-queue-stall-badges.md b/.changeset/merge-queue-stall-badges.md new file mode 100644 index 0000000000..fe7fa41554 --- /dev/null +++ b/.changeset/merge-queue-stall-badges.md @@ -0,0 +1,4 @@ +"@runfusion/fusion": patch +--- + +Suppress in-review stall and merge-stalled signals for tasks already owned by the merge queue. diff --git a/packages/core/src/__tests__/store-in-review-stall.test.ts b/packages/core/src/__tests__/store-in-review-stall.test.ts index e5292c19cd..a417c5eac2 100644 --- a/packages/core/src/__tests__/store-in-review-stall.test.ts +++ b/packages/core/src/__tests__/store-in-review-stall.test.ts @@ -52,6 +52,27 @@ describe("TaskStore inReviewStall hydration", () => { expect(task?.inReviewStall?.reason).toContain("no active merger"); }); + it("omits merge-stalled hydration while the task is already queued for merge", async () => { + await seedTask("FN-6088", {}); + await store.enqueueMergeQueue("FN-6088"); + + const listed = (await store.listTasks({ slim: true })).find((entry) => entry.id === "FN-6088"); + expect(listed?.inReviewStall).toBeUndefined(); + expect(listed?.inReviewStalled).toBeUndefined(); + + const detailed = await store.getTask("FN-6088"); + expect(detailed.inReviewStall).toBeUndefined(); + expect(detailed.inReviewStalled).toBeUndefined(); + + const modified = (await store.listTasksModifiedSince("1970-01-01T00:00:00.000Z")).tasks.find((entry) => entry.id === "FN-6088"); + expect(modified?.inReviewStall).toBeUndefined(); + expect(modified?.inReviewStalled).toBeUndefined(); + + const searched = (await store.searchTasks("FN-6088", { slim: true })).find((entry) => entry.id === "FN-6088"); + expect(searched?.inReviewStall).toBeUndefined(); + expect(searched?.inReviewStalled).toBeUndefined(); + }); + it("omits inReviewStall for paused in-review task", async () => { await seedTask("FN-4217-PAUSED", { paused: true }); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index c35cda22cb..7dff4a89dd 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -2738,6 +2738,11 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} return this.rowToTask(row); } + private getMergeQueuedTaskIds(): Set { + const rows = this.db.prepare("SELECT taskId FROM mergeQueue").all() as Array<{ taskId: string }>; + return new Set(rows.map((row) => row.taskId)); + } + private isTaskIdPresentInArchivedTasksTable(id: string): boolean { try { const row = this.db.prepare("SELECT 1 as found FROM archivedTasks WHERE id = ? LIMIT 1").get(id) as { found?: number } | undefined; @@ -4808,7 +4813,27 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} }; } - task.stalledReview = detectStalledReview(task, { now: Date.now() }); + const now = Date.now(); + const settings = await this.getSettingsFast(); + const mergeQueuedTaskIds = this.getMergeQueuedTaskIds(); + task.inReviewStall = mergeQueuedTaskIds.has(task.id) + ? undefined + : getInReviewStallReason(task, { + now, + autoMerge: allowsAutoMergeProcessing(task, settings), + engineActiveSinceMs: settings.engineActiveSinceMs, + engineActivationGraceMs: settings.engineActivationGraceMs, + }); + task.inReviewStalled = mergeQueuedTaskIds.has(task.id) + ? undefined + : getInReviewStalledSignal(task, { + now, + thresholdMs: settings.inReviewStalledThresholdMs, + autoMerge: allowsAutoMergeProcessing(task, settings), + engineActiveSinceMs: settings.engineActiveSinceMs, + engineActivationGraceMs: settings.engineActivationGraceMs, + }); + task.stalledReview = detectStalledReview(task, { now }); // Derived at read time only; retrySummary is never persisted to SQLite. task.retrySummary = computeRetrySummary(task); @@ -5311,9 +5336,11 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} inReviewCriticalMs: settings.staleInReviewCriticalMs, }; let disableAgeStalenessHydration = false; + const mergeQueuedTaskIds = this.getMergeQueuedTaskIds(); const activeTasks = await Promise.all((rows as unknown as TaskRow[]).map(async (row) => { const task = this.rowToTask(row); - task.inReviewStall = getInReviewStallReason(task, { + const isMergeQueued = mergeQueuedTaskIds.has(task.id); + task.inReviewStall = isMergeQueued ? undefined : getInReviewStallReason(task, { now, autoMerge: allowsAutoMergeProcessing(task, settings), engineActiveSinceMs: settings.engineActiveSinceMs, @@ -5325,7 +5352,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} engineActiveSinceMs: settings.engineActiveSinceMs, engineActivationGraceMs: settings.engineActivationGraceMs, }); - task.inReviewStalled = getInReviewStalledSignal(task, { + task.inReviewStalled = isMergeQueued ? undefined : getInReviewStalledSignal(task, { now, thresholdMs: settings.inReviewStalledThresholdMs, autoMerge: allowsAutoMergeProcessing(task, settings), @@ -5831,9 +5858,11 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} inReviewCriticalMs: settings.staleInReviewCriticalMs, }; let disableAgeStalenessHydration = false; + const mergeQueuedTaskIds = this.getMergeQueuedTaskIds(); const tasks = rows.slice(0, resolvedLimit).map((row) => { const task = this.rowToTask(row); - task.inReviewStall = getInReviewStallReason(task, { + const isMergeQueued = mergeQueuedTaskIds.has(task.id); + task.inReviewStall = isMergeQueued ? undefined : getInReviewStallReason(task, { now, autoMerge: allowsAutoMergeProcessing(task, settings), engineActiveSinceMs: settings.engineActiveSinceMs, @@ -5845,7 +5874,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} engineActiveSinceMs: settings.engineActiveSinceMs, engineActivationGraceMs: settings.engineActivationGraceMs, }); - task.inReviewStalled = getInReviewStalledSignal(task, { + task.inReviewStalled = isMergeQueued ? undefined : getInReviewStalledSignal(task, { now, thresholdMs: settings.inReviewStalledThresholdMs, autoMerge: allowsAutoMergeProcessing(task, settings), @@ -5994,9 +6023,11 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} inReviewCriticalMs: settings.staleInReviewCriticalMs, }; let disableAgeStalenessHydration = false; + const mergeQueuedTaskIds = this.getMergeQueuedTaskIds(); const activeMatches = await Promise.all(rows.map(async (row) => { const task = this.rowToTask(row); - task.inReviewStall = getInReviewStallReason(task, { + const isMergeQueued = mergeQueuedTaskIds.has(task.id); + task.inReviewStall = isMergeQueued ? undefined : getInReviewStallReason(task, { now, autoMerge: allowsAutoMergeProcessing(task, settings), engineActiveSinceMs: settings.engineActiveSinceMs, @@ -6008,7 +6039,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} engineActiveSinceMs: settings.engineActiveSinceMs, engineActivationGraceMs: settings.engineActivationGraceMs, }); - task.inReviewStalled = getInReviewStalledSignal(task, { + task.inReviewStalled = isMergeQueued ? undefined : getInReviewStalledSignal(task, { now, thresholdMs: settings.inReviewStalledThresholdMs, autoMerge: allowsAutoMergeProcessing(task, settings), 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 1c91fcc2f9..706fd38dac 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 @@ -22,7 +22,7 @@ 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) { +function createStore(task: Task, mergeQueuedTaskIds: string[] = []) { let current = { ...task } as Task; return { getSettings: vi.fn(async () => ({ globalPause: false, enginePaused: false })), @@ -33,6 +33,16 @@ function createStore(task: Task) { }), moveTask: vi.fn(async () => undefined), enqueueMergeQueue: vi.fn(async () => undefined), + peekMergeQueue: vi.fn(() => mergeQueuedTaskIds.map((taskId) => ({ + taskId, + enqueuedAt: new Date().toISOString(), + priority: "normal", + leasedBy: null, + leasedAt: null, + leaseExpiresAt: null, + attemptCount: 0, + lastError: null, + }))), logEntry: vi.fn(async () => undefined), recordRunAuditEvent: vi.fn(async () => undefined), _get: () => current, @@ -167,19 +177,42 @@ describe("FN-4999 reliability interactions: completion-handoff-limbo", () => { }); it("skips tasks already held by the merge queue without incrementing the count", async () => { - const store = createStore(limboTask({ completionHandoffLimboRecoveryCount: 1 })); + const store = createStore(limboTask({ completionHandoffLimboRecoveryCount: 1 }), ["FN-4999-T"]); const requeueForAutoMerge = vi.fn(() => false); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); await manager.recoverCompletionHandoffLimbo(); await manager.recoverCompletionHandoffLimbo(); - expect(requeueForAutoMerge).toHaveBeenCalledTimes(2); + expect(requeueForAutoMerge).not.toHaveBeenCalled(); + expect(store.enqueueMergeQueue).not.toHaveBeenCalled(); expect(store._get().completionHandoffLimboRecoveryCount).toBe(1); expect(store._get().status).toBeUndefined(); expect(store.logEntry).not.toHaveBeenCalled(); }); + it("clears false handoff exhaustion for tasks already held by the merge queue", async () => { + const store = createStore(limboTask({ + status: "failed", + error: "Completion handoff limbo recovery exhausted", + completionHandoffLimboRecoveryCount: MAX_COMPLETION_HANDOFF_LIMBO_RECOVERIES, + }), ["FN-4999-T"]); + const requeueForAutoMerge = vi.fn(() => false); + const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge }); + + await manager.recoverCompletionHandoffLimbo(); + + expect(requeueForAutoMerge).not.toHaveBeenCalled(); + expect(store.enqueueMergeQueue).not.toHaveBeenCalled(); + expect(store._get().status).toBeNull(); + expect(store._get().error).toBeNull(); + expect(store._get().completionHandoffLimboRecoveryCount).toBe(0); + expect(store.logEntry).toHaveBeenCalledWith( + "FN-4999-T", + "Auto-recovered: cleared false completion-handoff exhaustion while task is already owned by merge queue", + ); + }); + it("exhausts only after three accepted limbo recoveries", async () => { const store = createStore(limboTask()); const manager = new SelfHealingManager(store, { rootDir: "/repo", requeueForAutoMerge: vi.fn(() => true) }); diff --git a/packages/engine/src/__tests__/reliability-interactions/in-review-stalled-detector.test.ts b/packages/engine/src/__tests__/reliability-interactions/in-review-stalled-detector.test.ts index b8adc73d3a..690fca7e42 100644 --- a/packages/engine/src/__tests__/reliability-interactions/in-review-stalled-detector.test.ts +++ b/packages/engine/src/__tests__/reliability-interactions/in-review-stalled-detector.test.ts @@ -3,7 +3,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import type { Task, TaskStore } from "@fusion/core"; import { SelfHealingManager } from "../../self-healing.js"; -function createStore(task: Task, settings: Record = {}): TaskStore & EventEmitter { +function createStore(task: Task, settings: Record = {}, mergeQueuedTaskIds: string[] = []): TaskStore & EventEmitter { const emitter = new EventEmitter() as TaskStore & EventEmitter; (emitter as any).getSettings = vi.fn().mockResolvedValue({ autoMerge: true, @@ -30,6 +30,16 @@ function createStore(task: Task, settings: Record = {}): TaskSt task.column = column as any; task.updatedAt = new Date(Date.now()).toISOString(); }); + (emitter as any).peekMergeQueue = vi.fn(() => mergeQueuedTaskIds.map((taskId) => ({ + taskId, + enqueuedAt: new Date().toISOString(), + priority: "normal", + leasedBy: null, + leasedAt: null, + leaseExpiresAt: null, + attemptCount: 0, + lastError: null, + }))); (emitter as any).recordRunAuditEvent = vi.fn().mockResolvedValue(undefined); return emitter; } @@ -80,6 +90,49 @@ describe("reliability interactions: in-review-stalled detector", () => { manager.stop(); }); + it("does not surface merge-stalled badges for tasks already queued in the merge lane", async () => { + const task = baseTask({ + id: "FN-6088", + status: "merging", + error: null, + updatedAt: "2026-01-01T00:00:00.000Z", + columnMovedAt: "2026-01-01T00:00:00.000Z", + }); + const store = createStore(task, { inReviewStalledThresholdMs: 3_600_000 }, ["FN-6088"]); + const manager = new SelfHealingManager(store, { rootDir: "/tmp/repo", getExecutingTaskIds: () => new Set() }); + + vi.setSystemTime(new Date("2026-01-01T06:00:00.000Z")); + + expect(await manager.surfaceInReviewStalls()).toBe(0); + expect(await manager.surfaceInReviewStalled()).toBe(0); + expect(await manager.recoverGhostReviewTasks()).toBe(0); + expect(task.column).toBe("in-review"); + expect(task.log?.some((entry) => entry.action.includes("stall"))).toBe(false); + + manager.stop(); + }); + + it("does not kick queued completed review tasks back to todo as ghost reviews", async () => { + const task = baseTask({ + id: "FN-6086", + status: null, + error: null, + steps: [], + inReviewStall: undefined, + worktree: "/tmp/wt", + }); + const store = createStore(task, { inReviewStalledThresholdMs: 3_600_000, taskStuckTimeoutMs: 12 * 3_600_000 }, ["FN-6086"]); + const manager = new SelfHealingManager(store, { rootDir: "/tmp/repo", getExecutingTaskIds: () => new Set() }); + + vi.setSystemTime(new Date("2026-01-01T13:00:00.000Z")); + + expect(await manager.recoverGhostReviewTasks()).toBe(0); + expect(await manager.surfaceInReviewStalled()).toBe(0); + expect(task.column).toBe("in-review"); + + manager.stop(); + }); + it("paused in-review tasks are owned by stale-paused-review detector", async () => { const task = baseTask({ paused: true, pausedReason: "manual-hold", status: "failed", error: null }); const store = createStore(task); diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index 20c6902adb..83224b8e81 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -663,6 +663,26 @@ export class SelfHealingManager { return this.options.getActiveMergeTaskId?.() ?? null; } + private isMergeLaneOwned(taskId: string): boolean { + if (this.options.getActiveMergeTaskId?.() === taskId) return true; + + try { + return this.store.peekMergeQueue().some((entry) => entry.taskId === taskId); + } catch (err: unknown) { + const errorMessage = err instanceof Error ? err.message : String(err); + log.warn(`Unable to inspect merge queue ownership for ${taskId}: ${errorMessage}`); + return false; + } + } + + private isFalseCompletionHandoffExhaustionWhileMergeOwned(task: Task): boolean { + return task.column === "in-review" + && task.status === "failed" + && typeof task.error === "string" + && task.error.includes("Completion handoff limbo recovery exhausted") + && this.isMergeLaneOwned(task.id); + } + private emitTaskMerged(task: Task | undefined | null, overrides: Partial = {}): void { if (!task) return; this.store.emit("task:merged", { @@ -5391,6 +5411,7 @@ export class SelfHealingManager { engineActivationGraceMs: settings.engineActivationGraceMs, }); if (!signal) continue; + if (this.isMergeLaneOwned(task.id)) continue; if (Date.parse(task.updatedAt) >= cycleStartMs) { continue; @@ -5520,6 +5541,7 @@ export class SelfHealingManager { if (!allowsAutoMergeProcessing(task, settings)) continue; if (task.paused === true) continue; if (task.id === activeMergeTaskId || executingTaskIds.has(task.id)) continue; + if (this.isMergeLaneOwned(task.id)) continue; const signal = getInReviewStalledSignal(task, { now: cycleStartMs, @@ -5690,6 +5712,7 @@ export class SelfHealingManager { allowsAutoMergeProcessing(task, settings) && !task.paused && !executingIds.has(task.id) && + !this.isMergeLaneOwned(task.id) && !(task.status && GHOST_REVIEW_PRESERVED_STATUSES.has(task.status)) && // Confirmed merges belong in `done` (handled by `recoverMergedReviewTasks`). task.mergeDetails?.mergeConfirmed !== true && @@ -6970,8 +6993,21 @@ export class SelfHealingManager { for (const task of tasks) { if (task.column !== "in-review" || task.paused) continue; if (!allowsAutoMergeProcessing(task, settings)) continue; + if (this.isFalseCompletionHandoffExhaustionWhileMergeOwned(task)) { + await this.store.updateTask(task.id, { + status: null, + error: null, + completionHandoffLimboRecoveryCount: 0, + }); + await this.store.logEntry( + task.id, + "Auto-recovered: cleared false completion-handoff exhaustion while task is already owned by merge queue", + ); + continue; + } if (task.status != null || task.mergeDetails != null || task.review != null || task.reviewState != null) continue; if (this.options.isTaskActive?.(task.id)) continue; + if (this.isMergeLaneOwned(task.id)) continue; if (getTaskMergeBlocker(task) !== undefined) continue; const doneMarker = [...(task.log ?? [])].reverse().find((entry) => entry.action === "Task marked done by agent");