From 8bb02f7d3cb6a1a198cd764260a25c0cf2f0aa91 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Wed, 3 Jun 2026 17:53:40 -0700 Subject: [PATCH] fix(engine,core): merge-seam multi-waiter + autoMerge gate; selection lock; settings toggle MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Resolve the two needs-human findings from PR #1363 review, plus surface the flag. Merge seam (project-engine.ts): - manualMergeResolvers is now a per-task LIST of waiters. Both the dashboard "merge now" path and the interpreter merge seam call onMerge, so a single resolver per task let the second caller overwrite (and strand) the first. All resolve/reject/requeue/late-resolver/shutdown sites drain the whole list. - New requestInterpreterMerge() honors auto-merge eligibility: when autoMerge is off (or the task isn't merge-ready) it returns merged:false instead of forcing the merge, so a graph merge node can't override an autoMerge-off project — it parks the task in review for a human. setMergeRequester now wires the interpreter to this gate rather than the human bypass. Selection race (store.ts): - selectTaskWorkflow/clearTaskWorkflowSelection now hold one withTaskLock across their whole mutate sequence. Extracted updateTaskUnlocked() (the per-task lock is non-reentrant, so they couldn't wrap the public updateTask without deadlocking) and call that inside the lock. Settings: - Add "Workflow Graph Engine (run custom workflows)" to the Experimental Features list so the workflowGraphExecutor flag is a labeled toggle in Settings → Experimental, not just a raw key. Tests: interpreter-merge-seam.test.ts (multi-waiter resolve/reject + autoMerge eligibility gate); existing merge lifecycle/bypass/selection suites still pass. Co-Authored-By: Claude Opus 4.8 (1M context) --- packages/core/src/store.ts | 93 ++++++---- .../app/components/SettingsModal.tsx | 1 + .../__tests__/interpreter-merge-seam.test.ts | 88 ++++++++++ packages/engine/src/project-engine.ts | 162 ++++++++++++------ 4 files changed, 259 insertions(+), 85 deletions(-) create mode 100644 packages/engine/src/__tests__/interpreter-merge-seam.test.ts diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index a5f33e2865..443d3c9e2c 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -6000,7 +6000,24 @@ export class TaskStore extends EventEmitter { updates: { title?: string; description?: string; priority?: TaskPriority | null; prompt?: string; worktree?: string | null; status?: string | null; dependencies?: string[]; steps?: import("./types.js").TaskStep[]; currentStep?: number; blockedBy?: string | null; overlapBlockedBy?: string | null; assignedAgentId?: string | null; pausedByAgentId?: string | null; pausedReason?: string | null; tokenBudgetSoftAlertedAt?: string | null; worktrunkFallbackAlertedAt?: string | null; worktrunkFailure?: import("./types.js").Task["worktrunkFailure"] | null; tokenBudgetHardAlertedAt?: string | null; tokenBudgetOverride?: import("./types.js").TaskTokenBudgetOverride | null; dispatchStormCount?: number | null; lastDispatchAt?: string | null; assigneeUserId?: string | null; scopeOverride?: boolean | null; scopeOverrideReason?: string | null; scopeAutoWiden?: string[] | null; nodeId?: string | null; effectiveNodeId?: string | null; effectiveNodeSource?: string | null; checkedOutBy?: string | null; checkedOutAt?: string | null; checkoutNodeId?: string | null; checkoutRunId?: string | null; checkoutLeaseRenewedAt?: string | null; checkoutLeaseEpoch?: number | null; paused?: boolean; baseBranch?: string | null; autoMerge?: boolean | null; branch?: string | null; executionStartBranch?: string | null; baseCommitSha?: string | null; size?: "S" | "M" | "L"; reviewLevel?: number; executionMode?: import("./types.js").ExecutionMode | null; mergeRetries?: number; workflowStepRetries?: number; stuckKillCount?: number | null; resumeLimboCount?: number | null; resumeLimboTipSha?: string | null; resumeLimboStepSignature?: string | null; postReviewFixCount?: number | null; recoveryRetryCount?: number | null; taskDoneRetryCount?: number | null; worktreeSessionRetryCount?: number | null; completionHandoffLimboRecoveryCount?: number | null; verificationFailureCount?: number | null; mergeConflictBounceCount?: number | null; mergeAuditBounceCount?: number | null; mergeTransientRetryCount?: number | null; branchConflictRecoveryCount?: number | null; reviewerContextRetryCount?: number | null; reviewerFallbackRetryCount?: number | null; nextRecoveryAt?: string | null; enabledWorkflowSteps?: string[]; noCommitsExpected?: boolean | null; modelProvider?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningModelId?: string | null; thinkingLevel?: string | null; error?: string | null; summary?: string | null; sessionFile?: string | null; firstExecutionAt?: string | null; cumulativeActiveMs?: number | null; executionStartedAt?: string | null; executionCompletedAt?: string | null; review?: import("./types.js").TaskReview | null; reviewState?: import("./types.js").TaskReviewState | null; workflowStepResults?: import("./types.js").WorkflowStepResult[] | null; mergeDetails?: import("./types.js").MergeDetails | null; sourceIssue?: import("./types.js").TaskSourceIssue | null; sourceMetadataPatch?: Record | null; githubTracking?: import("./types.js").TaskGithubTracking | null; tokenUsage?: import("./types.js").TaskTokenUsage | null; modifiedFiles?: string[] | null; missionId?: string | null; sliceId?: string | null }, runContext?: RunMutationContext, ): Promise { - return this.withTaskLock(id, async () => { + return this.withTaskLock(id, () => this.updateTaskUnlocked(id, updates, runContext)); + } + + /** + * The body of {@link updateTask} WITHOUT acquiring the per-task lock. Callers + * that already hold `withTaskLock(id)` — e.g. workflow-selection mutations + * that bundle a `task_workflow_selection`/`workflow_steps` write with the + * `enabledWorkflowSteps` update — invoke this directly so the whole sequence + * runs under a single lock acquisition. The per-task lock is non-reentrant, + * so calling the public `updateTask` from inside an outer `withTaskLock(id)` + * would deadlock; this variant exists to avoid that. + */ + private async updateTaskUnlocked( + id: string, + updates: Parameters[1], + runContext?: RunMutationContext, + ): Promise { + { if (updates.dependencies !== undefined) { await this.assertNoDependencyCycle( id, @@ -6663,7 +6680,7 @@ export class TaskStore extends EventEmitter { } this.emitTaskLifecycleEventSafely("task:updated", [task]); return task; - }); + } } /** @@ -11431,47 +11448,55 @@ ${stepsSection}`; * before any state is written. */ async selectTaskWorkflow(taskId: string, workflowId: string): Promise { - const def = await this.getWorkflowDefinition(workflowId); - if (!def) throw new Error(`Workflow '${workflowId}' not found`); - // Compile once up front: a non-linear graph aborts before any mutation. - const inputs = compileWorkflowToSteps(def.ir); + // Hold the task lock across the whole sequence (materialize → owner write → + // prior-step cleanup) so it can't interleave with a concurrent select/clear + // or executor updateTask on the same task. updateTaskUnlocked is used inside + // because the per-task lock is non-reentrant. + return this.withTaskLock(taskId, async () => { + const def = await this.getWorkflowDefinition(workflowId); + if (!def) throw new Error(`Workflow '${workflowId}' not found`); + // Compile once up front: a non-linear graph aborts before any mutation. + const inputs = compileWorkflowToSteps(def.ir); - // Materialize the new steps and point the task at them BEFORE deleting the - // prior selection's rows, so a mid-flight failure never leaves the task - // referencing already-deleted step ids. - const priorSelection = this.getTaskWorkflowSelection(taskId); - const ids = await this.materializeWorkflowSteps(workflowId, inputs); - try { - await this.updateTask(taskId, { enabledWorkflowSteps: ids }); - this.writeTaskWorkflowSelection(taskId, workflowId, ids); - } catch (err) { - // The owner write (updateTask / selection upsert) failed, so the steps we - // just materialized would orphan with no selection row pointing at them. - // Delete them before propagating; the prior selection is left untouched. - for (const stepId of ids) { - try { - this.db.prepare("DELETE FROM workflow_steps WHERE id = ?").run(stepId); - } catch { - // Best-effort cleanup; surface the original error below. + // Materialize the new steps and point the task at them BEFORE deleting the + // prior selection's rows, so a mid-flight failure never leaves the task + // referencing already-deleted step ids. + const priorSelection = this.getTaskWorkflowSelection(taskId); + const ids = await this.materializeWorkflowSteps(workflowId, inputs); + try { + await this.updateTaskUnlocked(taskId, { enabledWorkflowSteps: ids }); + this.writeTaskWorkflowSelection(taskId, workflowId, ids); + } catch (err) { + // The owner write (updateTask / selection upsert) failed, so the steps we + // just materialized would orphan with no selection row pointing at them. + // Delete them before propagating; the prior selection is left untouched. + for (const stepId of ids) { + try { + this.db.prepare("DELETE FROM workflow_steps WHERE id = ?").run(stepId); + } catch { + // Best-effort cleanup; surface the original error below. + } } + this.workflowStepsCache = null; + throw err; } - this.workflowStepsCache = null; - throw err; - } - if (priorSelection) { - for (const stepId of priorSelection.stepIds) { - this.db.prepare("DELETE FROM workflow_steps WHERE id = ?").run(stepId); + if (priorSelection) { + for (const stepId of priorSelection.stepIds) { + this.db.prepare("DELETE FROM workflow_steps WHERE id = ?").run(stepId); + } + this.workflowStepsCache = null; } - this.workflowStepsCache = null; - } - return ids; + return ids; + }); } /** Clear a task's workflow selection and its enabled steps. */ async clearTaskWorkflowSelection(taskId: string): Promise { - this.removeMaterializedSelection(taskId); - await this.updateTask(taskId, { enabledWorkflowSteps: [] }); + await this.withTaskLock(taskId, async () => { + this.removeMaterializedSelection(taskId); + await this.updateTaskUnlocked(taskId, { enabledWorkflowSteps: [] }); + }); } /** diff --git a/packages/dashboard/app/components/SettingsModal.tsx b/packages/dashboard/app/components/SettingsModal.tsx index f2923876c5..2fa8729550 100644 --- a/packages/dashboard/app/components/SettingsModal.tsx +++ b/packages/dashboard/app/components/SettingsModal.tsx @@ -344,6 +344,7 @@ const KNOWN_EXPERIMENTAL_FEATURES: Record = { sandbox: "Sandbox (command isolation)", chatRooms: "Chat Rooms", agentOnboarding: "Planning-style Agent Onboarding", + workflowGraphExecutor: "Workflow Graph Engine (run custom workflows)", }; const EXPERIMENTAL_FEATURE_LEGACY_ALIASES: Record = { diff --git a/packages/engine/src/__tests__/interpreter-merge-seam.test.ts b/packages/engine/src/__tests__/interpreter-merge-seam.test.ts new file mode 100644 index 0000000000..b4a42c1c1f --- /dev/null +++ b/packages/engine/src/__tests__/interpreter-merge-seam.test.ts @@ -0,0 +1,88 @@ +import { describe, expect, it, vi } from "vitest"; +import { ProjectEngine } from "../project-engine.js"; + +// These exercise two PR-1363 review fixes on the merge seam without standing up +// a full ProjectEngine: the per-task resolver LIST (multiple awaiters can't +// overwrite each other) and requestInterpreterMerge honoring auto-merge +// eligibility (a graph merge node must not override an autoMerge-off project). + +describe("interpreter merge seam", () => { + // Build a bare object whose prototype is ProjectEngine's, so the helper + // methods can call each other (this.takeMergeResolvers etc.) without + // constructing a full engine. + function bareEngine(extra: Record = {}): any { + return Object.assign(Object.create(ProjectEngine.prototype), { + manualMergeResolvers: new Map(), + ...extra, + }); + } + + describe("per-task resolver list", () => { + it("resolves every awaiter for a task, not just the last one", () => { + const engine = bareEngine(); + + const a = { resolve: vi.fn(), reject: vi.fn() }; + const b = { resolve: vi.fn(), reject: vi.fn() }; + engine.addMergeResolver("FN-1", a); + engine.addMergeResolver("FN-1", b); // would have overwritten 'a' before the fix + + const result = { merged: true } as any; + engine.resolveMergeResolvers("FN-1", result); + + expect(a.resolve).toHaveBeenCalledWith(result); + expect(b.resolve).toHaveBeenCalledWith(result); + expect(engine.manualMergeResolvers.has("FN-1")).toBe(false); + }); + + it("rejects every awaiter and clears the task entry", () => { + const engine = bareEngine(); + + const a = { resolve: vi.fn(), reject: vi.fn() }; + const b = { resolve: vi.fn(), reject: vi.fn() }; + engine.addMergeResolver("FN-2", a); + engine.addMergeResolver("FN-2", b); + + const err = new Error("boom"); + engine.rejectMergeResolvers("FN-2", err); + + expect(a.reject).toHaveBeenCalledWith(err); + expect(b.reject).toHaveBeenCalledWith(err); + expect(engine.hasMergeResolvers("FN-2")).toBe(false); + }); + }); + + describe("requestInterpreterMerge auto-merge eligibility", () => { + function fakeEngineWith(opts: { autoEligible: boolean; onMerge?: ReturnType }) { + return { + runtime: { + getTaskStore: () => ({ + getSettings: async () => ({ autoMerge: opts.autoEligible, globalPause: false, enginePaused: false }), + getTask: async () => ({ id: "FN-3", column: "in-review", branch: "feat", paused: false, mergeDetails: undefined }), + }), + }, + // The real allowInReviewMergeProcessing depends on branch-group context; + // stub it to isolate the autoMerge-off gate this test cares about. + allowInReviewMergeProcessing: () => opts.autoEligible, + onMerge: opts.onMerge ?? vi.fn(), + }; + } + + it("returns merged:false and does NOT force a merge when autoMerge is off", async () => { + const onMerge = vi.fn(); + const fakeEngine = fakeEngineWith({ autoEligible: false, onMerge }); + const result = await (ProjectEngine.prototype as any).requestInterpreterMerge.call(fakeEngine, "FN-3"); + + expect(result.merged).toBe(false); + expect(onMerge).not.toHaveBeenCalled(); // the human "merge now" bypass is never invoked + }); + + it("routes through onMerge when the task is auto-merge eligible", async () => { + const onMerge = vi.fn(async () => ({ merged: true, branch: "feat" }) as any); + const fakeEngine = fakeEngineWith({ autoEligible: true, onMerge }); + const result = await (ProjectEngine.prototype as any).requestInterpreterMerge.call(fakeEngine, "FN-3"); + + expect(onMerge).toHaveBeenCalledWith("FN-3"); + expect(result.merged).toBe(true); + }); + }); +}); diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index 2219fdebe6..76ab6fac6a 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -236,6 +236,8 @@ export interface ProjectEngineOptions { * via ProjectManager) gets the full subsystem set, eliminating the class of * bugs where a subsystem is forgotten in one code path. */ +type MergeResolver = { resolve: (result: MergeResult) => void; reject: (err: Error) => void }; + export class ProjectEngine { private runtime: InProcessRuntime; private started = false; @@ -280,12 +282,40 @@ export class ProjectEngine { * When `onMerge` is called, the task is enqueued like auto-merge but a * Promise is stored here so the caller can await the result. */ - private manualMergeResolvers = new Map< - string, - { resolve: (result: MergeResult) => void; reject: (err: Error) => void } - >(); + // Per-task LIST of waiters, not a single resolver: both the dashboard "merge + // now" path and the workflow interpreter's merge seam call onMerge, so a task + // can have more than one caller awaiting the same merge. A single-entry map + // would let a second caller overwrite the first, stranding its promise. + private manualMergeResolvers = new Map>(); private shuttingDown = false; + private addMergeResolver(taskId: string, r: MergeResolver): void { + const list = this.manualMergeResolvers.get(taskId); + if (list) list.push(r); + else this.manualMergeResolvers.set(taskId, [r]); + } + + /** Remove and return all waiters for a task (empty array if none). */ + private takeMergeResolvers(taskId: string): MergeResolver[] { + const list = this.manualMergeResolvers.get(taskId); + this.manualMergeResolvers.delete(taskId); + return list ?? []; + } + + private hasMergeResolvers(taskId: string): boolean { + return (this.manualMergeResolvers.get(taskId)?.length ?? 0) > 0; + } + + /** Resolve every waiter for a task with the same result, then clear them. */ + private resolveMergeResolvers(taskId: string, result: MergeResult): void { + for (const r of this.takeMergeResolvers(taskId)) r.resolve(result); + } + + /** Reject every waiter for a task with the same error, then clear them. */ + private rejectMergeResolvers(taskId: string, err: Error): void { + for (const r of this.takeMergeResolvers(taskId)) r.reject(err); + } + private static readonly MAX_AUTO_MERGE_RETRIES = 3; /** FN-5697/FN-5674: cap transient provider/network abort retries in auto-merge. * Examples: "This operation was aborted", "socket hang up", `server_error`. @@ -353,9 +383,10 @@ export class ProjectEngine { this.runtime.setMergeActiveClearer?.((taskId) => { this.mergeActive.delete(taskId); }); - // Workflow-graph interpreter merge seam: resolves with the merge outcome - // through the same serialized merge queue the legacy pipeline uses. - this.runtime.setMergeRequester?.((taskId) => this.onMerge(taskId)); + // Workflow-graph interpreter merge seam: routes through the auto-merge + // eligibility gate (requestInterpreterMerge), NOT the human "merge now" + // bypass, so a graph merge node can't override an autoMerge-off project. + this.runtime.setMergeRequester?.((taskId) => this.requestInterpreterMerge(taskId)); } getActiveMergeTaskId(): string | null { @@ -637,9 +668,11 @@ export class ProjectEngine { this.activeMergeSession = null; } - // Reject any pending manual merge promises - for (const [taskId, resolver] of this.manualMergeResolvers) { - resolver.reject(new Error(`Engine shutting down — merge for ${taskId} aborted`)); + // Reject any pending manual merge promises (every waiter per task) + for (const [taskId, resolvers] of this.manualMergeResolvers) { + for (const resolver of resolvers) { + resolver.reject(new Error(`Engine shutting down — merge for ${taskId} aborted`)); + } } this.manualMergeResolvers.clear(); @@ -939,20 +972,57 @@ export class ProjectEngine { // existing merge to finish rather than starting a second one. if (this.mergeActive.has(taskId)) { return new Promise((resolve, reject) => { - this.manualMergeResolvers.set(taskId, { resolve, reject }); + this.addMergeResolver(taskId, { resolve, reject }); // Don't re-enqueue — the task is already in the queue/active }); } return new Promise((resolve, reject) => { - this.manualMergeResolvers.set(taskId, { resolve, reject }); + this.addMergeResolver(taskId, { resolve, reject }); if (!this.internalEnqueueMerge(taskId)) { - this.manualMergeResolvers.delete(taskId); - reject(new Error(`Merge enqueue rejected for ${taskId}`)); + // Drop just-added waiter(s) for this task and fail them. + this.rejectMergeResolvers(taskId, new Error(`Merge enqueue rejected for ${taskId}`)); } }); } + /** + * Merge entry point for the workflow graph interpreter's `merge` seam. Unlike + * onMerge (the human "merge now" bypass), this honors the project's auto-merge + * eligibility: when autoMerge is off (or the task isn't merge-eligible), it + * does NOT force the merge. It resolves with `merged: false` so the seam treats + * it as "manual merge required" and parks the task in review — preserving the + * contract that autoMerge-off leaves in-review terminal until a human merges. + */ + async requestInterpreterMerge(taskId: string): Promise { + let task: Task | null = null; + let settings: Settings | undefined; + try { + const store = this.runtime.getTaskStore(); + settings = await store.getSettings(); + task = await store.getTask(taskId); + } catch { + // Fall through to the not-eligible response below. + } + const eligible = !!task && !!settings + && task.column === "in-review" + && !settings.globalPause && !settings.enginePaused + && this.allowInReviewMergeProcessing(task, settings) + && !(task.paused && !task.mergeDetails?.mergeConfirmed); + if (!eligible) { + runtimeLog.log(`Interpreter merge for ${taskId} not auto-eligible (autoMerge off / not ready) — manual merge required`); + return { + task: task as Task, + branch: task?.branch ?? "", + merged: false, + worktreeRemoved: false, + branchDeleted: false, + } as MergeResult; + } + // Eligible: route through the normal serialized merge path. + return this.onMerge(taskId); + } + private setRestoreDiagnostics( outcome: TunnelRestoreDiagnostics["outcome"], reason: TunnelRestoreReasonCode, @@ -1490,10 +1560,10 @@ export class ProjectEngine { // pickNextMergeTaskId awaits store.getTask; re-check shutdown so we // don't start a merge whose queue entry was cleared by stop(). if (this.shuttingDown) break; - const manualResolver = this.manualMergeResolvers.get(taskId); + const hasManualResolver = this.hasMergeResolvers(taskId); try { // Manual merges (onMerge) skip auto-merge eligibility checks - if (!manualResolver) { + if (!hasManualResolver) { // Re-check autoMerge and pause before each merge const settings = await store.getSettings(); if (settings.globalPause || settings.enginePaused) { @@ -1820,20 +1890,16 @@ export class ProjectEngine { 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); - } + // Temporarily stash the waiters so the finally block doesn't + // prematurely resolve them. The re-enqueue restores them. + const stashedResolvers = this.takeMergeResolvers(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")); + for (const r of stashedResolvers) r.reject(new Error("Engine shutting down")); return; } - if (manualResolver) { - this.manualMergeResolvers.set(taskId, manualResolver); - } + for (const r of stashedResolvers) this.addMergeResolver(taskId, r); this.internalEnqueueMerge(taskId); }, retryMs); continue; @@ -1876,7 +1942,7 @@ export class ProjectEngine { if (mergeStrategy === "pull-request" && this.options.processPullRequestMerge) { this.activeMergeTaskId = taskId; - runtimeLog.log(`${manualResolver ? "Manual" : "Auto"}-merge processing PR flow for ${taskId}...`); + runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge processing PR flow for ${taskId}...`); const result = await this.options.processPullRequestMerge( store, cwd, @@ -1884,7 +1950,7 @@ export class ProjectEngine { (this.runtime as any).worktreePool, ); if (result === "merged") { - runtimeLog.log(`${manualResolver ? "Manual" : "Auto"}-merge PR merged: ${taskId}`); + runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR merged: ${taskId}`); const mergedTask = await store.getTask(taskId).catch(() => null); if (mergedTask) { store.emit("task:merged", { @@ -1900,14 +1966,13 @@ export class ProjectEngine { } await attemptBranchGroupPromotion(mergedTask); } else if (result === "waiting") { - runtimeLog.log(`${manualResolver ? "Manual" : "Auto"}-merge PR waiting: ${taskId}`); + runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR waiting: ${taskId}`); } - if (manualResolver) { + if (hasManualResolver) { // PR merge path doesn't produce a full MergeResult — fetch the task // and construct one so the dashboard endpoint can respond. const prTask = await store.getTask(taskId).catch(() => null); - this.manualMergeResolvers.delete(taskId); - manualResolver.resolve({ + this.resolveMergeResolvers(taskId, { task: prTask!, branch: prTask?.branch ?? "", merged: result === "merged", @@ -1917,7 +1982,7 @@ export class ProjectEngine { } } else { // Direct merge via AI agent, gated by semaphore - runtimeLog.log(`${manualResolver ? "Manual" : "Auto"}-merge merging ${taskId}...`); + runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge merging ${taskId}...`); const semaphore = (this.runtime as any).globalSemaphore; @@ -1931,7 +1996,7 @@ export class ProjectEngine { this.activeMergeTaskId = taskId; this.mergeAbortController = new AbortController(); const mergerOptions = { - manual: !!manualResolver, + manual: hasManualResolver, pool, usageLimitPauser, agentStore, @@ -1962,11 +2027,10 @@ export class ProjectEngine { } this.activeMergeSession = null; - runtimeLog.log(`${manualResolver ? "Manual" : "Auto"}-merge merged: ${taskId}`); + runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge merged: ${taskId}`); - if (manualResolver) { - this.manualMergeResolvers.delete(taskId); - manualResolver.resolve(result); + if (hasManualResolver) { + this.resolveMergeResolvers(taskId, result); } // Reset retries on success @@ -1983,25 +2047,24 @@ export class ProjectEngine { const mergeWasAborted = err instanceof Error && err.name === "MergeAbortedError"; if (mergeWasAborted) { - runtimeLog.log(`${manualResolver ? "Manual" : "Auto"}-merge aborted for ${taskId}: ${errorMsg}`); + runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge aborted for ${taskId}: ${errorMsg}`); this.mergeAbortController = null; - if (manualResolver) { - this.manualMergeResolvers.delete(taskId); - manualResolver.reject(err instanceof Error ? err : new Error(errorMsg)); + if (hasManualResolver) { + this.rejectMergeResolvers(taskId, err instanceof Error ? err : new Error(errorMsg)); } else { await store.updateTask(taskId, { status: null }).catch(() => undefined); } continue; } - runtimeLog.error(`${manualResolver ? "Manual" : "Auto"}-merge failed for ${taskId}: ${errorMsg}`); + runtimeLog.error(`${hasManualResolver ? "Manual" : "Auto"}-merge failed for ${taskId}: ${errorMsg}`); // Surface every merge failure on the task log so the dashboard shows // *why* a merge didn't complete instead of silently looping. await store .logEntry( taskId, - `${manualResolver ? "Manual" : "Auto"}-merge failed: ${errorMsg}`, + `${hasManualResolver ? "Manual" : "Auto"}-merge failed: ${errorMsg}`, err instanceof Error ? err.name : undefined, ) .catch((logErr: unknown) => { @@ -2011,9 +2074,8 @@ export class ProjectEngine { }); // If this was a manual merge, reject the promise and skip auto-retry logic - if (manualResolver) { - this.manualMergeResolvers.delete(taskId); - manualResolver.reject(err instanceof Error ? err : new Error(errorMsg)); + if (hasManualResolver) { + this.rejectMergeResolvers(taskId, err instanceof Error ? err : new Error(errorMsg)); continue; } @@ -2544,12 +2606,10 @@ export class ProjectEngine { this.mergeAbortController = null; this.mergeActive.delete(taskId); // If a manual merge was requested while this task was already in-flight, - // the resolver was set but not consumed above. Resolve it now. - const lateResolver = this.manualMergeResolvers.get(taskId); - if (lateResolver) { - this.manualMergeResolvers.delete(taskId); + // the waiter(s) were set but not consumed above. Resolve them now. + if (this.hasMergeResolvers(taskId)) { const finalTask = await store.getTask(taskId).catch(() => null); - lateResolver.resolve({ + this.resolveMergeResolvers(taskId, { task: finalTask!, branch: finalTask?.branch ?? "", merged: finalTask?.column === "done",