fix(engine): reclaim wedged single-flight merge pump automatically
AI-merge review hangs left activeMergeTaskId/mergeRunning set while status=reviewing and overseer logEntry noise kept updatedAt fresh, so self-healing never reclaimed the owner and the board showed no merging badge. Race merge work with abort, force-abort on pause/reclaim, treat reviewing as merge-active, and recover on merger agent silence; also forward PluginRunner into AI merge so grok-cli merger matches chat.
This commit is contained in:
@@ -53,6 +53,16 @@ vi.mock("../runtimes/in-process-runtime.js", () => ({
|
||||
getRoutineRunner: vi.fn(),
|
||||
getHeartbeatMonitor: vi.fn(),
|
||||
getTriggerScheduler: vi.fn(),
|
||||
configurePrMonitoring: vi.fn(),
|
||||
setActiveMergeTaskIdProvider: vi.fn(),
|
||||
setActiveMergeStartedAtMsProvider: vi.fn(),
|
||||
setActiveMergeAborter: vi.fn(),
|
||||
setMergeEnqueuer: vi.fn(),
|
||||
setMergeActiveClearer: vi.fn(),
|
||||
setMergePendingProvider: vi.fn(),
|
||||
setMergeRequester: vi.fn(),
|
||||
resumeAfterUnpause: vi.fn(async () => undefined),
|
||||
getPluginRunner: vi.fn(() => undefined),
|
||||
};
|
||||
}),
|
||||
}));
|
||||
@@ -1112,4 +1122,88 @@ describe("ProjectEngine merge error recovery", () => {
|
||||
true,
|
||||
);
|
||||
});
|
||||
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:41:
|
||||
Repro for board-wide hung merge pump: active AI merge ignores AbortSignal (wedged tool), operator pauses the card, and without an outer race drainMergeQueue never settles so no later task gets status=merging.
|
||||
*/
|
||||
it("unblocks the merge pump when a paused active merge ignores abort", async () => {
|
||||
const listeners = new Map<string, Set<(...args: unknown[]) => void>>();
|
||||
let wedgedPaused = false;
|
||||
const store = makeStore();
|
||||
store.on = vi.fn((event: string, listener: (...args: unknown[]) => void) => {
|
||||
const set = listeners.get(event) ?? new Set();
|
||||
set.add(listener);
|
||||
listeners.set(event, set);
|
||||
});
|
||||
store.getTask = vi.fn(async (id: string) =>
|
||||
makeTask({
|
||||
id,
|
||||
paused: id === "FN-wedged" ? wedgedPaused : false,
|
||||
status: null,
|
||||
mergeRetries: 0,
|
||||
}),
|
||||
);
|
||||
|
||||
let wedgedStarted = false;
|
||||
const disposeSession = vi.fn();
|
||||
vi.mocked(runAiMerge).mockImplementation(async (...args: unknown[]) => {
|
||||
const taskId = args[2] as string;
|
||||
const options = args[3] as { signal?: AbortSignal; onSession?: (session: { dispose: () => void }) => void };
|
||||
options.onSession?.({ dispose: disposeSession });
|
||||
if (taskId === "FN-wedged") {
|
||||
wedgedStarted = true;
|
||||
// Hang forever — do not observe abort (matches wedged agent tool).
|
||||
await new Promise<never>(() => {});
|
||||
}
|
||||
return {
|
||||
merged: true,
|
||||
task: makeTask({ id: taskId }),
|
||||
branch: `fusion/${taskId.toLowerCase()}`,
|
||||
} as never;
|
||||
});
|
||||
|
||||
const engine = createEngine(store);
|
||||
const privateEngine = engine as unknown as {
|
||||
mergeRunning: boolean;
|
||||
activeMergeTaskId: string | null;
|
||||
mergeActive: Set<string>;
|
||||
enqueueMerge: (taskId: string) => void;
|
||||
};
|
||||
|
||||
await engine.start();
|
||||
privateEngine.enqueueMerge("FN-wedged");
|
||||
|
||||
await vi.waitFor(() => {
|
||||
expect(wedgedStarted).toBe(true);
|
||||
expect(privateEngine.activeMergeTaskId).toBe("FN-wedged");
|
||||
});
|
||||
|
||||
const updatedHandlers = [...(listeners.get("task:updated") ?? [])];
|
||||
expect(updatedHandlers.length).toBeGreaterThan(0);
|
||||
wedgedPaused = true;
|
||||
for (const handler of updatedHandlers) {
|
||||
await handler({ id: "FN-wedged", column: "in-review", paused: true });
|
||||
}
|
||||
|
||||
await vi.waitFor(() => {
|
||||
expect(disposeSession).toHaveBeenCalled();
|
||||
expect(privateEngine.activeMergeTaskId).toBeNull();
|
||||
expect(privateEngine.mergeActive.has("FN-wedged")).toBe(false);
|
||||
expect(privateEngine.mergeRunning).toBe(false);
|
||||
});
|
||||
|
||||
privateEngine.enqueueMerge("FN-next");
|
||||
|
||||
await vi.waitFor(() => {
|
||||
expect(runAiMerge).toHaveBeenCalledWith(
|
||||
expect.anything(),
|
||||
expect.any(String),
|
||||
"FN-next",
|
||||
expect.anything(),
|
||||
);
|
||||
});
|
||||
|
||||
await engine.stop();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
import { readFileSync } from "node:fs";
|
||||
import { dirname, resolve } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import * as fusionCore from "@fusion/core";
|
||||
import { createResolvedAgentSession } from "../agent-session-helpers.js";
|
||||
|
||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||
|
||||
/*
|
||||
FNXC:GrokCliRouting 2026-07-15-09:45:
|
||||
Auto-merge was failing with "Grok CLI models require the bundled Grok CLI runtime" while dashboard chat worked, because project-engine's runAiMerge options omitted pluginRunner. ChatManager already receives engine.getPluginRunner(); the merge door must forward the same runner so createResolvedAgentSession can resolve getRuntimeById("grok") for grok-cli/no-key selections.
|
||||
*/
|
||||
describe("AI merge PluginRunner wiring for Grok CLI", () => {
|
||||
beforeEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it("project-engine merge door forwards this.getPluginRunner() into mergerOptions", () => {
|
||||
const source = readFileSync(resolve(__dirname, "../project-engine.ts"), "utf8");
|
||||
const optionsIndex = source.indexOf("const mergerOptions = {");
|
||||
const pluginRunnerIndex = source.indexOf("pluginRunner: this.getPluginRunner()", optionsIndex);
|
||||
const runAiMergeIndex = source.indexOf("return runAiMerge(store, cwd, taskId, mergeOptionsWithSettings)", optionsIndex);
|
||||
|
||||
expect(optionsIndex).toBeGreaterThanOrEqual(0);
|
||||
expect(pluginRunnerIndex).toBeGreaterThan(optionsIndex);
|
||||
expect(runAiMergeIndex).toBeGreaterThan(pluginRunnerIndex);
|
||||
});
|
||||
|
||||
it("createResolvedAgentSession routes merger grok-cli selections through the provided PluginRunner", async () => {
|
||||
vi.spyOn(fusionCore, "isGrokApiKeyFusionVisible").mockReturnValue(false);
|
||||
|
||||
const createSession = vi.fn().mockResolvedValue({
|
||||
session: {
|
||||
model: "grok-4.5",
|
||||
messages: [],
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
});
|
||||
const grokRuntime = {
|
||||
id: "grok",
|
||||
name: "Grok Runtime",
|
||||
createSession,
|
||||
promptWithFallback: vi.fn(),
|
||||
describeModel: vi.fn(() => "grok/grok-4.5"),
|
||||
};
|
||||
const registration = {
|
||||
pluginId: "fusion-plugin-grok-runtime",
|
||||
runtime: {
|
||||
metadata: { runtimeId: "grok", name: "Grok Runtime" },
|
||||
factory: vi.fn().mockResolvedValue(grokRuntime),
|
||||
},
|
||||
};
|
||||
const pluginRunner = {
|
||||
getPluginRuntimes: vi.fn().mockReturnValue([registration]),
|
||||
getRuntimeById: vi.fn().mockReturnValue(registration),
|
||||
createRuntimeContext: vi.fn().mockResolvedValue({
|
||||
pluginId: "fusion-plugin-grok-runtime",
|
||||
taskStore: {},
|
||||
settings: {},
|
||||
logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() },
|
||||
emitEvent: vi.fn(),
|
||||
}),
|
||||
};
|
||||
|
||||
const result = await createResolvedAgentSession({
|
||||
sessionPurpose: "merger",
|
||||
pluginRunner: pluginRunner as never,
|
||||
cwd: "/tmp/project",
|
||||
defaultProvider: "grok-cli",
|
||||
defaultModelId: "grok-4.5",
|
||||
systemPrompt: "merge",
|
||||
});
|
||||
|
||||
expect(pluginRunner.getRuntimeById).toHaveBeenCalledWith("grok");
|
||||
expect(result.runtimeId).toBe("grok");
|
||||
expect(createSession).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("throws the dual-remediation error for merger grok-cli when pluginRunner is omitted", async () => {
|
||||
vi.spyOn(fusionCore, "isGrokApiKeyFusionVisible").mockReturnValue(false);
|
||||
|
||||
await expect(createResolvedAgentSession({
|
||||
sessionPurpose: "merger",
|
||||
// Intentionally omit pluginRunner — the pre-fix auto-merge wiring bug.
|
||||
cwd: "/tmp/project",
|
||||
defaultProvider: "grok-cli",
|
||||
defaultModelId: "grok-4.5",
|
||||
systemPrompt: "merge",
|
||||
})).rejects.toThrow(/Install and enable the Grok CLI runtime plugin, or set GROK_API_KEY/);
|
||||
});
|
||||
});
|
||||
@@ -52,6 +52,7 @@ const BATCH2_METHODS = [
|
||||
"recoverStaleIncompleteReviewTasks",
|
||||
"recoverReviewTasksWithFailedPreMergeSteps",
|
||||
"recoverInterruptedMergingTasks",
|
||||
"recoverWedgedActiveMerge",
|
||||
"recoverDoneTaskMergeMetadata",
|
||||
"recoverStaleMergingStatus",
|
||||
"finalizeNoOpReviewTasks",
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Settings, TaskStore } from "@fusion/core";
|
||||
import { SelfHealingManager } from "../self-healing.js";
|
||||
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:50:
|
||||
Self-healing must reclaim a wedged in-process active merge when the AI merge review pass hangs (status=reviewing / merger agent silence) so the single-flight pump is not stuck with no merging badge on the board.
|
||||
*/
|
||||
|
||||
function createTask(id: string, overrides: Record<string, unknown> = {}) {
|
||||
return {
|
||||
id,
|
||||
title: id,
|
||||
description: id,
|
||||
column: "in-review",
|
||||
status: "reviewing",
|
||||
paused: false,
|
||||
blockedBy: null,
|
||||
dependencies: [],
|
||||
steps: [{ name: "Ship", status: "done" }],
|
||||
log: [],
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
describe("SelfHealingManager wedged active merge recovery", () => {
|
||||
let tasks: Map<string, Record<string, unknown>>;
|
||||
let store: TaskStore;
|
||||
let agentLogs: Array<{ agent?: string; timestamp?: string; type?: string; text?: string }>;
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date("2026-01-01T01:00:00.000Z"));
|
||||
tasks = new Map();
|
||||
agentLogs = [];
|
||||
|
||||
const wedged = createTask("FN-WEDGE");
|
||||
tasks.set("FN-WEDGE", wedged);
|
||||
|
||||
store = {
|
||||
getSettings: vi.fn().mockResolvedValue({
|
||||
globalPause: false,
|
||||
enginePaused: false,
|
||||
autoMerge: true,
|
||||
taskStuckTimeoutMs: 15 * 60_000,
|
||||
} as unknown as Settings),
|
||||
listTasks: vi.fn().mockImplementation(async (options?: { column?: string }) => {
|
||||
const all = Array.from(tasks.values());
|
||||
if (!options?.column) return all;
|
||||
return all.filter((task) => task.column === options.column);
|
||||
}),
|
||||
getTask: vi.fn().mockImplementation(async (id: string) => tasks.get(id) ?? null),
|
||||
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),
|
||||
getAgentLogs: vi.fn().mockImplementation(async () => agentLogs),
|
||||
getCompletionHandoffAcceptedMarker: vi.fn().mockReturnValue(null),
|
||||
parseFileScopeFromPrompt: vi.fn().mockResolvedValue([]),
|
||||
} as unknown as TaskStore;
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("does not reclaim a live active merge with recent merger agent activity", async () => {
|
||||
agentLogs = [
|
||||
{ agent: "merger", timestamp: "2026-01-01T00:55:00.000Z", type: "tool", text: "bash" },
|
||||
];
|
||||
const abortActiveMerge = vi.fn().mockReturnValue(true);
|
||||
const manager = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
getActiveMergeTaskId: () => "FN-WEDGE",
|
||||
getActiveMergeStartedAtMs: () => Date.parse("2026-01-01T00:00:00.000Z"),
|
||||
abortActiveMerge,
|
||||
});
|
||||
|
||||
const recovered = await manager.recoverWedgedActiveMerge();
|
||||
expect(recovered).toBe(0);
|
||||
expect(abortActiveMerge).not.toHaveBeenCalled();
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("reclaims active merge after merger agent silence past stuck timeout", async () => {
|
||||
// Last merger activity 30 minutes ago; stuck timeout is 15 minutes.
|
||||
agentLogs = [
|
||||
{ agent: "merger", timestamp: "2026-01-01T00:30:00.000Z", type: "tool", text: "fn_task_show" },
|
||||
{ agent: "executor", timestamp: "2026-01-01T00:59:00.000Z", type: "text", text: "noise" },
|
||||
];
|
||||
const abortActiveMerge = vi.fn().mockReturnValue(true);
|
||||
const enqueueMerge = vi.fn().mockReturnValue(true);
|
||||
const clearMergeActive = vi.fn();
|
||||
const manager = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
getActiveMergeTaskId: () => "FN-WEDGE",
|
||||
getActiveMergeStartedAtMs: () => Date.parse("2026-01-01T00:00:00.000Z"),
|
||||
abortActiveMerge,
|
||||
enqueueMerge,
|
||||
clearMergeActive,
|
||||
});
|
||||
|
||||
const recovered = await manager.recoverWedgedActiveMerge();
|
||||
expect(recovered).toBe(1);
|
||||
expect(abortActiveMerge).toHaveBeenCalledWith("FN-WEDGE", "wedged-active-merge-no-merger-progress");
|
||||
expect(clearMergeActive).toHaveBeenCalledWith("FN-WEDGE");
|
||||
expect(enqueueMerge).toHaveBeenCalledWith("FN-WEDGE");
|
||||
expect(tasks.get("FN-WEDGE")?.status).toBeNull();
|
||||
expect(store.logEntry).toHaveBeenCalledWith(
|
||||
"FN-WEDGE",
|
||||
expect.stringContaining("wedged active merge reclaimed"),
|
||||
);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("recoverInterruptedMergingTasks includes AI-merge reviewing status and aborts the live owner", async () => {
|
||||
agentLogs = [
|
||||
{ agent: "merger", timestamp: "2026-01-01T00:20:00.000Z", type: "tool", text: "bash" },
|
||||
];
|
||||
const abortActiveMerge = vi.fn().mockReturnValue(true);
|
||||
const enqueueMerge = vi.fn().mockReturnValue(true);
|
||||
const manager = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
getActiveMergeTaskId: () => "FN-WEDGE",
|
||||
getActiveMergeStartedAtMs: () => Date.parse("2026-01-01T00:00:00.000Z"),
|
||||
abortActiveMerge,
|
||||
enqueueMerge,
|
||||
clearMergeActive: vi.fn(),
|
||||
});
|
||||
|
||||
const recovered = await manager.recoverInterruptedMergingTasks();
|
||||
expect(recovered).toBe(1);
|
||||
expect(abortActiveMerge).toHaveBeenCalledWith("FN-WEDGE", "recover-interrupted-merging-wedged-owner");
|
||||
expect(enqueueMerge).toHaveBeenCalledWith("FN-WEDGE");
|
||||
expect(tasks.get("FN-WEDGE")?.status).toBeNull();
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("reclaims when agent logs are empty but claim wall-clock exceeds stuck timeout", async () => {
|
||||
agentLogs = [];
|
||||
const abortActiveMerge = vi.fn().mockReturnValue(true);
|
||||
const manager = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
getActiveMergeTaskId: () => "FN-WEDGE",
|
||||
// Claimed 40 minutes ago at system time 01:00
|
||||
getActiveMergeStartedAtMs: () => Date.parse("2026-01-01T00:20:00.000Z"),
|
||||
abortActiveMerge,
|
||||
enqueueMerge: vi.fn().mockReturnValue(true),
|
||||
clearMergeActive: vi.fn(),
|
||||
});
|
||||
|
||||
const recovered = await manager.recoverWedgedActiveMerge();
|
||||
expect(recovered).toBe(1);
|
||||
expect(abortActiveMerge).toHaveBeenCalled();
|
||||
manager.stop();
|
||||
});
|
||||
});
|
||||
@@ -460,6 +460,8 @@ export class ProjectEngine {
|
||||
private mergeRunning = false;
|
||||
private activeMergeSession: { dispose: () => void } | null = null;
|
||||
private activeMergeTaskId: string | null = null;
|
||||
/** Wall-clock when `activeMergeTaskId` was claimed; self-healing uses this when agent logs are silent. */
|
||||
private activeMergeStartedAtMs: number | null = null;
|
||||
private mergeAbortController: AbortController | null = null;
|
||||
private mergeRetryTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
private autostashSweepTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
@@ -579,16 +581,14 @@ export class ProjectEngine {
|
||||
//
|
||||
// Tests substitute a minimal runtime mock that may not implement this hook.
|
||||
this.runtime.setActiveMergeTaskIdProvider?.(() => this.getActiveMergeTaskId());
|
||||
this.runtime.setActiveMergeStartedAtMsProvider?.(() => this.activeMergeStartedAtMs);
|
||||
this.runtime.setActiveMergeAborter?.((taskId, reason) => this.abortActiveMerge(taskId, reason));
|
||||
this.runtime.setMergeEnqueuer?.((taskId) => {
|
||||
// If the wedged attempt was the active one, abort its in-flight signal
|
||||
// and dispose its session so subsequent code paths can release file
|
||||
// handles / child processes promptly.
|
||||
if (this.activeMergeTaskId === taskId) {
|
||||
this.mergeAbortController?.abort();
|
||||
this.mergeAbortController = null;
|
||||
this.activeMergeSession?.dispose();
|
||||
this.activeMergeSession = null;
|
||||
this.activeMergeTaskId = null;
|
||||
this.abortActiveMerge(taskId, "merge-enqueuer-reclaim");
|
||||
}
|
||||
this.mergeActive.delete(taskId);
|
||||
return this.internalEnqueueMerge(taskId);
|
||||
@@ -610,6 +610,43 @@ export class ProjectEngine {
|
||||
return this.activeMergeTaskId;
|
||||
}
|
||||
|
||||
getActiveMergeStartedAtMs(): number | null {
|
||||
return this.activeMergeStartedAtMs;
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:50:
|
||||
Self-healing reclaim path for a wedged single-flight merge. Always abort + dispose + clear identity so raceMergeWithAbort can settle drainMergeQueue even when the agent ignores cooperative abort.
|
||||
*/
|
||||
abortActiveMerge(taskId: string, reason: string): boolean {
|
||||
if (this.activeMergeTaskId !== taskId) return false;
|
||||
runtimeLog.log(`Aborting active merge for ${taskId} (${reason})`);
|
||||
this.mergeAbortController?.abort();
|
||||
this.mergeAbortController = null;
|
||||
if (this.activeMergeSession) {
|
||||
this.activeMergeSession.dispose();
|
||||
this.activeMergeSession = null;
|
||||
}
|
||||
this.mergeActive.delete(taskId);
|
||||
this.activeMergeTaskId = null;
|
||||
this.activeMergeStartedAtMs = null;
|
||||
return true;
|
||||
}
|
||||
|
||||
private claimActiveMerge(taskId: string): AbortSignal {
|
||||
this.activeMergeTaskId = taskId;
|
||||
this.activeMergeStartedAtMs = Date.now();
|
||||
this.mergeAbortController = new AbortController();
|
||||
return this.mergeAbortController.signal;
|
||||
}
|
||||
|
||||
private clearActiveMergeClaim(taskId: string): void {
|
||||
if (this.activeMergeTaskId === taskId) {
|
||||
this.activeMergeTaskId = null;
|
||||
this.activeMergeStartedAtMs = null;
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:Workspace 2026-06-22-16:40 (Phase D P1 TOCTOU — merge-queue dispatch blind spot):
|
||||
A workspace task is "merge-pending" if it sits ANYWHERE in this engine's in-memory merge
|
||||
@@ -1040,6 +1077,7 @@ export class ProjectEngine {
|
||||
this.mergeAbortController?.abort();
|
||||
this.mergeAbortController = null;
|
||||
this.activeMergeTaskId = null;
|
||||
this.activeMergeStartedAtMs = null;
|
||||
this.pausedReviewTaskIds.clear();
|
||||
|
||||
const queuedTaskIds = [...this.mergeQueue];
|
||||
@@ -2413,6 +2451,39 @@ export class ProjectEngine {
|
||||
return true;
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:41:
|
||||
Operator-visible hang: pause/cancel aborted the merge AbortController and disposed the session, but runAiMerge/promptWithFallback often keeps awaiting a wedged agent tool (e.g. fn_task_show) that never observes the signal. drainMergeQueue stayed parked on that await with mergeRunning=true, so no card got a merging badge and later enqueues no-op'd. Race the merge work with the abort signal so pause always unblocks the single-flight pump even when the agent ignores abort; dispose remains best-effort cleanup for the orphan session.
|
||||
*/
|
||||
private createMergeAbortedError(taskId: string): Error {
|
||||
// Name-tagged Error (not a class import) so test mocks of merger.js stay compatible and catch paths that match err.name === "MergeAbortedError" keep working.
|
||||
const err = new Error(`Merge aborted for ${taskId}: pause or cancel requested`);
|
||||
err.name = "MergeAbortedError";
|
||||
return err;
|
||||
}
|
||||
|
||||
private raceMergeWithAbort<T>(work: Promise<T>, signal: AbortSignal, taskId: string): Promise<T> {
|
||||
if (signal.aborted) {
|
||||
return Promise.reject(this.createMergeAbortedError(taskId));
|
||||
}
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
const onAbort = (): void => {
|
||||
reject(this.createMergeAbortedError(taskId));
|
||||
};
|
||||
signal.addEventListener("abort", onAbort, { once: true });
|
||||
work.then(
|
||||
(value) => {
|
||||
signal.removeEventListener("abort", onAbort);
|
||||
resolve(value);
|
||||
},
|
||||
(err: unknown) => {
|
||||
signal.removeEventListener("abort", onAbort);
|
||||
reject(err);
|
||||
},
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Filter a sweep's listTasks() result to merge-eligible tasks, sort by
|
||||
* priority (urgent → low, then createdAt ASC, then id ASC), and enqueue.
|
||||
@@ -3339,7 +3410,7 @@ export class ProjectEngine {
|
||||
const routeWorkspaceDirect = !!mergeCandidate && isWorkspaceTask(mergeCandidate);
|
||||
|
||||
if (mergeStrategy === "pull-request" && this.options.processPullRequestMerge && !routeWorkspaceDirect) {
|
||||
this.activeMergeTaskId = taskId;
|
||||
this.claimActiveMerge(taskId);
|
||||
runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge processing PR flow for ${taskId}...`);
|
||||
const result = await this.options.processPullRequestMerge(
|
||||
store,
|
||||
@@ -3391,19 +3462,28 @@ export class ProjectEngine {
|
||||
const usageLimitPauser = (this.runtime as any).usageLimitPauser;
|
||||
|
||||
const rawMerge = async () => {
|
||||
this.activeMergeTaskId = taskId;
|
||||
this.mergeAbortController = new AbortController();
|
||||
const abortSignal = this.claimActiveMerge(taskId);
|
||||
/*
|
||||
FNXC:GrokCliRouting 2026-07-15-09:45:
|
||||
AI merge creates sessions via createResolvedAgentSession with the same Grok CLI no-visible-key auto-derive seam as chat/executor. Without pluginRunner, getRuntimeById("grok") is unavailable and grok-cli merger/fallback selections throw "Grok CLI models require the bundled Grok CLI runtime" even when chat works (ChatManager already receives engine.getPluginRunner()). Forward the engine PluginRunner so merge can route to the logged-in grok CLI like every other lane.
|
||||
*/
|
||||
const mergerOptions = {
|
||||
manual: hasManualResolver,
|
||||
pool,
|
||||
usageLimitPauser,
|
||||
agentStore,
|
||||
signal: this.mergeAbortController.signal,
|
||||
pluginRunner: this.getPluginRunner(),
|
||||
signal: abortSignal,
|
||||
syncGroupPr: this.options.syncGroupPr,
|
||||
onSession: (session: { dispose: () => void }) => {
|
||||
this.activeMergeSession = session;
|
||||
},
|
||||
};
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:41:
|
||||
Always race the merge body with the pause/cancel abort signal. Cooperative abort inside runAiMerge is best-effort; without this outer race a wedged agent tool parks drainMergeQueue forever (no merging badge board-wide).
|
||||
*/
|
||||
return this.raceMergeWithAbort((async () => {
|
||||
// FNXC:Workspace 2026-06-21-23:40 (Phase C U1, KTD2):
|
||||
// Engine merge dispatch door. A workspace-mode task (non-empty
|
||||
// `workspaceWorktrees`) routes to the per-repo merge loop
|
||||
@@ -3486,6 +3566,7 @@ export class ProjectEngine {
|
||||
allowDirtyLocalCheckoutSync: settings.merger?.allowDirtyLocalCheckoutSync === true,
|
||||
};
|
||||
return runAiMerge(store, cwd, taskId, mergeOptionsWithSettings);
|
||||
})(), abortSignal, taskId);
|
||||
};
|
||||
|
||||
let result: MergeResult;
|
||||
@@ -4269,9 +4350,7 @@ export class ProjectEngine {
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (this.activeMergeTaskId === taskId) {
|
||||
this.activeMergeTaskId = null;
|
||||
}
|
||||
this.clearActiveMergeClaim(taskId);
|
||||
this.mergeAbortController = null;
|
||||
this.mergeActive.delete(taskId);
|
||||
// If a manual merge was requested while this task was already in-flight,
|
||||
@@ -4493,16 +4572,11 @@ export class ProjectEngine {
|
||||
return;
|
||||
}
|
||||
|
||||
runtimeLog.log(`Paused in-review task interrupting active merge: ${task.id}`);
|
||||
this.mergeAbortController?.abort();
|
||||
this.mergeAbortController = null;
|
||||
|
||||
if (this.activeMergeSession) {
|
||||
this.activeMergeSession.dispose();
|
||||
this.activeMergeSession = null;
|
||||
}
|
||||
|
||||
this.mergeActive.delete(task.id);
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:41:
|
||||
Pause of the active merge must free the single-flight lane the same way soft-delete does. Abort + dispose alone left activeMergeTaskId set and, when the agent ignored abort, drainMergeQueue wedged with mergeRunning=true (no merging badge on any card). Clear identity now; raceMergeWithAbort rejects the parked await so the drain finally settles and later enqueues can start.
|
||||
*/
|
||||
this.abortActiveMerge(task.id, "task-paused");
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -4547,17 +4621,7 @@ export class ProjectEngine {
|
||||
return;
|
||||
}
|
||||
|
||||
runtimeLog.log(`Soft-deleted task interrupting active merge: ${task.id}`);
|
||||
this.mergeAbortController?.abort();
|
||||
this.mergeAbortController = null;
|
||||
|
||||
if (this.activeMergeSession) {
|
||||
this.activeMergeSession.dispose();
|
||||
this.activeMergeSession = null;
|
||||
}
|
||||
|
||||
this.mergeActive.delete(task.id);
|
||||
this.activeMergeTaskId = null;
|
||||
this.abortActiveMerge(task.id, "task-soft-deleted");
|
||||
};
|
||||
|
||||
store.on("task:updated", this.taskUpdatedHandler);
|
||||
|
||||
@@ -245,6 +245,8 @@ export class InProcessRuntime
|
||||
) => Promise<import("@fusion/core").MergeResult>;
|
||||
private clearMergeActive?: (taskId: string) => void;
|
||||
private activeMergeTaskIdProvider?: () => string | null;
|
||||
private activeMergeStartedAtMsProvider?: () => number | null;
|
||||
private activeMergeAborter?: (taskId: string, reason: string) => boolean;
|
||||
/**
|
||||
* FNXC:Workspace 2026-06-22-16:40 (Phase D P1 TOCTOU): predicate that reports whether a task is
|
||||
* anywhere in ProjectEngine's in-memory merge pipeline (queued OR dequeued-and-merging). Set by
|
||||
@@ -1050,6 +1052,10 @@ export class InProcessRuntime
|
||||
isTaskActive: (taskId: string) => this.executor.isTaskActive(taskId),
|
||||
clearMergeActive: this.clearMergeActive ? (taskId: string) => this.clearMergeActive?.(taskId) : undefined,
|
||||
getActiveMergeTaskId: () => this.activeMergeTaskIdProvider?.() ?? null,
|
||||
getActiveMergeStartedAtMs: () => this.activeMergeStartedAtMsProvider?.() ?? null,
|
||||
abortActiveMerge: this.activeMergeAborter
|
||||
? (taskId: string, reason: string) => this.activeMergeAborter?.(taskId, reason) ?? false
|
||||
: undefined,
|
||||
// FNXC:Workspace 2026-06-22-16:40 (Phase D P1 TOCTOU): undefined provider → "not pending"
|
||||
// (graceful when unwired; existing guards still apply). In production it is always wired.
|
||||
isMergePending: this.mergePendingProvider ? (taskId: string) => this.mergePendingProvider?.(taskId) ?? false : undefined,
|
||||
@@ -1498,6 +1504,14 @@ export class InProcessRuntime
|
||||
this.activeMergeTaskIdProvider = getActiveMergeTaskId;
|
||||
}
|
||||
|
||||
setActiveMergeStartedAtMsProvider(getActiveMergeStartedAtMs: () => number | null): void {
|
||||
this.activeMergeStartedAtMsProvider = getActiveMergeStartedAtMs;
|
||||
}
|
||||
|
||||
setActiveMergeAborter(abortActiveMerge: (taskId: string, reason: string) => boolean): void {
|
||||
this.activeMergeAborter = abortActiveMerge;
|
||||
}
|
||||
|
||||
setMergePendingProvider(isMergePending: (taskId: string) => boolean): void {
|
||||
this.mergePendingProvider = isMergePending;
|
||||
}
|
||||
|
||||
@@ -368,6 +368,19 @@ export interface SelfHealingOptions {
|
||||
* Used to avoid clearing a transient merge status mid-merge.
|
||||
*/
|
||||
getActiveMergeTaskId?: () => string | null;
|
||||
/**
|
||||
* Wall-clock ms when the current active merge session was claimed
|
||||
* (`activeMergeTaskId` set). Used to reclaim wedged merges without trusting
|
||||
* `task.updatedAt` (overseer/logEntry refresh that field continuously).
|
||||
*/
|
||||
getActiveMergeStartedAtMs?: () => number | null;
|
||||
/**
|
||||
* Force-abort the in-process active merge for `taskId` (AbortController +
|
||||
* session dispose + clear identity) so drainMergeQueue can settle even when
|
||||
* the agent ignores cooperative abort. Returns true when this task was the
|
||||
* active owner and abort was issued.
|
||||
*/
|
||||
abortActiveMerge?: (taskId: string, reason: string) => boolean;
|
||||
/*
|
||||
FNXC:Workspace 2026-06-22-16:40 (Phase D P1 TOCTOU — merge-queue dispatch blind spot):
|
||||
Returns true if the task is ANYWHERE in ProjectEngine's in-memory merge pipeline — queued in
|
||||
@@ -436,7 +449,11 @@ const ORPHANED_EXECUTION_RECOVERY_GRACE_MS = 60_000;
|
||||
*/
|
||||
const LEAKED_WORKTREE_SLOT_GRACE_MS = 60_000;
|
||||
export const VALIDATOR_RUN_STALE_MAX_AGE_MS = 6 * 60 * 60 * 1000;
|
||||
const ACTIVE_MERGE_STATUSES = new Set(["merging", "merging-pr", "merging-fix"]);
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:50:
|
||||
AI merge sets status="reviewing" during the clean-room review pass (merger-ai mergeAndReview). That is still live merge activity. Without it here, recoverInterruptedMergingTasks ignored hung review-phase merges and the single-flight pump stayed wedged while the board showed no merging badge.
|
||||
*/
|
||||
const ACTIVE_MERGE_STATUSES = new Set(["merging", "merging-pr", "merging-fix", "reviewing"]);
|
||||
const NON_TERMINAL_STEP_STATUSES = new Set(["pending", "in-progress"]);
|
||||
const STRANDED_COMPLETED_TODO_ACTIVE_STATUSES = new Set([
|
||||
"in-progress",
|
||||
@@ -1358,6 +1375,7 @@ export class SelfHealingManager {
|
||||
{ name: "failed-pre-merge-steps", fn: () => this.recoverReviewTasksWithFailedPreMergeSteps().then(() => undefined) },
|
||||
{ name: "missing-worktree-review-failures", fn: () => this.recoverMissingWorktreeReviewFailures().then(() => undefined) },
|
||||
{ name: "interrupted-merging", fn: () => this.recoverInterruptedMergingTasks().then(() => undefined) },
|
||||
{ name: "wedged-active-merge", fn: () => this.recoverWedgedActiveMerge().then(() => undefined) },
|
||||
{ name: "transient-merge-failures", fn: () => this.recoverTransientMergeFailures().then(() => undefined) },
|
||||
{ name: "done-merge-metadata", fn: () => this.recoverDoneTaskMergeMetadata().then(() => undefined) },
|
||||
{ name: "reconcile-done-task-integrity", fn: () => this.reconcileDoneTaskIntegrity().then(() => undefined) },
|
||||
@@ -1872,6 +1890,53 @@ export class SelfHealingManager {
|
||||
return Date.now() - updatedAt >= timeoutMs;
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:50:
|
||||
Do not use task.updatedAt as the sole stall clock for an in-flight merge. Overseer "progressing" lines and other logEntry writes refresh updatedAt while the merger agent is completely wedged (e.g. hung fn_task_show), so age gates never fire. Prefer last merger agent-log timestamp, then active-merge claim wall clock, then updatedAt only as last resort.
|
||||
*/
|
||||
private async getLastMergerAgentActivityMs(taskId: string): Promise<number | null> {
|
||||
if (typeof this.store.getAgentLogs !== "function") return null;
|
||||
try {
|
||||
const logs = await this.store.getAgentLogs(taskId, { limit: 80 });
|
||||
if (!Array.isArray(logs) || logs.length === 0) return null;
|
||||
for (let i = logs.length - 1; i >= 0; i--) {
|
||||
const entry = logs[i] as { agent?: string; timestamp?: string };
|
||||
if (entry.agent !== "merger" || !entry.timestamp) continue;
|
||||
const ms = Date.parse(entry.timestamp);
|
||||
if (Number.isFinite(ms) && ms > 0) return ms;
|
||||
}
|
||||
return null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private async isActiveMergeWedged(taskId: string, timeoutMs: number): Promise<boolean> {
|
||||
if (!Number.isFinite(timeoutMs) || timeoutMs <= 0) return false;
|
||||
const activeId = this.options.getActiveMergeTaskId?.() ?? null;
|
||||
if (activeId !== taskId) return false;
|
||||
|
||||
const now = Date.now();
|
||||
const lastMergerMs = await this.getLastMergerAgentActivityMs(taskId);
|
||||
if (lastMergerMs != null) {
|
||||
return now - lastMergerMs >= timeoutMs;
|
||||
}
|
||||
const startedAt = this.options.getActiveMergeStartedAtMs?.() ?? null;
|
||||
if (startedAt != null && Number.isFinite(startedAt) && startedAt > 0) {
|
||||
return now - startedAt >= timeoutMs;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
private async isPastInterruptedMergeGraceAsync(task: Task, timeoutMs: number): Promise<boolean> {
|
||||
const activeId = this.options.getActiveMergeTaskId?.() ?? null;
|
||||
if (activeId === task.id) {
|
||||
// Live owner: require merger silence / claim-age, not updatedAt (overseer noise).
|
||||
return this.isActiveMergeWedged(task.id, timeoutMs);
|
||||
}
|
||||
return this.isPastInterruptedMergeGrace(task, timeoutMs);
|
||||
}
|
||||
|
||||
private async findLandedTaskCommit(
|
||||
task: Task,
|
||||
options?: { preferEarliestOwnedCommit?: boolean },
|
||||
@@ -2567,6 +2632,7 @@ export class SelfHealingManager {
|
||||
{ name: "recover-failed-pre-merge-steps", fn: () => this.recoverReviewTasksWithFailedPreMergeSteps() },
|
||||
{ name: "recover-missing-worktree-review-failures", fn: () => this.recoverMissingWorktreeReviewFailures() },
|
||||
{ name: "recover-interrupted-merging", fn: () => this.recoverInterruptedMergingTasks() },
|
||||
{ name: "recover-wedged-active-merge", fn: () => this.recoverWedgedActiveMerge() },
|
||||
{ name: "recover-transient-merge-failures", fn: () => this.recoverTransientMergeFailures() },
|
||||
{ name: "recover-done-merge-metadata", fn: () => this.recoverDoneTaskMergeMetadata() },
|
||||
{ name: "recover-stale-merging-status", fn: () => this.recoverStaleMergingStatus() },
|
||||
@@ -2892,7 +2958,10 @@ export class SelfHealingManager {
|
||||
const tasks = await this.store.listTasks({ column: "in-review", slim: true });
|
||||
const stale = tasks.filter((task) => {
|
||||
if (task.column !== "in-review" || task.paused) return false;
|
||||
if (!task.status || (task.status !== "merging" && task.status !== "merging-pr")) return false;
|
||||
// Include AI-merge "reviewing" (same ACTIVE_MERGE_STATUSES as interrupted recovery).
|
||||
if (!task.status || !ACTIVE_MERGE_STATUSES.has(task.status)) return false;
|
||||
// Live owner is reclaimed by recoverWedgedActiveMerge / recoverInterruptedMergingTasks
|
||||
// using merger-silence clocks; this path only clears orphan status without a live owner.
|
||||
if (activeMergeTaskId && activeMergeTaskId === task.id) return false;
|
||||
|
||||
const updatedAtMs = task.updatedAt ? Date.parse(task.updatedAt) : Number.NaN;
|
||||
@@ -2937,6 +3006,79 @@ export class SelfHealingManager {
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:50:
|
||||
Reclaim the in-process single-flight merge owner when it is wedged with no merger agent progress. Covers:
|
||||
- status="reviewing" (AI merge review pass) that recoverStaleMergingStatus skipped because getActiveMergeTaskId still pointed at the owner
|
||||
- status already cleared/null while activeMergeTaskId + mergeRunning still hold the pump (board shows no merging badge, nothing starts)
|
||||
- overseer logEntry noise keeping updatedAt fresh so status-age gates never fire
|
||||
Uses merger agent-log silence (preferred) or active-merge claim wall clock — never updatedAt alone.
|
||||
*/
|
||||
async recoverWedgedActiveMerge(): Promise<number> {
|
||||
try {
|
||||
const settings = await this.store.getSettings();
|
||||
if (settings.globalPause || settings.enginePaused) return 0;
|
||||
const timeoutMs = settings.taskStuckTimeoutMs;
|
||||
if (!timeoutMs || timeoutMs <= 0) return 0;
|
||||
|
||||
const activeId = this.options.getActiveMergeTaskId?.() ?? null;
|
||||
if (!activeId) return 0;
|
||||
|
||||
if (!(await this.isActiveMergeWedged(activeId, timeoutMs))) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
const task = await this.store.getTask(activeId).catch(() => null);
|
||||
if (task?.paused) {
|
||||
// Pause should already abort; if identity remains, force-clear the lane.
|
||||
const aborted = this.options.abortActiveMerge?.(activeId, "wedged-active-merge-while-paused") ?? false;
|
||||
if (aborted) {
|
||||
log.warn(`Force-aborted wedged active merge ${activeId} that remained after pause`);
|
||||
return 1;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
if (task && task.column !== "in-review") {
|
||||
const aborted = this.options.abortActiveMerge?.(activeId, "wedged-active-merge-left-in-review") ?? false;
|
||||
if (aborted) {
|
||||
log.warn(`Force-aborted wedged active merge ${activeId}: task column is ${task.column}`);
|
||||
return 1;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
log.warn(`Reclaiming wedged active merge ${activeId} (no merger agent progress past stuck timeout)`);
|
||||
this.options.abortActiveMerge?.(activeId, "wedged-active-merge-no-merger-progress");
|
||||
this.options.clearMergeActive?.(activeId);
|
||||
if (task) {
|
||||
if (task.status && ACTIVE_MERGE_STATUSES.has(task.status)) {
|
||||
await this.store.updateTask(activeId, { status: null, error: null }).catch(() => undefined);
|
||||
}
|
||||
await this.store
|
||||
.logEntry(
|
||||
activeId,
|
||||
"Auto-recovered: wedged active merge reclaimed after merger agent silence — merge will be retried",
|
||||
)
|
||||
.catch(() => undefined);
|
||||
}
|
||||
if (task && allowsAutoMergeProcessing(task, settings) && !task.paused && task.column === "in-review") {
|
||||
try {
|
||||
this.options.enqueueMerge?.(activeId);
|
||||
} catch (enqueueErr: unknown) {
|
||||
log.warn(
|
||||
`Failed to re-enqueue ${activeId} after wedged-active-merge recovery: ${enqueueErr instanceof Error ? enqueueErr.message : String(enqueueErr)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
return 1;
|
||||
} catch (err: unknown) {
|
||||
const errorMessage = err instanceof Error ? err.message : String(err);
|
||||
log.error(`Wedged active merge recovery failed: ${errorMessage}`);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
async reclaimPrConflicts(): Promise<number> {
|
||||
const tasks = await this.store.listTasks({ slim: true });
|
||||
const candidates = tasks.filter((task) => {
|
||||
@@ -7472,13 +7614,18 @@ export class SelfHealingManager {
|
||||
if (!timeoutMs || timeoutMs <= 0) return 0;
|
||||
|
||||
const tasks = await this.store.listTasks({ column: "in-review", slim: true });
|
||||
const candidates = tasks.filter((task) =>
|
||||
const statusCandidates = tasks.filter((task) =>
|
||||
task.column === "in-review" &&
|
||||
allowsAutoMergeProcessing(task, settings) &&
|
||||
!task.paused &&
|
||||
Boolean(task.status && ACTIVE_MERGE_STATUSES.has(task.status)) &&
|
||||
this.isPastInterruptedMergeGrace(task, timeoutMs),
|
||||
Boolean(task.status && ACTIVE_MERGE_STATUSES.has(task.status)),
|
||||
);
|
||||
const candidates: Task[] = [];
|
||||
for (const task of statusCandidates) {
|
||||
if (await this.isPastInterruptedMergeGraceAsync(task, timeoutMs)) {
|
||||
candidates.push(task);
|
||||
}
|
||||
}
|
||||
|
||||
if (candidates.length === 0) return 0;
|
||||
|
||||
@@ -7487,6 +7634,13 @@ export class SelfHealingManager {
|
||||
let recovered = 0;
|
||||
for (const task of candidates) {
|
||||
try {
|
||||
/*
|
||||
FNXC:MergeQueue 2026-07-15-09:50:
|
||||
If this task still owns the in-process active merge, force-abort first. Cooperative abort alone left mergeRunning=true when the agent ignored the signal; abortActiveMerge + raceMergeWithAbort free the single-flight drain so re-enqueue can actually start.
|
||||
*/
|
||||
if (this.options.getActiveMergeTaskId?.() === task.id) {
|
||||
this.options.abortActiveMerge?.(task.id, "recover-interrupted-merging-wedged-owner");
|
||||
}
|
||||
/*
|
||||
FNXC:Workspace 2026-06-22-09:30 (Phase D U1, KTD1 — P0 workspace gate):
|
||||
A workspace task lands PER-REPO and `landWorkspaceTask` sets status:"merging". The
|
||||
|
||||
Reference in New Issue
Block a user