fix(engine): close active-slot reservation handoff gaps

This commit is contained in:
gsxdsm
2026-07-31 22:10:14 -07:00
parent 012729cf2b
commit dcf1b921a3
3 changed files with 83 additions and 55 deletions

View File

@@ -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");
});

View File

@@ -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,

View File

@@ -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),