fix(engine): recover stuck in-review merges
This commit is contained in:
@@ -148,6 +148,7 @@ function createMockStore() {
|
||||
}),
|
||||
updateTask: vi.fn().mockResolvedValue({}),
|
||||
moveTask: vi.fn().mockResolvedValue({}),
|
||||
mergeTask: vi.fn().mockResolvedValue({}),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
parseStepsFromPrompt: vi.fn().mockResolvedValue([]),
|
||||
updateSettings: vi.fn().mockResolvedValue({}),
|
||||
@@ -5695,6 +5696,39 @@ describe("Invalid transition error handling", () => {
|
||||
expect(onComplete).toHaveBeenCalled();
|
||||
expect(onComplete).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-002" }));
|
||||
});
|
||||
|
||||
it("finalizes an already-reviewed task when it is ready to merge", async () => {
|
||||
const store = createMockStore();
|
||||
store.getTask.mockResolvedValue({
|
||||
id: "FN-003",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status: null,
|
||||
error: null,
|
||||
worktree: "/tmp/test/.worktrees/fn-003",
|
||||
dependencies: [],
|
||||
steps: [{ name: "Done", status: "done" }],
|
||||
workflowStepResults: [{ id: "ws-1", status: "passed", phase: "pre-merge" }],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test");
|
||||
const result = await (executor as any).finalizeAlreadyReviewedTask("FN-003");
|
||||
|
||||
expect(result).toBe("merged");
|
||||
expect(store.mergeTask).toHaveBeenCalledWith("FN-003");
|
||||
expect(store.logEntry).toHaveBeenCalledWith(
|
||||
"FN-003",
|
||||
"Task already in-review after completion — finalizing merge",
|
||||
undefined,
|
||||
undefined,
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
describe("TaskExecutor task_done with summary", () => {
|
||||
|
||||
@@ -3,7 +3,7 @@ import { join } from "node:path";
|
||||
import { existsSync } from "node:fs";
|
||||
import type { TaskStore, Task, TaskDetail, StepStatus, Settings, WorkflowStep, MissionStore, Slice, AgentState, AgentCapability, RunMutationContext } from "@fusion/core";
|
||||
import type { AgentStore } from "@fusion/core";
|
||||
import { buildExecutionMemoryInstructions, resolveAgentPrompt } from "@fusion/core";
|
||||
import { buildExecutionMemoryInstructions, getTaskMergeBlocker, resolveAgentPrompt } from "@fusion/core";
|
||||
import { findWorktreeUser } from "./merger.js";
|
||||
import { generateWorktreeName, slugify } from "./worktree-names.js";
|
||||
import { Type, type Static } from "@mariozechner/pi-ai";
|
||||
@@ -290,6 +290,28 @@ export class TaskExecutor {
|
||||
private loopRecoveryState = new Map<string, { attempts: number; pending: boolean }>();
|
||||
/** Spawned child agent IDs per parent task ID. Used for lifecycle tracking. */
|
||||
private spawnedAgents = new Map<string, Set<string>>();
|
||||
|
||||
private async finalizeAlreadyReviewedTask(taskId: string): Promise<"merged" | "blocked" | "missing"> {
|
||||
const latestTask = await this.store.getTask(taskId);
|
||||
if (!latestTask || latestTask.column !== "in-review") {
|
||||
return "missing";
|
||||
}
|
||||
|
||||
const blocker = getTaskMergeBlocker(latestTask);
|
||||
if (blocker) {
|
||||
await this.store.logEntry(taskId, "Task already in-review; merge deferred", blocker, this.currentRunContext);
|
||||
return "blocked";
|
||||
}
|
||||
|
||||
await this.store.logEntry(
|
||||
taskId,
|
||||
"Task already in-review after completion — finalizing merge",
|
||||
undefined,
|
||||
this.currentRunContext,
|
||||
);
|
||||
await this.store.mergeTask(taskId);
|
||||
return "merged";
|
||||
}
|
||||
/** Child agent sessions keyed by agent ID. Used for termination. */
|
||||
private childSessions = new Map<string, AgentSession>();
|
||||
/** Total count of currently spawned agents (across all parents). */
|
||||
@@ -1585,6 +1607,14 @@ export class TaskExecutor {
|
||||
const logMessage = `Task already moved from '${fromColumn}' — skipping transition to '${toColumn}'`;
|
||||
executorLog.log(`${task.id} ${logMessage}`);
|
||||
await this.store.logEntry(task.id, logMessage, err.message, this.currentRunContext);
|
||||
if (fromColumn === "in-review" && toColumn === "in-review") {
|
||||
try {
|
||||
const finalizeResult = await this.finalizeAlreadyReviewedTask(task.id);
|
||||
executorLog.log(`${task.id} duplicate in-review finalization result: ${finalizeResult}`);
|
||||
} catch (finalizeErr: any) {
|
||||
executorLog.warn(`${task.id} failed to finalize duplicate in-review transition: ${finalizeErr.message}`);
|
||||
}
|
||||
}
|
||||
// Task finished successfully (just already moved), so call onComplete
|
||||
this.options.onComplete?.(task);
|
||||
} else if (this.pausedAborted.has(task.id)) {
|
||||
|
||||
@@ -54,6 +54,7 @@ function createMockStore(overrides: Record<string, unknown> = {}): TaskStore & E
|
||||
updateTask: vi.fn().mockResolvedValue({} as Task),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
moveTask: vi.fn().mockResolvedValue(undefined),
|
||||
mergeTask: vi.fn().mockResolvedValue(undefined),
|
||||
walCheckpoint: vi.fn().mockReturnValue({ busy: 0, log: 5, checkpointed: 5 }),
|
||||
listTasks: vi.fn().mockResolvedValue([]),
|
||||
getRootDir: vi.fn().mockReturnValue("/tmp/test-project"),
|
||||
@@ -688,6 +689,66 @@ describe("SelfHealingManager", () => {
|
||||
});
|
||||
|
||||
describe("recoverMergedReviewTasks", () => {
|
||||
it("merges eligible in-review tasks that still have an unmerged worktree", async () => {
|
||||
const managerWithRecovery = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
});
|
||||
|
||||
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
|
||||
{
|
||||
id: "FN-352",
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status: null,
|
||||
error: null,
|
||||
worktree: "/tmp/test-project/.worktrees/fn-352",
|
||||
steps: [{ name: "Ship it", status: "done" }],
|
||||
workflowStepResults: [{ id: "ws-1", status: "passed", phase: "pre-merge" }],
|
||||
mergeDetails: undefined,
|
||||
log: [],
|
||||
},
|
||||
]);
|
||||
|
||||
const result = await managerWithRecovery.recoverMergeableReviewTasks();
|
||||
|
||||
expect(result).toBe(1);
|
||||
expect(store.mergeTask).toHaveBeenCalledWith("FN-352");
|
||||
expect(store.logEntry).toHaveBeenCalledWith(
|
||||
"FN-352",
|
||||
expect.stringContaining("eligible in-review task was merged"),
|
||||
);
|
||||
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("ignores in-review tasks that are not yet mergeable", async () => {
|
||||
const managerWithRecovery = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
});
|
||||
|
||||
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
|
||||
{
|
||||
id: "FN-353",
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status: null,
|
||||
error: null,
|
||||
worktree: "/tmp/test-project/.worktrees/fn-353",
|
||||
steps: [{ name: "Ship it", status: "in-progress" }],
|
||||
workflowStepResults: [],
|
||||
mergeDetails: undefined,
|
||||
log: [],
|
||||
},
|
||||
]);
|
||||
|
||||
const result = await managerWithRecovery.recoverMergeableReviewTasks();
|
||||
|
||||
expect(result).toBe(0);
|
||||
expect(store.mergeTask).not.toHaveBeenCalled();
|
||||
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("moves merged in-review tasks to done and clears transient merge state", async () => {
|
||||
const managerWithRecovery = new SelfHealingManager(store, {
|
||||
rootDir: "/tmp/test-project",
|
||||
|
||||
@@ -16,7 +16,7 @@
|
||||
import { execSync } from "node:child_process";
|
||||
import { existsSync, readdirSync, statSync } from "node:fs";
|
||||
import { join } from "node:path";
|
||||
import type { TaskStore, Settings, Task } from "@fusion/core";
|
||||
import { getTaskMergeBlocker, type TaskStore, type Settings, type Task } from "@fusion/core";
|
||||
import { createLogger } from "./logger.js";
|
||||
import { scanIdleWorktrees, scanOrphanedBranches } from "./worktree-pool.js";
|
||||
|
||||
@@ -325,6 +325,7 @@ export class SelfHealingManager {
|
||||
this.checkpointWal();
|
||||
await this.enforceWorktreeCap();
|
||||
await this.recoverCompletedTasks();
|
||||
await this.recoverMergeableReviewTasks();
|
||||
await this.recoverMergedReviewTasks();
|
||||
await this.recoverMisclassifiedFailures();
|
||||
await this.recoverOrphanedExecutions();
|
||||
@@ -386,6 +387,55 @@ export class SelfHealingManager {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Recover `in-review` tasks that are fully mergeable but never had
|
||||
* `mergeTask()` invoked.
|
||||
*
|
||||
* This catches races where a task reached review, retained its worktree,
|
||||
* and then got stranded without a merger loop to finish the branch.
|
||||
*
|
||||
* @returns Number of tasks merged or finalized to done
|
||||
*/
|
||||
async recoverMergeableReviewTasks(): Promise<number> {
|
||||
try {
|
||||
const tasks = await this.store.listTasks();
|
||||
|
||||
const mergeable = tasks.filter((t) =>
|
||||
t.column === "in-review" &&
|
||||
Boolean(t.worktree) &&
|
||||
t.mergeDetails?.mergeConfirmed !== true &&
|
||||
getTaskMergeBlocker(t) === undefined,
|
||||
);
|
||||
|
||||
if (mergeable.length === 0) return 0;
|
||||
|
||||
log.warn(`Found ${mergeable.length} mergeable review task(s) stuck in in-review`);
|
||||
|
||||
let recovered = 0;
|
||||
for (const task of mergeable) {
|
||||
try {
|
||||
await this.store.mergeTask(task.id);
|
||||
await this.store.logEntry(
|
||||
task.id,
|
||||
"Auto-recovered: eligible in-review task was merged and moved to done",
|
||||
);
|
||||
log.log(`Recovered mergeable review task ${task.id}: merged to done`);
|
||||
recovered++;
|
||||
} catch (err: any) {
|
||||
log.error(`Failed to recover mergeable review task ${task.id}: ${err.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
if (recovered > 0) {
|
||||
log.log(`Recovered ${recovered} mergeable review task(s) → done`);
|
||||
}
|
||||
return recovered;
|
||||
} catch (err: any) {
|
||||
log.error(`Mergeable review recovery failed: ${err.message}`);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
// ── Misclassified failure recovery ───────────────────────────────
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user