fix(FN-6086): preserve merge-queued review tasks

Treat persisted merge queue ownership as an active merge-lane owner when hydrating in-review stall state and when self-healing scans for ghost, stalled, or completion-handoff-limbo tasks.

Also clear false completion-handoff exhaustion for tasks that are still owned by the merge queue so previously poisoned review tasks can continue merging.

Fusion-Task-Id: FN-6086
This commit is contained in:
gsxdsm
2026-06-09 09:41:28 -07:00
parent a937bc6ab0
commit 7a9d2b0f30
6 changed files with 189 additions and 11 deletions

View File

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

View File

@@ -2738,6 +2738,11 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS}
return this.rowToTask(row);
}
private getMergeQueuedTaskIds(): Set<string> {
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),

View File

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

View File

@@ -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<string, unknown> = {}): TaskStore & EventEmitter {
function createStore(task: Task, settings: Record<string, unknown> = {}, 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<string, unknown> = {}): 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);

View File

@@ -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<MergeResult> = {}): 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");