feat(FN-5094): merge fusion/fn-5094

Commits merged:
- merge fusion/fn-5094

Files changed:
.../self-healing-stale-merger-status.test.ts       | 200 +++++++++++++++++++++
 packages/engine/src/branch-conflicts.ts            |   6 +-
 packages/engine/src/run-audit.ts                   |   2 +
 packages/engine/src/self-healing.ts                |  67 ++++++-
 4 files changed, 271 insertions(+), 4 deletions(-)

Fusion-Task-Id: FN-5094

Fusion-Task-Lineage: c6c9710c-51fe-40e9-bfad-bfb6d589c7f4
This commit is contained in:
Fusion (runfusion.ai)
2026-05-18 21:51:07 -07:00
committed by gsxdsm
parent 0c35e6db6f
commit a0b75579b8
4 changed files with 271 additions and 4 deletions

View File

@@ -0,0 +1,200 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { EventEmitter } from "node:events";
import type { Settings, Task, TaskStore } from "@fusion/core";
const { execMock } = vi.hoisted(() => ({
execMock: vi.fn(),
}));
vi.mock("node:child_process", () => ({ exec: execMock, execSync: vi.fn() }));
const { logger } = vi.hoisted(() => ({
logger: { log: vi.fn(), warn: vi.fn(), error: vi.fn() },
}));
vi.mock("../logger.js", () => ({ createLogger: vi.fn(() => logger) }));
const { recordRunAuditEventMock } = vi.hoisted(() => ({
recordRunAuditEventMock: vi.fn(async () => undefined),
}));
vi.mock("../run-audit.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../run-audit.js")>();
return {
...actual,
createRunAuditor: vi.fn(() => ({
database: recordRunAuditEventMock,
git: vi.fn(),
filesystem: vi.fn(),
sandbox: vi.fn(),
})),
};
});
import { SelfHealingManager } from "../self-healing.js";
function makeTask(id: string, overrides: Partial<Task> = {}): Task {
return {
id,
title: id,
description: id,
column: "done",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: new Date(Date.now() - 12 * 60 * 60 * 1000).toISOString(),
updatedAt: new Date(Date.now() - 12 * 60 * 60 * 1000).toISOString(),
...overrides,
} as Task;
}
function createStore(tasks: Task[]): TaskStore & EventEmitter {
const map = new Map(tasks.map((t) => [t.id, t]));
const emitter = new EventEmitter();
return Object.assign(emitter, {
getSettings: vi.fn(async () => ({ globalPause: false, enginePaused: false } as Settings)),
listTasks: vi.fn(async (opts?: { column?: Task["column"]; slim?: boolean }) => {
const all = [...map.values()];
if (!opts?.column) return all;
return all.filter((t) => t.column === opts.column);
}),
getTask: vi.fn(async (id: string) => map.get(id)),
updateTask: vi.fn(async (id: string, patch: Partial<Task>) => {
const task = map.get(id)!;
const merged = { ...task, ...patch } as Task;
map.set(id, merged);
return merged;
}),
logEntry: vi.fn(async () => undefined),
}) as unknown as TaskStore & EventEmitter;
}
describe("FN-5092: reconcileStaleMergerStatus watchdog", () => {
beforeEach(() => {
vi.clearAllMocks();
});
it("clears status=\"merging\" on a done task (the FN-5052 stranding case)", async () => {
const stranded = makeTask("FN-5052", {
column: "done",
status: "merging" as Task["status"],
mergeDetails: { commitSha: "abc123", mergeConfirmed: true } as Task["mergeDetails"],
});
const store = createStore([stranded]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(1);
const after = await store.getTask("FN-5052");
expect(after?.status).toBeNull();
expect(after?.column).toBe("done");
// mergeDetails preserved
expect(after?.mergeDetails?.commitSha).toBe("abc123");
// Audit event recorded
expect(recordRunAuditEventMock).toHaveBeenCalledWith(
expect.objectContaining({
type: "task:auto-recover-stale-merger-status",
target: "FN-5052",
metadata: expect.objectContaining({
previousColumn: "done",
previousStatus: "merging",
mergeConfirmed: true,
commitSha: "abc123",
}),
}),
);
// Log entry recorded for forensics
expect((store as any).logEntry).toHaveBeenCalledWith(
"FN-5052",
expect.stringContaining("Auto-recovered: cleared stale status=\"merging\""),
);
});
it("clears status=\"merging-pr\" on a done task", async () => {
const stranded = makeTask("FN-9999", {
column: "done",
status: "merging-pr" as Task["status"],
});
const store = createStore([stranded]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(1);
expect((await store.getTask("FN-9999"))?.status).toBeNull();
});
it("also catches the same leak on archived tasks", async () => {
const stranded = makeTask("FN-ARCH", {
column: "archived",
status: "merging" as Task["status"],
});
const store = createStore([stranded]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(1);
expect((await store.getTask("FN-ARCH"))?.status).toBeNull();
});
it("does not touch in-review tasks that legitimately have status=\"merging\"", async () => {
const legit = makeTask("FN-INREVIEW", {
column: "in-review",
status: "merging" as Task["status"],
});
const store = createStore([legit]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(0);
expect((await store.getTask("FN-INREVIEW"))?.status).toBe("merging");
});
it("does not touch done tasks with null status (the healthy case)", async () => {
const healthy = makeTask("FN-HEALTHY", { column: "done", status: undefined });
const store = createStore([healthy]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(0);
expect(recordRunAuditEventMock).not.toHaveBeenCalled();
});
it("handles multiple leaked tasks in one sweep", async () => {
const tasks = [
makeTask("FN-A", { column: "done", status: "merging" as Task["status"] }),
makeTask("FN-B", { column: "done", status: "merging-pr" as Task["status"] }),
makeTask("FN-C", { column: "archived", status: "merging" as Task["status"] }),
makeTask("FN-D", { column: "done", status: undefined }), // healthy
];
const store = createStore(tasks);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(3);
});
it("continues sweep when a single task update fails", async () => {
const a = makeTask("FN-FAIL", { column: "done", status: "merging" as Task["status"] });
const b = makeTask("FN-OK", { column: "done", status: "merging" as Task["status"] });
const store = createStore([a, b]);
let failOnce = true;
(store as any).updateTask = vi.fn(async (id: string, patch: Partial<Task>) => {
if (id === "FN-FAIL" && failOnce) {
failOnce = false;
throw new Error("simulated update failure");
}
const task = (store as any).getTask.mock.results[0]?.value;
return { ...(task || a), ...patch };
});
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
// FN-FAIL fails, FN-OK succeeds
expect(cleared).toBe(1);
});
it("returns 0 on empty board (no done/archived tasks at all)", async () => {
const store = createStore([makeTask("FN-INPROGRESS", { column: "in-progress" })]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
const cleared = await mgr.reconcileStaleMergerStatus();
expect(cleared).toBe(0);
});
});

View File

@@ -613,12 +613,12 @@ export async function classifyForeignOnlyContamination(
} catch {
// fall back to persisted baseSha on any git failure
}
const output = await runGit(repoDir, `git log --format=%H%x1f%s%x1f%b ${quoteShellArg(`${effectiveBaseSha}..${branchName}`)}`)
const persistedRangeOutput = await runGit(repoDir, `git log --format=%H%x1f%s%x1f%b ${quoteShellArg(`${baseSha}..${branchName}`)}`)
.catch(() => "");
const subjectPattern = /^(feat|fix|test|chore|docs|refactor|perf|build)\((FN-\d+)\):/i;
const trailerPattern = /(?:^|\n)Fusion-Task-Id:\s*(FN-\d+)\s*(?:\n|$)/i;
const foreignCommits: BranchCrossContaminationCommit[] = [];
for (const line of output.split("\n").map((entry) => entry.trim()).filter(Boolean)) {
for (const line of persistedRangeOutput.split("\n").map((entry) => entry.trim()).filter(Boolean)) {
const [sha, subject, body] = line.split("\u001f");
const subjectMatch = (subject ?? "").match(subjectPattern);
const trailerMatch = (body ?? "").match(trailerPattern);
@@ -650,7 +650,7 @@ export async function classifyForeignOnlyContamination(
const foreignClassification = await classifyForeignCommits({
repoDir,
branchName,
baseSha,
baseSha: effectiveBaseSha,
foreignCommits,
mainRef,
});

View File

@@ -249,6 +249,8 @@ export type DatabaseMutationType =
| "task:finalize-unproven-blocked"
| "task:integrity-reconcile-modified-files"
| "task:integrity-warning"
/** FN-5092 watchdog: stale `status: "merging"` / `"merging-pr"` cleared on a done/archived task. Metadata: { previousColumn, previousStatus, ageMs, mergeConfirmed?: boolean } */
| "task:auto-recover-stale-merger-status"
| "auto-recovery:classify-decision"
| "auto-recovery:retry-issued"
| "auto-recovery:ai-session-spawned"

View File

@@ -632,6 +632,9 @@ export class SelfHealingManager {
{ name: "interrupted-merging", fn: () => this.recoverInterruptedMergingTasks().then(() => undefined) },
{ name: "done-merge-metadata", fn: () => this.recoverDoneTaskMergeMetadata().then(() => undefined) },
{ name: "reconcile-done-task-integrity", fn: () => this.reconcileDoneTaskIntegrity().then(() => undefined) },
// FN-5092: must run BEFORE any merger pickup path so the merger queue is
// not stalled by a leaked `status: "merging"` on an already-done task.
{ name: "reconcile-stale-merger-status", fn: () => this.reconcileStaleMergerStatus().then(() => undefined) },
{ name: "recover-already-merged-review", fn: () => this.recoverAlreadyMergedReviewTasks().then(() => undefined) },
{ name: "recover-completion-handoff-limbo", fn: () => this.recoverCompletionHandoffLimbo().then(() => undefined) },
{ name: "recover-branch-misbound-in-review", fn: () => this.recoverBranchMisboundInReviewTasks().then(() => undefined) },
@@ -1209,6 +1212,7 @@ export class SelfHealingManager {
{ name: "recover-stale-merging-status", fn: () => this.recoverStaleMergingStatus() },
{ name: "finalize-noop-review", fn: () => this.finalizeNoOpReviewTasks() },
{ name: "reconcile-done-task-integrity", fn: () => this.reconcileDoneTaskIntegrity() },
{ name: "reconcile-stale-merger-status", fn: () => this.reconcileStaleMergerStatus() },
{ name: "recover-mergeable-review", fn: () => this.recoverMergeableReviewTasks() },
{ name: "recover-merged-review", fn: () => this.recoverMergedReviewTasks() },
{ name: "recover-already-merged-review", fn: () => this.recoverAlreadyMergedReviewTasks() },
@@ -3344,7 +3348,7 @@ export class SelfHealingManager {
return recovered;
}
private async recordIntegrityAudit(taskId: string, mutationType: "task:finalize-unproven-blocked" | "task:integrity-reconcile-modified-files" | "task:integrity-warning", metadata: Record<string, unknown>): Promise<void> {
private async recordIntegrityAudit(taskId: string, mutationType: "task:finalize-unproven-blocked" | "task:integrity-reconcile-modified-files" | "task:integrity-warning" | "task:auto-recover-stale-merger-status", metadata: Record<string, unknown>): Promise<void> {
const auditor = createRunAuditor(this.store, {
runId: generateSyntheticRunId("self-healing-integrity", taskId),
agentId: "self-healing",
@@ -3354,6 +3358,67 @@ export class SelfHealingManager {
await auditor.database({ type: mutationType, target: taskId, metadata });
}
/**
* FN-5092 watchdog: detect and repair tasks left in an impossible state where
* `column ∈ {done, archived}` but `status ∈ {merging, merging-pr}`.
*
* Cause: a recovery path (FN-4499 misbinding, FN-4500 already-on-main, manual
* finalization) moved the task to done WITHOUT going through the merger's
* `completeTask()`. The merger had previously set `status = "merging"` when it
* claimed `mergeActive[taskId]`; that slot is single-threaded and now leaks,
* stalling the entire merger queue for every subsequent in-review task.
*
* This watchdog catches the persistent-state half of the leak. The runtime
* in-memory `mergeActive` Map also has a periodic reconciler in
* `ProjectEngine.reconcileStaleMergeActive()`, but it skips entries that match
* the currently-active merge task; a status leak that survives across engine
* restarts can only be cleared at the storage layer.
*/
async reconcileStaleMergerStatus(): Promise<number> {
try {
const done = await this.store.listTasks({ column: "done", slim: true });
const archived = await this.store.listTasks({ column: "archived", slim: true });
const candidates = [...done, ...archived].filter((task) => {
const s = task.status;
return s === "merging" || s === "merging-pr";
});
if (candidates.length === 0) return 0;
let cleared = 0;
for (const task of candidates) {
try {
const previousStatus = task.status;
const updatedAtMs = Date.parse(task.updatedAt ?? "") || Date.now();
const ageMs = Math.max(0, Date.now() - updatedAtMs);
await this.store.updateTask(task.id, { status: null });
await this.recordIntegrityAudit(task.id, "task:auto-recover-stale-merger-status", {
previousColumn: task.column,
previousStatus,
ageMs,
mergeConfirmed: task.mergeDetails?.mergeConfirmed === true,
commitSha: task.mergeDetails?.commitSha ?? null,
});
await this.store.logEntry(
task.id,
`Auto-recovered: cleared stale status="${previousStatus}" on ${task.column} task (age ${Math.round(ageMs / 1000)}s) — was blocking merger queue`,
);
cleared++;
} catch (err: unknown) {
const errorMessage = err instanceof Error ? err.message : String(err);
log.warn(`reconcileStaleMergerStatus: failed for ${task.id}: ${errorMessage}`);
}
}
if (cleared > 0) {
log.warn(`Cleared ${cleared} stale merger-status leak${cleared === 1 ? "" : "s"} on done/archived tasks (FN-5092)`);
}
return cleared;
} catch (err: unknown) {
const errorMessage = err instanceof Error ? err.message : String(err);
log.error(`reconcileStaleMergerStatus failed: ${errorMessage}`);
return 0;
}
}
async finalizeNoOpReviewTasks(): Promise<number> {
try {
const settings = await this.store.getSettings();