From dcf1b921a3c4fec5dd8accdffedcd8048f258d3e Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Fri, 31 Jul 2026 22:10:14 -0700 Subject: [PATCH] fix(engine): close active-slot reservation handoff gaps --- .../engine/src/__tests__/concurrency.test.ts | 11 ++ packages/engine/src/project-engine.ts | 14 ++- packages/engine/src/scheduler.ts | 113 ++++++++++-------- 3 files changed, 83 insertions(+), 55 deletions(-) diff --git a/packages/engine/src/__tests__/concurrency.test.ts b/packages/engine/src/__tests__/concurrency.test.ts index fc2acd8705..f8871dd096 100644 --- a/packages/engine/src/__tests__/concurrency.test.ts +++ b/packages/engine/src/__tests__/concurrency.test.ts @@ -1106,6 +1106,17 @@ describe("ProjectAdmissionCoordinator", () => { })).toBe(false); expect(started).toEqual(["FN-PLANNING"]); + // Once the selected task is durably live, its matching reservation is the + // same slot—not a second occupant—so the next real slot remains usable. + expect(await coordinator.reserveIfAvailable({ + projectId: "project-a", + taskId: "FN-DIRECT-SCHEDULER", + maxConcurrent: 10, + claimed: () => 9, + claimedTaskIds: () => ["FN-PLANNING"], + })).toBe(true); + + coordinator.releaseReservation("FN-DIRECT-SCHEDULER"); coordinator.releaseReservation("FN-PLANNING"); }); diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index 030267dc87..82e62d7aa3 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -76,7 +76,7 @@ import { sweepStaleAutostashes, VerificationError } from "./merger.js"; import { runAiMerge, landWorkspaceTask, WorkspacePartialLandError, WorkspaceRepoLandBusyError } from "./merger-ai.js"; import { promoteBranchGroup, type BranchGroupPromotionResult, type CreateGroupPrFn, type SyncGroupPrFn } from "./group-merge-coordinator.js"; import { - computeTopLevelConcurrencyClaimedFromStore, + persistedTopLevelAgentTaskIdsFromStore, projectAdmissionCoordinator, resolveActiveTaskCapacityLimit, } from "./concurrency.js"; @@ -3893,6 +3893,12 @@ export class ProjectEngine { } let selected = false; const admissionSettings = await store.getSettings(); + let mergeClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined; + const getMergeClaimSnapshot = () => mergeClaimSnapshot ??= (async () => { + const tasks = await store.listTasks({ slim: true, includeArchived: false }); + const ids = await persistedTopLevelAgentTaskIdsFromStore(store, tasks); + return { count: ids.length, ids }; + })(); /* FNXC:ConcurrencyAdmission 2026-08-01-01:50 (ROOT CAUSE — triage admission died during every merge): This lane previously ran `value = await start()` INSIDE its admission `start()` callback — @@ -3918,10 +3924,8 @@ export class ProjectEngine { maxWorktrees: admissionSettings.maxWorktrees ?? 4, worktreeLimitEnabled: admissionSettings.worktreeLimitEnabled, }), - claimed: async () => computeTopLevelConcurrencyClaimedFromStore({ - store, - tasks: await store.listTasks({ slim: true, includeArchived: false }), - }), + claimed: async () => (await getMergeClaimSnapshot()).count, + claimedTaskIds: async () => (await getMergeClaimSnapshot()).ids, refresh: async () => [{ taskId, projectId: cwd, diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts index 9ec35dcf70..90df991f12 100644 --- a/packages/engine/src/scheduler.ts +++ b/packages/engine/src/scheduler.ts @@ -2888,6 +2888,7 @@ export class Scheduler { taskId: task.id, maxConcurrent: activeTaskLimit, claimed: () => activeWorktreeTaskIds.length, + claimedTaskIds: () => activeWorktreeTaskIds, }); if (!projectSlotReserved) { if (reservedScope) { @@ -2930,65 +2931,77 @@ export class Scheduler { } let acquiredSymbols: string[] | undefined; - if (missionAdmission.kind === "symbol-lock") { + try { + if (missionAdmission.kind === "symbol-lock") { /* FNXC:MissionSymbolAdmission 2026-07-31-12:00: Acquire after all capacity gates and immediately before hold release; the reservation release path below returns this lock if moveTask rejects, while a successful move transfers ownership to the task. */ - const lockResult = await this.store.acquireSymbolLocks( - missionAdmission.symbols, - { ownerTaskId: task.id, missionId: freshTask.missionId, featureId: missionAdmission.feature.id, agentId: "scheduler" }, - SYMBOL_LOCK_LEASE_MS, - ); - if (!lockResult.acquired) { - if (dropPreHeldExecutorSlot(task.id)) sem?.release(); - if (reservedScope) { - activeScopes.delete(task.id); - activeScopeColumns.delete(task.id); - } - const conflict = lockResult.conflicts[0]; - await this.store.updateTask(task.id, { status: "queued", blockedBy: null, overlapBlockedBy: null }); - await this.logDispatchQueuedReason( - task.id, - `queued — symbol contention: symbol=${conflict?.symbolKey ?? "unknown"} holder=${conflict?.ownerTaskId ?? "unknown"}`, + const lockResult = await this.store.acquireSymbolLocks( + missionAdmission.symbols, + { ownerTaskId: task.id, missionId: freshTask.missionId, featureId: missionAdmission.feature.id, agentId: "scheduler" }, + SYMBOL_LOCK_LEASE_MS, ); - return null; + if (!lockResult.acquired) { + if (dropPreHeldExecutorSlot(task.id)) sem?.release(); + if (reservedScope) { + activeScopes.delete(task.id); + activeScopeColumns.delete(task.id); + } + const conflict = lockResult.conflicts[0]; + await this.store.updateTask(task.id, { status: "queued", blockedBy: null, overlapBlockedBy: null }); + await this.logDispatchQueuedReason( + task.id, + `queued — symbol contention: symbol=${conflict?.symbolKey ?? "unknown"} holder=${conflict?.ownerTaskId ?? "unknown"}`, + ); + return null; + } + acquiredSymbols = missionAdmission.symbols; + await this.store.logEntry(task.id, `symbol-lock admission acquired: ${missionAdmission.symbols.join(", ")}`); } - acquiredSymbols = missionAdmission.symbols; - await this.store.logEntry(task.id, `symbol-lock admission acquired: ${missionAdmission.symbols.join(", ")}`); + + dispatchPrepByTaskId.set(task.id, { + baseBranch: this.resolveBaseBranch(freshTask, tasks, isReviewColumnTask), + dispatchStormCount: nextDispatchStormCount, + dispatchTimestamp, + effectiveNodeId: effectiveNode.nodeId ?? null, + effectiveNodeSource: effectiveNode.source, + task: freshTask, + }); + + reservedWorktreeSlots += 1; + reservedConcurrentSlots += 1; + let released = false; + return { + release: () => { + if (released) return; + released = true; + if (reservedScope) { + activeScopes.delete(task.id); + activeScopeColumns.delete(task.id); + } + reservedWorktreeSlots = releaseReservedSlot(reservedWorktreeSlots); + reservedConcurrentSlots = releaseReservedSlot(reservedConcurrentSlots); + dispatchPrepByTaskId.delete(task.id); + if (dropPreHeldExecutorSlot(task.id)) sem?.release(); + if (acquiredSymbols) { + void this.store.releaseSymbolLocks(acquiredSymbols, task.id); + } + }, + }; + } catch (error) { + if (dropPreHeldExecutorSlot(task.id)) sem?.release(); + if (reservedScope) { + activeScopes.delete(task.id); + activeScopeColumns.delete(task.id); + } + if (acquiredSymbols) { + await this.store.releaseSymbolLocks(acquiredSymbols, task.id).catch(() => undefined); + } + throw error; } - - dispatchPrepByTaskId.set(task.id, { - baseBranch: this.resolveBaseBranch(freshTask, tasks, isReviewColumnTask), - dispatchStormCount: nextDispatchStormCount, - dispatchTimestamp, - effectiveNodeId: effectiveNode.nodeId ?? null, - effectiveNodeSource: effectiveNode.source, - task: freshTask, - }); - - reservedWorktreeSlots += 1; - reservedConcurrentSlots += 1; - let released = false; - return { - release: () => { - if (released) return; - released = true; - if (reservedScope) { - activeScopes.delete(task.id); - activeScopeColumns.delete(task.id); - } - reservedWorktreeSlots = releaseReservedSlot(reservedWorktreeSlots); - reservedConcurrentSlots = releaseReservedSlot(reservedConcurrentSlots); - dispatchPrepByTaskId.delete(task.id); - if (dropPreHeldExecutorSlot(task.id)) sem?.release(); - if (acquiredSymbols) { - void this.store.releaseSymbolLocks(acquiredSymbols, task.id); - } - }, - }; }, allocateWorktree: (task, reservedNames) => this.planWorktreePath(task, settings.worktreeNaming, reservedNames, settings),