feat(FN-4084): recover tasks blocked by stalled merger self-healing

Merge starvation recovery was hardened across the engine: self-healing now detects and clears blocked-in-review tasks that starve the merger, the project engine gains defensive recovery hooks for stalled merges, and tests cover the new recovery paths.

Fusion-Task-Id: FN-4084
This commit is contained in:
Fusion
2026-05-12 07:30:46 -07:00
committed by gsxdsm
parent 4a421d8ea8
commit 5804b8651c
7 changed files with 264 additions and 31 deletions

View File

@@ -0,0 +1,9 @@
---
"@runfusion/fusion": patch
---
Fix merger starvation where eligible in-review tasks looped in auto-recovery
without ever merging. Leaked in-memory merge-queue entries are now reconciled
automatically, and tasks whose re-enqueue is repeatedly dropped escalate to a
clear `status=failed` with an `Auto-merge starvation:` error instead of
looping indefinitely.

View File

@@ -97,6 +97,7 @@ Fusion task columns:
- Persisted executor session state is resumed only when it still matches the task's current worktree context. If a retry fails with `Refusing to start coding agent in missing worktree: ...` and the persisted session points at stale worktree metadata, recovery clears stale session pointers and retries fresh so review retries do not reopen deleted worktree paths.
- Merge-confirmed tasks still respect `getTaskMergeBlocker()` before the final `in-review``done` move. If merge is confirmed but a blocker remains (for example, incomplete steps), Fusion parks the task in `in-review` with `status: "failed"` and an explicit blocker error instead of retry-looping auto-finalization.
- Self-healing can still auto-finalize retry-exhausted failed review tasks when it can prove their branch content already landed on the merge target, so already-merged work does not deadlock in `in-review`.
- Repeated engine merge-queue drops now escalate to an explicit recoverable review failure: if auto-recovery hits `Auto-merge starvation:` in the task `error`, Fusion has already seen three consecutive enqueue attempts rejected by the engine merge queue. Operators can recover by clearing the failed state from the dashboard, which lets the usual unpause/clear flow re-attempt merge once the underlying queue wedge is resolved.
- Non-recoverable state-machine errors during finalization (for example `Invalid transition: 'todo' → 'done'`) are treated as terminal review failures: recovery must not re-enqueue these tasks for merge unless task state changes prove they are recoverable.
5. **done** — merged/finalized
6. **archived** — preserved history, optionally cleaned from filesystem

View File

@@ -2053,7 +2053,7 @@ describe("ProjectEngine stale mergeActive rescue (FN-3900)", () => {
await engine.stop();
});
it("internalEnqueueMerge warns and skips direct leaked mergeActive entries", async () => {
it("FN-4084: internalEnqueueMerge reconciles leaked mergeActive entries", async () => {
const mockStore = createMockStore({ ...baseSettings, autoMerge: true });
mocks.currentStore = mockStore.store;
@@ -2062,6 +2062,7 @@ describe("ProjectEngine stale mergeActive rescue (FN-3900)", () => {
mergeQueue: string[];
mergeActive: Set<string>;
activeMergeTaskId: string | null;
mergeRunning: boolean;
internalEnqueueMerge: (taskId: string) => void;
};
const warnSpy = vi.spyOn(runtimeLog, "warn").mockImplementation(() => {});
@@ -2071,12 +2072,13 @@ describe("ProjectEngine stale mergeActive rescue (FN-3900)", () => {
privateEngine.mergeActive = new Set(["FN-leaked2"]);
privateEngine.mergeQueue = [];
privateEngine.activeMergeTaskId = null;
privateEngine.mergeRunning = true;
privateEngine.internalEnqueueMerge("FN-leaked2");
expect(warnSpy).toHaveBeenCalledWith(expect.stringContaining("mergeActive entry is leaked"));
expect(privateEngine.mergeQueue).toEqual([]);
expect(mocks.aiMergeTask).not.toHaveBeenCalled();
expect(privateEngine.mergeQueue).toEqual(["FN-leaked2"]);
expect(privateEngine.mergeActive.has("FN-leaked2")).toBe(true);
await engine.stop();
});

View File

@@ -2243,6 +2243,38 @@ describe("SelfHealingManager", () => {
managerWithRecovery.stop();
});
it("FN-4084: stale merging recovery clears mergeActive via callback", async () => {
const clearMergeActive = vi.fn();
const managerWithRecovery = new SelfHealingManager(store, {
rootDir: "/tmp/test-project",
clearMergeActive,
});
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
globalPause: false,
enginePaused: false,
});
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
{
id: "FN-4084-stale",
column: "in-review",
paused: false,
status: "merging-pr",
updatedAt: new Date(Date.now() - 10 * 60_000).toISOString(),
steps: [{ name: "Ship it", status: "done" }],
workflowStepResults: [],
log: [],
},
]);
const result = await managerWithRecovery.recoverStaleMergingStatus();
expect(result).toBe(1);
expect(clearMergeActive).toHaveBeenCalledTimes(1);
expect(clearMergeActive).toHaveBeenCalledWith("FN-4084-stale");
managerWithRecovery.stop();
});
it("keeps transient merge status when task is actively merging", async () => {
const managerWithRecovery = new SelfHealingManager(store, {
rootDir: "/tmp/test-project",
@@ -2340,7 +2372,7 @@ describe("SelfHealingManager", () => {
});
it("routes through enqueueMerge when wired so mergeStrategy is honored", async () => {
const enqueueMerge = vi.fn();
const enqueueMerge = vi.fn().mockReturnValue(true);
const managerWithRecovery = new SelfHealingManager(store, {
rootDir: "/tmp/test-project",
enqueueMerge,
@@ -2375,6 +2407,105 @@ describe("SelfHealingManager", () => {
managerWithRecovery.stop();
});
it("FN-4084: recoverMergeableReviewTasks escalates after repeated no-op re-enqueues", async () => {
const enqueueMerge = vi.fn().mockReturnValue(false);
const managerWithRecovery = new SelfHealingManager(store, {
rootDir: "/tmp/test-project",
enqueueMerge,
});
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
autoMerge: true,
globalPause: false,
enginePaused: false,
});
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
{
id: "FN-4084-starved",
column: "in-review",
paused: false,
status: null,
error: null,
mergeRetries: 0,
worktree: "/tmp/test-project/.worktrees/fn-4084-starved",
steps: [{ name: "Ship it", status: "done" }],
workflowStepResults: [{ id: "ws-1", status: "passed", phase: "pre-merge" }],
mergeDetails: undefined,
log: [],
},
]);
expect(await managerWithRecovery.recoverMergeableReviewTasks()).toBe(0);
expect(await managerWithRecovery.recoverMergeableReviewTasks()).toBe(0);
expect(await managerWithRecovery.recoverMergeableReviewTasks()).toBe(1);
expect(enqueueMerge).toHaveBeenCalledTimes(3);
expect(store.updateTask).toHaveBeenCalledWith(
"FN-4084-starved",
expect.objectContaining({
status: "failed",
error: expect.stringContaining("Auto-merge starvation: 3 consecutive enqueue attempts"),
}),
);
expect(store.logEntry).toHaveBeenCalledTimes(1);
expect(store.logEntry).toHaveBeenCalledWith(
"FN-4084-starved",
expect.stringContaining("Auto-merge starvation"),
);
expect(store.logEntry).not.toHaveBeenCalledWith(
"FN-4084-starved",
expect.stringContaining("re-enqueued for merge"),
);
managerWithRecovery.stop();
});
it("FN-4084: recoverMergeableReviewTasks resets starvation counters after successful enqueue", async () => {
const enqueueMerge = vi.fn().mockReturnValue(true);
const managerWithRecovery = new SelfHealingManager(store, {
rootDir: "/tmp/test-project",
enqueueMerge,
});
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
autoMerge: true,
globalPause: false,
enginePaused: false,
});
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
{
id: "FN-4084-healthy",
column: "in-review",
paused: false,
status: null,
error: null,
mergeRetries: 0,
worktree: "/tmp/test-project/.worktrees/fn-4084-healthy",
steps: [{ name: "Ship it", status: "done" }],
workflowStepResults: [{ id: "ws-1", status: "passed", phase: "pre-merge" }],
mergeDetails: undefined,
log: [],
},
]);
for (let i = 0; i < 5; i++) {
expect(await managerWithRecovery.recoverMergeableReviewTasks()).toBe(1);
}
expect(enqueueMerge).toHaveBeenCalledTimes(5);
expect(store.updateTask).not.toHaveBeenCalledWith(
"FN-4084-healthy",
expect.objectContaining({ status: "failed" }),
);
expect(store.logEntry).toHaveBeenCalledTimes(5);
expect(store.logEntry).not.toHaveBeenCalledWith(
"FN-4084-healthy",
expect.stringContaining("Auto-merge starvation"),
);
managerWithRecovery.stop();
});
it("skips entirely when autoMerge is disabled (respects PR-based review flow)", async () => {
const enqueueMerge = vi.fn();
const managerWithRecovery = new SelfHealingManager(store, {

View File

@@ -177,6 +177,7 @@ export class ProjectEngine {
private mergeAbortController: AbortController | null = null;
private mergeRetryTimer: ReturnType<typeof setTimeout> | null = null;
private autostashSweepTimer: ReturnType<typeof setTimeout> | null = null;
private mergeActiveReconcileTimer: ReturnType<typeof setInterval> | null = null;
/**
* Pending manual merge resolvers — keyed by taskId.
@@ -245,7 +246,10 @@ export class ProjectEngine {
this.activeMergeTaskId = null;
}
this.mergeActive.delete(taskId);
this.internalEnqueueMerge(taskId);
return this.internalEnqueueMerge(taskId);
});
this.runtime.setMergeActiveClearer?.((taskId) => {
this.mergeActive.delete(taskId);
});
}
@@ -452,6 +456,7 @@ export class ProjectEngine {
// 8. Start periodic merge retry sweep
this.scheduleMergeRetry(store);
this.scheduleMergeActiveReconciliation(settings.maintenanceIntervalMs ?? 900_000);
// 9. Startup + periodic stale autostash sweeps (independent of autoMerge)
void this.runStaleAutostashSweep(store, "startup");
@@ -484,6 +489,10 @@ export class ProjectEngine {
clearTimeout(this.autostashSweepTimer);
this.autostashSweepTimer = null;
}
if (this.mergeActiveReconcileTimer) {
clearInterval(this.mergeActiveReconcileTimer);
this.mergeActiveReconcileTimer = null;
}
// Abort active/pending merge work before tearing down sessions.
this.mergeAbortController?.abort();
@@ -758,8 +767,8 @@ export class ProjectEngine {
* Exposed publicly so callers can integrate the engine's merge queue with
* an external `onMerge` callback (e.g. dashboard's createServer call).
*/
enqueueMerge(taskId: string): void {
this.internalEnqueueMerge(taskId);
enqueueMerge(taskId: string): boolean {
return this.internalEnqueueMerge(taskId);
}
/**
@@ -781,7 +790,10 @@ export class ProjectEngine {
return new Promise<MergeResult>((resolve, reject) => {
this.manualMergeResolvers.set(taskId, { resolve, reject });
this.internalEnqueueMerge(taskId);
if (!this.internalEnqueueMerge(taskId)) {
this.manualMergeResolvers.delete(taskId);
reject(new Error(`Merge enqueue rejected for ${taskId}`));
}
});
}
@@ -1145,23 +1157,23 @@ export class ProjectEngine {
return undefined;
}
private internalEnqueueMerge(taskId: string): void {
if (this.shuttingDown) return;
private internalEnqueueMerge(taskId: string): boolean {
if (this.shuttingDown) return false;
if (this.mergeActive.has(taskId)) {
// Distinguish "actually being processed" (queued or active) from a
// leaked entry. Leaks are dropped by reconcileStaleMergeActive() on the
// next 15s sweep, so we only log the genuinely-busy case at debug
// verbosity. Without this log the de-dup was invisible — a leaked
// entry made every subsequent enqueue silently no-op until the 15-min
// maintenance loop woke up.
// leaked entry. Reconcile leaks immediately so recovery paths and fresh
// in-review handoffs can make forward progress without waiting for the
// periodic maintenance sweep.
const isActuallyLive =
this.mergeQueue.includes(taskId) || this.activeMergeTaskId === taskId;
if (!isActuallyLive) {
runtimeLog.warn(
`internalEnqueueMerge(${taskId}): skipped — mergeActive entry is leaked (not queued, not active). reconcileStaleMergeActive() will clear it on the next sweep.`,
`internalEnqueueMerge(${taskId}): skipped — mergeActive entry is leaked (not queued, not active). Reconciling stale entry and retrying enqueue now.`,
);
this.mergeActive.delete(taskId);
} else {
return false;
}
return;
}
this.mergeActive.add(taskId);
this.mergeQueue.push(taskId);
@@ -1170,6 +1182,7 @@ export class ProjectEngine {
`Merge queue drain failed unexpectedly: ${err instanceof Error ? err.message : String(err)}`,
);
});
return true;
}
/**
@@ -1190,6 +1203,29 @@ export class ProjectEngine {
return eligible.length;
}
private reconcileStaleMergeActive(): number {
let cleared = 0;
for (const taskId of [...this.mergeActive]) {
if (taskId === this.activeMergeTaskId) continue;
if (this.mergeQueue.includes(taskId)) continue;
this.mergeActive.delete(taskId);
cleared++;
}
return cleared;
}
private scheduleMergeActiveReconciliation(intervalMs: number): void {
if (!Number.isFinite(intervalMs) || intervalMs <= 0) {
return;
}
this.mergeActiveReconcileTimer = setInterval(() => {
const cleared = this.reconcileStaleMergeActive();
if (cleared > 0) {
runtimeLog.warn(`Reconciled ${cleared} stale mergeActive entr${cleared === 1 ? "y" : "ies"}`);
}
}, intervalMs);
}
private async findActiveRecoveryFollowUp(
store: TaskStore,
parentTaskId: string,
@@ -1221,6 +1257,7 @@ export class ProjectEngine {
this.mergeRunning = true;
try {
this.reconcileStaleMergeActive();
const store = this.runtime.getTaskStore();
const cwd = this.config.workingDirectory;
@@ -1872,11 +1909,9 @@ export class ProjectEngine {
runtimeLog.log(`Auto-merge handoff (${task.id}) skipped: autoMerge disabled`);
return;
}
// Belt-and-braces: clear any stale mergeActive entry from a wedged
// prior attempt so this enqueue isn't silently no-op'd. The 15s
// sweep also reconciles via reconcileStaleMergeActive(), but waiting
// up to 15s for a fresh in-review task to start merging is the
// exact regression we're fixing.
// Belt-and-braces: eager handoff still clears a stale mergeActive
// entry before enqueue so freshly completed review tasks do not wait
// for a later queue reconciliation pass before their merge starts.
if (
this.mergeActive.has(task.id) &&
!this.mergeQueue.includes(task.id) &&
@@ -2228,7 +2263,27 @@ export class ProjectEngine {
store.on("settings:updated", onEngineUnpause);
this.settingsHandlers.push(onEngineUnpause);
// 5. Stuck task timeout change — trigger immediate check
// 5. Maintenance interval change — reschedule mergeActive reconciliation
const onMaintenanceIntervalChange = ({
settings: s,
previous: prev,
}: {
settings: Settings;
previous: Settings;
}) => {
if (s.maintenanceIntervalMs === prev.maintenanceIntervalMs) {
return;
}
if (this.mergeActiveReconcileTimer) {
clearInterval(this.mergeActiveReconcileTimer);
this.mergeActiveReconcileTimer = null;
}
this.scheduleMergeActiveReconciliation(s.maintenanceIntervalMs ?? 900_000);
};
store.on("settings:updated", onMaintenanceIntervalChange);
this.settingsHandlers.push(onMaintenanceIntervalChange);
// 6. Stuck task timeout change — trigger immediate check
const onStuckTimeoutChange = async ({
settings: s,
previous: prev,
@@ -2254,7 +2309,7 @@ export class ProjectEngine {
store.on("settings:updated", onStuckTimeoutChange);
this.settingsHandlers.push(onStuckTimeoutChange);
// 5. Memory maintenance settings change — sync automations
// 7. Memory maintenance settings change — sync automations
const onInsightSettingsChange = async ({
settings: s,
previous: prev,
@@ -2300,7 +2355,7 @@ export class ProjectEngine {
store.on("settings:updated", onInsightSettingsChange);
this.settingsHandlers.push(onInsightSettingsChange);
// 6. Auto-summarize settings change — sync automation
// 8. Auto-summarize settings change — sync automation
const onAutoSummarizeSettingsChange = async ({
settings: s,
previous: prev,
@@ -2336,7 +2391,7 @@ export class ProjectEngine {
store.on("settings:updated", onAutoSummarizeSettingsChange);
this.settingsHandlers.push(onAutoSummarizeSettingsChange);
// 7. Scheduled eval settings change — sync automation
// 9. Scheduled eval settings change — sync automation
const onScheduledEvalSettingsChange = async ({
settings: s,
previous: prev,

View File

@@ -120,7 +120,8 @@ export class InProcessRuntime
* stale-merge recovery can re-enqueue tasks immediately. Set by ProjectEngine
* before `start()` via `setMergeEnqueuer`.
*/
private mergeEnqueuer?: (taskId: string) => void;
private mergeEnqueuer?: (taskId: string) => boolean;
private clearMergeActive?: (taskId: string) => void;
private activeMergeTaskIdProvider?: () => string | null;
/** Tracks whether startup recovery was intentionally deferred due to pause state. */
private startupRecoveryDeferred = false;
@@ -632,7 +633,8 @@ export class InProcessRuntime
recoverApprovedTriageTask: (task) => this.triageProcessor?.recoverApprovedTask(task) ?? Promise.resolve(false),
getPlanningTaskIds: () => this.triageProcessor?.getProcessingTaskIds() ?? new Set<string>(),
evictStaleTriageProcessing: () => this.triageProcessor?.evictStaleProcessing() ?? new Set<string>(),
enqueueMerge: this.mergeEnqueuer ? (taskId: string) => this.mergeEnqueuer?.(taskId) : undefined,
enqueueMerge: this.mergeEnqueuer ? (taskId: string) => this.mergeEnqueuer?.(taskId) ?? false : undefined,
clearMergeActive: this.clearMergeActive ? (taskId: string) => this.clearMergeActive?.(taskId) : undefined,
getActiveMergeTaskId: () => this.activeMergeTaskIdProvider?.() ?? null,
leaseManager: this.leaseManager,
hasActiveAgentExecution: (agentId: string) => this.heartbeatMonitor?.getTrackedAgents().includes(agentId) ?? false,
@@ -887,10 +889,14 @@ export class InProcessRuntime
* auto-merge after clearing a stale `merging` status. Must be called before
* `start()` because SelfHealingManager is constructed during startup.
*/
setMergeEnqueuer(enqueueMerge: (taskId: string) => void): void {
setMergeEnqueuer(enqueueMerge: (taskId: string) => boolean): void {
this.mergeEnqueuer = enqueueMerge;
}
setMergeActiveClearer(clearMergeActive: (taskId: string) => void): void {
this.clearMergeActive = clearMergeActive;
}
setActiveMergeTaskIdProvider(getActiveMergeTaskId: () => string | null): void {
this.activeMergeTaskIdProvider = getActiveMergeTaskId;
}

View File

@@ -81,7 +81,8 @@ export interface SelfHealingOptions {
* refreshed (otherwise a leftover entry from a SIGKILL'd merge would cause
* the polling sweep's enqueue to silently no-op).
*/
enqueueMerge?: (taskId: string) => void;
enqueueMerge?: (taskId: string) => boolean;
clearMergeActive?: (taskId: string) => void;
/**
* Minimum age before a transient merge status is considered stale when no
* active merge session is associated with that task.
@@ -125,6 +126,7 @@ const ORPHANED_WITH_WORKTREE_GRACE_MS = 300_000;
*/
const MAX_TASK_DONE_RETRIES = 3;
const MAX_AUTO_MERGE_RETRIES = 3;
const MAX_STARVATION_DROPS = 3;
const DEADLOCK_RECOVERY_COOLDOWN_MS = 15 * 60_000;
const DEFAULT_STALE_MERGING_STATUS_MIN_AGE_MS = 5 * 60_000;
const DURABLE_ERROR_RECOVERY_MAX_RETRIES = 5;
@@ -256,6 +258,7 @@ export class SelfHealingManager {
// ── Per-task deadlock recovery cooldown ─────────────────────────────
private deadlockRecoveryCooldown: Map<string, number> = new Map();
private mergeStarvationDrops: Map<string, number> = new Map();
constructor(
private store: TaskStore,
@@ -1158,6 +1161,7 @@ export class SelfHealingManager {
try {
log.warn(`Clearing stale merge status for ${task.id}: ${previousStatus}`);
await this.store.updateTask(task.id, { status: null });
this.options.clearMergeActive?.(task.id);
await this.store.logEntry(
task.id,
`Auto-recovered: cleared stale '${previousStatus}' status (no active merger)`,
@@ -1388,6 +1392,7 @@ export class SelfHealingManager {
const mergeable = tasks.filter((t) =>
t.column === "in-review" &&
!t.paused &&
t.status !== "failed" &&
// Exclude transient merge statuses. Active merges should be left alone;
// stale ones are handled by recoverStaleMergingStatus().
t.status !== "merging" &&
@@ -1405,6 +1410,14 @@ export class SelfHealingManager {
getTaskMergeBlocker(t) === undefined,
);
const inReviewIds = new Set(tasks.map((task) => task.id));
const mergeableIds = new Set(mergeable.map((task) => task.id));
for (const taskId of [...this.mergeStarvationDrops.keys()]) {
if (!inReviewIds.has(taskId) || !mergeableIds.has(taskId)) {
this.mergeStarvationDrops.delete(taskId);
}
}
if (mergeable.length === 0) return 0;
log.warn(`Found ${mergeable.length} mergeable review task(s) stuck in in-review`);
@@ -1417,7 +1430,23 @@ export class SelfHealingManager {
for (const task of mergeable) {
try {
if (enqueueMerge) {
enqueueMerge(task.id);
const queued = enqueueMerge(task.id);
if (!queued) {
const drops = (this.mergeStarvationDrops.get(task.id) ?? 0) + 1;
this.mergeStarvationDrops.set(task.id, drops);
log.warn(
`Auto-recovery enqueue dropped for ${task.id} (${drops}/${MAX_STARVATION_DROPS}); engine merge queue rejected re-enqueue`,
);
if (drops >= MAX_STARVATION_DROPS) {
const error = `Auto-merge starvation: ${MAX_STARVATION_DROPS} consecutive enqueue attempts were dropped by the engine merge queue; task requires manual intervention.`;
await this.store.updateTask(task.id, { status: "failed", error });
await this.store.logEntry(task.id, error);
this.mergeStarvationDrops.delete(task.id);
recovered++;
}
continue;
}
this.mergeStarvationDrops.delete(task.id);
} else {
await this.store.mergeTask(task.id);
}