feat(FN-4241): close stale merging fanout recovery gap in self-healing
Closes a gap in the stale merging fanout recovery path by adding recovery logic in the self-healing module, backed by unit tests covering the new behavior. Fusion-Task-Id: FN-4241
This commit is contained in:
5
.changeset/fn-4241-stale-merging-fanout-recovery.md
Normal file
5
.changeset/fn-4241-stale-merging-fanout-recovery.md
Normal file
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
Self-healing now clears downstream `blockedBy` fan-out when an `in-review` blocker has been stuck in `status=merging` past a configurable threshold, preventing a single hung merge-verification from freezing the todo lane.
|
||||
@@ -0,0 +1,92 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Settings, TaskStore } from "@fusion/core";
|
||||
import { SelfHealingManager } from "../self-healing.js";
|
||||
|
||||
function createTask(id: string, overrides: Record<string, unknown> = {}) {
|
||||
return {
|
||||
id,
|
||||
title: id,
|
||||
description: id,
|
||||
column: "todo",
|
||||
status: null,
|
||||
paused: false,
|
||||
blockedBy: null,
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
log: [],
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
describe("SelfHealingManager stale merge fanout recovery (FN-4241)", () => {
|
||||
let tasks: Map<string, Record<string, unknown>>;
|
||||
let store: TaskStore;
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date("2026-01-01T00:30:00.000Z"));
|
||||
tasks = new Map();
|
||||
|
||||
store = {
|
||||
getSettings: vi.fn().mockResolvedValue({
|
||||
globalPause: false,
|
||||
enginePaused: false,
|
||||
autoUnpauseEnabled: false,
|
||||
maintenanceIntervalMs: 0,
|
||||
} as unknown as Settings),
|
||||
listTasks: vi.fn().mockImplementation(async (options?: { column?: string; includeArchived?: boolean }) => {
|
||||
const all = Array.from(tasks.values());
|
||||
if (!options?.column) return all;
|
||||
return all.filter((task) => task.column === options.column);
|
||||
}),
|
||||
updateTask: vi.fn().mockImplementation(async (id: string, patch: Record<string, unknown>) => {
|
||||
const current = tasks.get(id);
|
||||
if (!current) throw new Error(`Task ${id} missing`);
|
||||
tasks.set(id, { ...current, ...patch });
|
||||
}),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
} as unknown as TaskStore;
|
||||
|
||||
const blocker = createTask("FN-4241-BLOCKER", {
|
||||
column: "in-review",
|
||||
status: "merging",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
});
|
||||
tasks.set(String(blocker.id), blocker);
|
||||
|
||||
for (let index = 1; index <= 5; index += 1) {
|
||||
const id = `FN-4241-DOWNSTREAM-${index}`;
|
||||
tasks.set(id, createTask(id, {
|
||||
column: "todo",
|
||||
blockedBy: "FN-4241-BLOCKER",
|
||||
status: "queued",
|
||||
}));
|
||||
}
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("FN-4241: clears stale merging status then unblocks downstream fanout", async () => {
|
||||
const manager = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
staleMergingStatusMinAgeMs: 5 * 60_000,
|
||||
staleMergingFanoutMinAgeMs: 15 * 60_000,
|
||||
});
|
||||
|
||||
const recoveredMerging = await manager.recoverStaleMergingStatus();
|
||||
expect(recoveredMerging).toBe(1);
|
||||
expect(tasks.get("FN-4241-BLOCKER")?.status).toBeNull();
|
||||
|
||||
const recoveredBlockedBy = await manager.clearStaleBlockedBy();
|
||||
expect(recoveredBlockedBy).toBe(5);
|
||||
|
||||
for (let index = 1; index <= 5; index += 1) {
|
||||
const downstream = tasks.get(`FN-4241-DOWNSTREAM-${index}`);
|
||||
expect(downstream?.blockedBy).toBeNull();
|
||||
}
|
||||
|
||||
manager.stop();
|
||||
});
|
||||
});
|
||||
@@ -4559,6 +4559,76 @@ describe("clearStaleBlockedBy", () => {
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it.each(["merging", "merging-pr"] as const)("clears stale blockedBy when blocker is stale in-review %s", async (status) => {
|
||||
vi.setSystemTime(new Date("2026-01-01T00:20:00.000Z"));
|
||||
const store = createRunningStore();
|
||||
const blockerId = "FN-510";
|
||||
const taskA = createTask("A", { blockedBy: blockerId });
|
||||
const taskB = createTask(blockerId, {
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status,
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
});
|
||||
mockSweepTasks(store, { todo: [taskA], inReview: [taskB], all: [taskA, taskB] });
|
||||
|
||||
const manager = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
const recovered = await manager.clearStaleBlockedBy();
|
||||
|
||||
expect(recovered).toBe(1);
|
||||
expect(store.updateTask).toHaveBeenCalledWith("A", { blockedBy: null, status: null });
|
||||
expect(store.logEntry).toHaveBeenCalledWith("A", expect.stringContaining(`blocker ${blockerId}`));
|
||||
expect(store.logEntry).toHaveBeenCalledWith("A", expect.stringContaining("stale for"));
|
||||
manager.stop();
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("does not clear stale merging blocker inside threshold", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-01T00:10:00.000Z"));
|
||||
const store = createRunningStore();
|
||||
const taskA = createTask("A", { blockedBy: "FN-511" });
|
||||
const taskB = createTask("FN-511", {
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status: "merging",
|
||||
updatedAt: "2026-01-01T00:00:01.000Z",
|
||||
});
|
||||
mockSweepTasks(store, { todo: [taskA], inReview: [taskB], all: [taskA, taskB] });
|
||||
|
||||
const manager = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
const recovered = await manager.clearStaleBlockedBy();
|
||||
|
||||
expect(recovered).toBe(0);
|
||||
expect(store.updateTask).not.toHaveBeenCalled();
|
||||
manager.stop();
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("honors staleMergingFanoutMinAgeMs option override", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-01T00:00:04.000Z"));
|
||||
const store = createRunningStore();
|
||||
const taskA = createTask("A", { blockedBy: "FN-512" });
|
||||
const taskB = createTask("FN-512", {
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status: "merging",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
});
|
||||
mockSweepTasks(store, { todo: [taskA], inReview: [taskB], all: [taskA, taskB] });
|
||||
|
||||
const manager = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
staleMergingStatusMinAgeMs: 1,
|
||||
staleMergingFanoutMinAgeMs: 2_000,
|
||||
});
|
||||
const recovered = await manager.clearStaleBlockedBy();
|
||||
|
||||
expect(recovered).toBe(1);
|
||||
expect(store.updateTask).toHaveBeenCalledWith("A", { blockedBy: null, status: null });
|
||||
manager.stop();
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("does not clear blockedBy when blocker failed retries are below threshold", async () => {
|
||||
const store = createRunningStore();
|
||||
const taskA = createTask("A", { blockedBy: "FN-600" });
|
||||
|
||||
@@ -93,6 +93,11 @@ export interface SelfHealingOptions {
|
||||
* Used to avoid clearing a transient merge status mid-merge.
|
||||
*/
|
||||
getActiveMergeTaskId?: () => string | null;
|
||||
/**
|
||||
* Minimum blocker age before stale merge fan-out is cleared from downstream
|
||||
* blockedBy pointers. Must be >= staleMergingStatusMinAgeMs.
|
||||
*/
|
||||
staleMergingFanoutMinAgeMs?: number;
|
||||
hasActiveAgentExecution?: (agentId: string) => boolean;
|
||||
restartDurableAgentHeartbeat?: (agentId: string, context: { reason: string; attempt: number }) => Promise<boolean>;
|
||||
}
|
||||
@@ -129,6 +134,7 @@ 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 DEFAULT_STALE_MERGING_FANOUT_MIN_AGE_MS = 15 * 60_000;
|
||||
const DURABLE_ERROR_RECOVERY_MAX_RETRIES = 5;
|
||||
const DURABLE_ERROR_RECOVERY_BASE_COOLDOWN_MS = 30_000;
|
||||
const DURABLE_ERROR_RECOVERY_MAX_COOLDOWN_MS = 15 * 60_000;
|
||||
@@ -1194,6 +1200,10 @@ export class SelfHealingManager {
|
||||
* 3. Blocker `column === "in-review"` and `paused === true`
|
||||
* 4. Blocker `column === "in-review"` and `status === "failed"`
|
||||
* and `(mergeRetries ?? 0) >= MAX_AUTO_MERGE_RETRIES`
|
||||
* 5. Blocker `column === "in-review"` and `status === "merging" | "merging-pr"`
|
||||
* (or a stale post-recovery `status === null` aftermath) with stale
|
||||
* `updatedAt` (older than `staleMergingFanoutMinAgeMs`) and no active
|
||||
* merger ownership in this process
|
||||
*
|
||||
* @returns Number of tasks unblocked
|
||||
*/
|
||||
@@ -1202,6 +1212,12 @@ export class SelfHealingManager {
|
||||
const settings = await this.store.getSettings();
|
||||
if (settings.globalPause || settings.enginePaused) return 0;
|
||||
|
||||
const staleMergingStatusMinAgeMs = this.options.staleMergingStatusMinAgeMs ?? DEFAULT_STALE_MERGING_STATUS_MIN_AGE_MS;
|
||||
const configuredFanoutMinAgeMs = this.options.staleMergingFanoutMinAgeMs ?? DEFAULT_STALE_MERGING_FANOUT_MIN_AGE_MS;
|
||||
const staleMergingFanoutMinAgeMs = Math.max(staleMergingStatusMinAgeMs, configuredFanoutMinAgeMs);
|
||||
const activeMergeTaskId = this.options.getActiveMergeTaskId?.() ?? null;
|
||||
const now = Date.now();
|
||||
|
||||
const todoTasks = await this.store.listTasks({ column: "todo", slim: true });
|
||||
const inProgressTasks = await this.store.listTasks({ column: "in-progress", slim: true });
|
||||
const inReviewTasks = await this.store.listTasks({ column: "in-review", slim: true });
|
||||
@@ -1261,6 +1277,19 @@ export class SelfHealingManager {
|
||||
isMissingWorktreeSessionStartFailure(blocker.error)
|
||||
) {
|
||||
reason = `blocker ${blockerId} in-review + failed (missing-worktree session start)`;
|
||||
} else if (
|
||||
blocker.column === "in-review" &&
|
||||
(blocker.status === "merging" || blocker.status === "merging-pr" || blocker.status == null) &&
|
||||
(!activeMergeTaskId || activeMergeTaskId !== blocker.id)
|
||||
) {
|
||||
const updatedAtMs = blocker.updatedAt ? Date.parse(blocker.updatedAt) : Number.NaN;
|
||||
if (Number.isFinite(updatedAtMs)) {
|
||||
const elapsedMs = now - updatedAtMs;
|
||||
if (elapsedMs >= staleMergingFanoutMinAgeMs) {
|
||||
const blockerStatus = blocker.status ?? "no-status";
|
||||
reason = `blocker ${blockerId} in-review + ${blockerStatus} stale for ${elapsedMs}ms (threshold ${staleMergingFanoutMinAgeMs}ms)`;
|
||||
}
|
||||
}
|
||||
} else if (task.dependencies.length > 0 && !unresolvedDeps.includes(blockerId)) {
|
||||
reason = `blocker ${blockerId} not among unresolved dependencies`;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user