fix(engine): add cross-process merge guard to prevent concurrent merges
Multiple engine processes (dashboard + serve) share the same SQLite database but each has its own in-memory merge queue. Without a cross-process check, two processes can start merging different tasks simultaneously. Added store.getActiveMergingTask() as a DB-level check before any merge starts. The drainMergeQueue defers with pollIntervalMs delay, and both aiMergeTask and processPullRequestMergeTask have safety-net checks. Also moved stale merge status cleanup to run regardless of autoMerge setting. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -1298,6 +1298,14 @@ export async function aiMergeTask(
|
||||
}
|
||||
|
||||
// 5. Execute merge with retry logic
|
||||
// Cross-process safety net: abort if another task is already mid-merge.
|
||||
// The engine's drainMergeQueue also checks, but this catches direct callers.
|
||||
const activeMerge = store.getActiveMergingTask(taskId);
|
||||
if (activeMerge) {
|
||||
throw new Error(
|
||||
`Cannot merge ${taskId}: task ${activeMerge} is already merging (cross-process conflict)`,
|
||||
);
|
||||
}
|
||||
await store.updateTask(taskId, { status: "merging" });
|
||||
|
||||
// Normalize explicit verification commands from settings
|
||||
|
||||
@@ -471,6 +471,36 @@ export class ProjectEngine {
|
||||
}
|
||||
|
||||
const settings = await store.getSettings();
|
||||
|
||||
// Cross-process guard: check if another process is already merging a
|
||||
// task for this project. The in-memory mergeQueue serializes within
|
||||
// this process, but multiple processes (e.g. dashboard + serve) share
|
||||
// the same SQLite database and can race.
|
||||
const activeMergingTask = store.getActiveMergingTask(taskId);
|
||||
if (activeMergingTask) {
|
||||
const retryMs = settings.pollIntervalMs ?? 15_000;
|
||||
runtimeLog.log(
|
||||
`Merge deferred for ${taskId} — ${activeMergingTask} is already merging (cross-process guard, retry in ${retryMs / 1000}s)`,
|
||||
);
|
||||
// Temporarily remove the manual resolver so the finally block
|
||||
// doesn't prematurely resolve it. The re-enqueue will restore it.
|
||||
if (manualResolver) {
|
||||
this.manualMergeResolvers.delete(taskId);
|
||||
}
|
||||
// Re-queue after the poll interval so we retry once the other merge finishes
|
||||
setTimeout(() => {
|
||||
if (this.shuttingDown) {
|
||||
manualResolver?.reject(new Error("Engine shutting down"));
|
||||
return;
|
||||
}
|
||||
if (manualResolver) {
|
||||
this.manualMergeResolvers.set(taskId, manualResolver);
|
||||
}
|
||||
this.internalEnqueueMerge(taskId);
|
||||
}, retryMs);
|
||||
continue;
|
||||
}
|
||||
|
||||
const mergeStrategy = this.options.getMergeStrategy?.(settings) ?? "direct";
|
||||
|
||||
if (mergeStrategy === "pull-request" && this.options.processPullRequestMerge) {
|
||||
@@ -679,14 +709,13 @@ export class ProjectEngine {
|
||||
|
||||
private async startupMergeSweep(store: TaskStore): Promise<void> {
|
||||
try {
|
||||
const settings = await store.getSettings();
|
||||
if (!settings.autoMerge) return;
|
||||
|
||||
const tasks = await store.listTasks({ column: "in-review" });
|
||||
|
||||
// Clear stale "merging"/"merging-pr" statuses left by a prior crash.
|
||||
// No merge is actually running at startup, so any task still marked
|
||||
// as merging is a leftover from a previous engine lifecycle.
|
||||
// This runs unconditionally (regardless of autoMerge setting) because
|
||||
// stale statuses block manual merges too.
|
||||
const staleStatuses = new Set(["merging", "merging-pr"]);
|
||||
for (const t of tasks) {
|
||||
if (t.status && staleStatuses.has(t.status)) {
|
||||
@@ -697,6 +726,9 @@ export class ProjectEngine {
|
||||
}
|
||||
}
|
||||
|
||||
const settings = await store.getSettings();
|
||||
if (!settings.autoMerge) return;
|
||||
|
||||
const eligible = tasks.filter((t) => this.canMergeTask(t as any));
|
||||
if (eligible.length > 0) {
|
||||
runtimeLog.log(`Auto-merge startup sweep: enqueueing ${eligible.length} task(s)`);
|
||||
|
||||
@@ -145,6 +145,7 @@ function createMockStore(overrides: Record<string, any> = {}) {
|
||||
return makeTask("FN-NEW", "triage");
|
||||
}),
|
||||
deleteTask: vi.fn().mockResolvedValue(undefined),
|
||||
getActiveMergingTask: vi.fn().mockReturnValue(undefined),
|
||||
_trigger(event: string, ...args: any[]) {
|
||||
for (const fn of listeners.get(event) || []) fn(...args);
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user