fix: count only actively running tasks against worktree capacity

Retained directories on queued, paused, blocked, or terminal tasks no longer
consume scheduler slots. Agent concurrency and worktree capacity now count the
same canonical live-task population through one project admission ceiling
(resolveActiveTaskCapacityLimit) with an atomic reserveIfAvailable claim, so
planning, execute, and merge lanes cannot each observe and claim the final
worktree slot independently.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
gsxdsm
2026-07-31 22:07:05 -07:00
parent 031d0f3e0b
commit 5a19d1da6e
11 changed files with 451 additions and 488 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Count only actively running tasks against worktree capacity.
category: fix
dev: Retained directories on queued, paused, blocked, or terminal tasks no longer consume scheduler slots.

View File

@@ -15,6 +15,7 @@ import {
persistedTopLevelAgentSlots,
recoverIdleSemaphoreLeakCandidate,
registerPreHeldExecutorSlot,
resolveActiveTaskCapacityLimit,
takePreHeldExecutorSlot,
} from "../concurrency.js";
@@ -1067,6 +1068,47 @@ describe("AgentSemaphore resilience (FN-978)", () => {
describe("ProjectAdmissionCoordinator", () => {
it("shares the final active-task slot across planning, execution, and merge lanes", async () => {
const coordinator = new ProjectAdmissionCoordinator();
const started: string[] = [];
const activeTaskLimit = resolveActiveTaskCapacityLimit({
maxConcurrent: 12,
maxWorktrees: 9,
worktreeLimitEnabled: true,
});
for (const [lane, taskId, createdAt] of [
["planning", "FN-PLANNING", "2026-01-01T00:00:00.000Z"],
["execute", "FN-EXECUTE", "2026-01-02T00:00:00.000Z"],
["merge", "FN-MERGE", "2026-01-03T00:00:00.000Z"],
] as const) {
coordinator.registerProvider(lane, {
projectId: "project-a",
refresh: async () => [{
taskId,
projectId: "project-a",
createdAt,
start: async () => { started.push(taskId); },
}],
});
}
expect(await coordinator.admitOldest({
projectId: "project-a",
maxConcurrent: activeTaskLimit,
claimed: () => 8,
})).toBe("FN-PLANNING");
expect(await coordinator.reserveIfAvailable({
projectId: "project-a",
taskId: "FN-DIRECT-SCHEDULER",
maxConcurrent: activeTaskLimit,
claimed: () => 8,
})).toBe(false);
expect(started).toEqual(["FN-PLANNING"]);
coordinator.releaseReservation("FN-PLANNING");
});
it("admits the oldest same-project candidate atomically and partitions projects", async () => {
const coordinator = new ProjectAdmissionCoordinator();
const started: string[] = [];

View File

@@ -1,30 +1,14 @@
// @vitest-environment node
/*
FNXC:WorkflowResolvedColumns 2026-08-01-09:20 (fleet):
THE INVARIANT: the worktree-capacity read excludes terminal lanes by ROLE, on any board.
The capacity gate counts every non-terminal task holding a worktree. Terminal cards are excluded
because their retained worktrees are cleanup-owned, not capacity. That exclusion arrived hand-rolled,
with `flags.complete === true || flags.archived === true` and a `"done" | "archived"` literal fallback
— the same shape the FNXC on `isWipColumnTask` records as already removed once from this file.
Getting it wrong is not symmetric. Under-counting admits work over the cap: the commit that added the
gate reports maxWorktrees=4 with four planning sessions and a fifth worktree admitted. Over-counting
merely starves dispatch. So the risk of a renamed board silently failing the exclusion is the
dangerous direction.
WHY THIS TESTS THE PREDICATE AND NOT THE DISPATCH PATH: same reason as
`scheduler-load-lane-union.test.ts` — the call site sits inside `schedule()`, which a unit test has no
business standing up. The predicate is the whole of the decision.
NOTE ON COVERAGE, recorded rather than implied: blinding this predicate to `false` leaves all 22
scheduler suites green (365 tests). The capacity logic it feeds has no behavioural coverage at all.
This pins the lane vocabulary; it does NOT pin the capacity arithmetic, which is still unguarded.
FNXC:WorktreeCapacity 2026-08-01-04:38:
The worktree-capacity ledger uses the canonical enriched live-task predicate. Renamed workflow roles
must classify active WIP/planning tasks correctly while excluding inactive retained directories and
terminal cards. This is the role-resolution half of the real scheduler regression test.
*/
import { describe, expect, it } from "vitest";
import { isTerminalColumnRole, resolveColumnFlags } from "@fusion/core";
import type { WorkflowIr } from "@fusion/core";
import { enrichRunningAgentTaskShape, isRunningAgentTask } from "@fusion/core";
import type { Task, WorkflowIr } from "@fusion/core";
const RENAMED_IR = {
version: "v2", id: "wf-renamed", name: "renamed", nodes: [], edges: [],
@@ -36,29 +20,36 @@ const RENAMED_IR = {
],
} as unknown as WorkflowIr;
const flagsFor = (columnId: string) => resolveColumnFlags(
RENAMED_IR.columns.find((c) => c.id === columnId) as never,
);
const task = (id: string, column: string, overrides: Partial<Task> = {}): Task => ({
id,
column,
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
...overrides,
} as Task);
describe("worktree capacity excludes terminal lanes by role", () => {
it("excludes BOTH renamed terminal lanes, not just the complete one", () => {
expect(isTerminalColumnRole(flagsFor("shipped"), "shipped")).toBe(true);
expect(isTerminalColumnRole(flagsFor("attic"), "attic")).toBe(true);
const isLive = (candidate: Task): boolean =>
isRunningAgentTask(enrichRunningAgentTaskShape(candidate, RENAMED_IR));
describe("worktree capacity follows enriched live-task roles", () => {
it("counts active execution and planning on renamed lanes", () => {
expect(isLive(task("FN-WIP", "building"))).toBe(true);
expect(isLive(task("FN-PLAN", "backlog", { status: "planning" }))).toBe(true);
});
it("counts renamed working and hold lanes toward capacity", () => {
expect(isTerminalColumnRole(flagsFor("building"), "building")).toBe(false);
expect(isTerminalColumnRole(flagsFor("backlog"), "backlog")).toBe(false);
it("ignores an inactive retained directory", () => {
expect(isLive(task("FN-QUEUED", "backlog", {
status: "queued",
worktree: "/tmp/project/.worktrees/queued",
}))).toBe(false);
});
/*
The degraded arm is why the shared helper is used rather than the hand-rolled version: an
unresolvable column must keep answering for the legacy ids, or a board mid-migration starts counting
its own done cards against the worktree cap.
*/
it("still recognises the legacy terminal ids when flags cannot be resolved", () => {
expect(isTerminalColumnRole(undefined, "done")).toBe(true);
expect(isTerminalColumnRole(undefined, "archived")).toBe(true);
expect(isTerminalColumnRole(undefined, "in-progress")).toBe(false);
it("ignores complete and archived cards even with stale live-looking metadata", () => {
expect(isLive(task("FN-DONE", "shipped", { status: "planning" }))).toBe(false);
expect(isLive(task("FN-ARCHIVED", "attic", { sessionFile: "/tmp/stale-session" }))).toBe(false);
});
});

View File

@@ -460,7 +460,7 @@ describe("Scheduler workflow cutover", () => {
expect(ready.column).toBe("todo");
});
it("releases a retained-worktree task while already over maxWorktrees", async () => {
it("does not let a retained directory bypass a full active-task worktree cap", async () => {
const active = task({ id: "FN-101", column: "in-progress", worktree: "/tmp/project/.worktrees/fn-101" });
const ready = task({ id: "FN-200", status: "queued", worktree: "/tmp/project/.worktrees/fn-200" });
const store = storeWith([active, ready], { maxConcurrent: 4, maxWorktrees: 1 });
@@ -470,8 +470,50 @@ describe("Scheduler workflow cutover", () => {
await scheduler.schedule();
expect(store.moveTaskIf).toHaveBeenCalledWith("FN-200", "in-progress", expect.any(Function), expect.anything());
expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-200", column: "in-progress" }));
expect(store.moveTaskIf).not.toHaveBeenCalledWith("FN-200", "in-progress", expect.anything(), expect.anything());
expect(onSchedule).not.toHaveBeenCalled();
});
it("counts only active tasks when retained queued worktrees would strand execution slots", async () => {
/*
FNXC:WorktreeCapacity 2026-08-01-04:38:
Live regression: seven active tasks plus two dependency-blocked queued cards with retained
worktrees filled a nine-slot ledger. The inactive holders then blocked both dependency-free
roots from starting. Slots represent live task execution, not directories retained on disk.
*/
const planners = Array.from({ length: 6 }, (_, index) => task({
id: `FN-PLAN-${index}`,
status: "planning",
worktree: `/tmp/project/.worktrees/plan-${index}`,
}));
const executing = task({
id: "FN-EXECUTING",
column: "in-progress",
worktree: "/tmp/project/.worktrees/executing",
});
const parkedDependents = [
task({ id: "FN-PARKED-1", status: "queued", worktree: "/tmp/project/.worktrees/parked-1", dependencies: ["FN-ROOT-1"] }),
task({ id: "FN-PARKED-2", status: "queued", worktree: "/tmp/project/.worktrees/parked-2", dependencies: ["FN-ROOT-1"] }),
];
const roots = [
task({ id: "FN-ROOT-1", status: "queued" }),
task({ id: "FN-ROOT-2", status: "queued" }),
task({ id: "FN-ROOT-3", status: "queued" }),
];
const store = storeWith([...planners, executing, ...parkedDependents, ...roots], {
maxConcurrent: 12,
maxWorktrees: 9,
});
const onSchedule = vi.fn();
const scheduler = new Scheduler(store, { onSchedule });
(scheduler as unknown as { running: boolean }).running = true;
await scheduler.schedule();
expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-ROOT-1", column: "in-progress" }));
expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-ROOT-2", column: "in-progress" }));
expect(onSchedule).not.toHaveBeenCalledWith(expect.objectContaining({ id: "FN-ROOT-3" }));
expect(roots[2]?.column).toBe("todo");
});
/*
@@ -614,7 +656,7 @@ describe("Scheduler workflow cutover", () => {
expect(ready.status).toBe("queued");
});
it("does not invent capacity when a retained-worktree transfer rejects", async () => {
it("returns a rejected active-task reservation so the next candidate can start", async () => {
const retained = task({ id: "FN-001", status: "queued", worktree: "/tmp/project/.worktrees/fn-001" });
const fresh = task({ id: "FN-002", status: "queued" });
const store = storeWith([retained, fresh], { maxConcurrent: 4, maxWorktrees: 1 });
@@ -635,11 +677,11 @@ describe("Scheduler workflow cutover", () => {
await scheduler.schedule();
expect(store.moveTaskIf).toHaveBeenCalledTimes(1);
expect(store.moveTaskIf).toHaveBeenCalledTimes(2);
expect(store.moveTaskIf).toHaveBeenCalledWith("FN-001", "in-progress", expect.any(Function), expect.anything());
expect(store.moveTaskIf).not.toHaveBeenCalledWith("FN-002", "in-progress", expect.anything(), expect.anything());
expect(onSchedule).not.toHaveBeenCalled();
expect(fresh.column).toBe("todo");
expect(store.moveTaskIf).toHaveBeenCalledWith("FN-002", "in-progress", expect.any(Function), expect.anything());
expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-002", column: "in-progress" }));
expect(fresh.column).toBe("in-progress");
});
it("does not release work when the shared semaphore is saturated", async () => {

View File

@@ -1,159 +0,0 @@
// @vitest-environment node
/*
FNXC:WorktreeCapacity 2026-08-01-00:50:
THE WORKTREE-CAPACITY ARITHMETIC — the numbers behind two live defects, previously unpinned.
#3262 pinned the terminal PREDICATE and said plainly what it did not cover: blinding that predicate
to `false` leaves all 22 scheduler suites green, because the arithmetic it feeds had no behavioural
coverage at all. This is that gap.
Both observed failures live in these few lines, and they fail in OPPOSITE directions:
UNDER-COUNT admits work over the cap. The commit that added this gate reports maxWorktrees=4 with
four planning sessions each holding a worktree, and a replan dispatch admitted as the FIFTH,
because the ledger counted WIP cards only and never learned to count planners.
OVER-COUNT self-deadlocks. A planned Ready card RETAINS its planning worktree for execution reuse,
so counting it as a holder blocks its own release — 2 wip + 3 idle-held = 5/4, and the first
unpause released only 2 of 4 slots' worth of work.
Asymmetry is why both are pinned: under-counting breaks the cap and lets real work over it;
over-counting only starves dispatch. A test that covered the "safe" direction alone would leave the
expensive one open.
WHY THE PREDICATE IS INJECTED rather than resolved here: the holder set's job is the SET ARITHMETIC —
who is excluded and how the total is formed. Which lanes are terminal is #3262's test, and resolving
it here would make this file fail for that reason instead of this one.
*/
import { describe, expect, it } from "vitest";
import type { Task } from "@fusion/core";
import { nonWipWorktreeHolderIdsOf, releaseReservedSlot, releaseWorktreeReservation, reserveWorktreeOnDispatch, resolveCandidateWorktreeCapacityLimit } from "../scheduler.js";
const task = (id: string, overrides: Partial<Task> = {}): Task => ({
id,
title: id,
description: "",
column: "todo",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
...overrides,
} as Task);
const neverTerminal = () => false;
const terminalIds = (ids: string[]) => (t: Task) => ids.includes(t.id);
describe("worktree-capacity holder set", () => {
it("counts a non-WIP card that HOLDS a worktree — the planner the ledger used to miss", () => {
/*
The under-count defect, minimally: four planning sessions holding worktrees, none of them WIP.
Before the gate learned to count them the reserved total was 0 and a fifth dispatch was admitted.
*/
const planners = ["P1", "P2", "P3", "P4"].map((id) => task(id, { worktree: `/wt/${id}` }));
const holders = nonWipWorktreeHolderIdsOf(planners, [], neverTerminal);
expect(holders).toEqual(["P1", "P2", "P3", "P4"]);
expect(holders.length).toBe(4);
});
it("does not double-count a WIP card that also holds a worktree", () => {
/* WIP membership already reserves; adding the same card as a holder would inflate the total. */
const tasks = [task("W1", { worktree: "/wt/W1" }), task("I1", { worktree: "/wt/I1" })];
expect(nonWipWorktreeHolderIdsOf(tasks, ["W1"], neverTerminal)).toEqual(["I1"]);
});
it("excludes a TERMINAL card's retained worktree — cleanup-owned, not capacity", () => {
const tasks = [task("DONE1", { worktree: "/wt/DONE1" }), task("LIVE1", { worktree: "/wt/LIVE1" })];
expect(nonWipWorktreeHolderIdsOf(tasks, [], terminalIds(["DONE1"]))).toEqual(["LIVE1"]);
});
it("ignores a card with no worktree — a WIP card without one still reserves via WIP, not here", () => {
const tasks = [task("N1"), task("N2", { worktree: "" }), task("H1", { worktree: "/wt/H1" })];
expect(nonWipWorktreeHolderIdsOf(tasks, [], neverTerminal)).toEqual(["H1"]);
});
it("reproduces the observed 5-of-4 total that the self-deadlock fix addresses", () => {
/*
2 wip + 3 idle-held = 5 against maxWorktrees=4. The total itself is CORRECT — every one of those
five is a real worktree. The bug was gating a candidate against a total that included the
candidate's OWN retained worktree, which the next case covers.
*/
const wipIds = ["W1", "W2"];
const tasks = [
task("W1", { worktree: "/wt/W1" }),
task("W2", { worktree: "/wt/W2" }),
task("R1", { worktree: "/wt/R1" }),
task("R2", { worktree: "/wt/R2" }),
task("R3", { worktree: "/wt/R3" }),
];
const holders = nonWipWorktreeHolderIdsOf(tasks, wipIds, neverTerminal);
const reservedWorktreeSlots = wipIds.length + holders.length;
expect(holders).toEqual(["R1", "R2", "R3"]);
expect(reservedWorktreeSlots).toBe(5);
});
});
describe("resolveCandidateWorktreeCapacityLimit", () => {
it("does not gate a candidate that reuses its retained worktree, even while the ledger is over cap", () => {
/*
FNXC:WorktreeCapacity 2026-08-01-04:07:
Live regression: 12 non-terminal cards retained worktrees against maxWorktrees=9. A queued
candidate already owned one of those trees, so dispatch would transfer the existing slot and
leave the ledger at 12. Subtracting only the candidate (12 - 1 = 11) still wedged every queued
card forever. The worktree dimension must be absent for a zero-allocation transfer; the separate
maxConcurrent gate continues to arbitrate whether another agent may run.
*/
expect(resolveCandidateWorktreeCapacityLimit(9, true)).toBeNull();
});
it("keeps the configured gate for a candidate that must allocate a worktree", () => {
expect(resolveCandidateWorktreeCapacityLimit(9, false)).toBe(9);
});
it("preserves a disabled worktree gate for every candidate", () => {
expect(resolveCandidateWorktreeCapacityLimit(null, false)).toBeNull();
expect(resolveCandidateWorktreeCapacityLimit(null, true)).toBeNull();
});
});
describe("ledger mutations", () => {
it("dispatch TRANSFERS a held slot rather than adding one", () => {
expect(reserveWorktreeOnDispatch(4, true)).toBe(4);
});
it("dispatch ADDS a slot for a candidate holding no worktree", () => {
expect(reserveWorktreeOnDispatch(4, false)).toBe(5);
});
it("a failed retained transfer leaves the worktree ledger unchanged", () => {
expect(releaseWorktreeReservation(5, true)).toBe(5);
});
it("a failed fresh allocation gives its worktree slot back", () => {
expect(releaseWorktreeReservation(5, false)).toBe(4);
});
it("a failed dispatch gives the slot back", () => {
expect(releaseReservedSlot(5)).toBe(4);
});
it("release FLOORS at zero — a negative count would read as free capacity", () => {
/*
The floor is the whole point of the Math.max. Every later comparison in the loop treats the
reserved count as "slots in use"; a negative value silently hands out capacity that does not
exist, which is the under-count direction that admits work over the cap.
*/
expect(releaseReservedSlot(0)).toBe(0);
expect(releaseReservedSlot(-3)).toBe(0);
});
});

View File

@@ -1,17 +1,11 @@
/*
FNXC:WorkflowResolvedColumns 2026-08-01-03:35:
PLANNING ADMISSION'S WORKTREE LEDGER EXCLUDES TERMINAL LANES BY ROLE, NOT BY NAME.
FNXC:WorktreeCapacity 2026-08-01-04:38:
PLANNING ADMISSION'S WORKTREE LEDGER COUNTS LIVE TASKS BY ROLE, NOT RETAINED DIRECTORIES.
`374956ef23` gave planning admission a maxWorktrees dimension — correctly, a live board was running
8 planners against a 4-slot budget — and counted holders with
`t.column !== "done" && t.column !== "archived"`.
On a RENAMED board neither literal matches, so every FINISHED card whose worktree has not been
reaped yet still counts as a live holder. `heldWorktrees` only grows, `worktreeRoom` reaches zero,
and planning admission is withheld forever on a board with free slots. That is the mirror of the
breach the commit set out to fix and strictly worse: the breach is visible as 8 planners, the stall
is silent — cards sit "Queued to plan" and the recorded reason names the worktree budget, which the
operator then checks and finds has room.
The worktree cap limits simultaneous task activity. Planning, WIP, and live review sessions count;
queued, paused, dependency-blocked, and terminal cards do not count even when their worktree path is
retained for reuse or cleanup. The canonical running-agent predicate resolves renamed workflow roles
before making that decision.
WHY THIS FILE RATHER THAN A CASE IN `triage-plan-admission-throttle-audit.test.ts`: that suite's
fixture deliberately exhausts `maxConcurrent` (1 slot, consumed by a running card) so the AGENT gate
@@ -19,10 +13,6 @@ binds. The agent gate is checked with the worktree gate via `Math.min`, so an ex
masks the worktree answer entirely — the conversion under test would never decide anything. This
fixture leaves agent headroom so the worktree ledger is the only gate that can bind.
MEASURED FIRST, per the blinding procedure: before this file existed, reverting the resolver to the
two literals left all 19 tests in the capacity suites green. Nothing in the tree could tell the
conversion from what it replaced.
The assertion is the run-audit event, not a return value: `task:plan-admission-throttled` is what
the withhold path DOES, and it names the binding gate. A boolean "did anything start" would also be
satisfied by an unrelated early return.
@@ -31,6 +21,7 @@ satisfied by an unrelated early return.
import { describe, expect, it, vi, beforeEach } from "vitest";
import type { Settings, Task, TaskStore } from "@fusion/core";
import { TriageProcessor } from "../triage.js";
import { projectAdmissionCoordinator } from "../concurrency.js";
vi.mock("@fusion/core", async (importOriginal) => {
const { createEngineCoreMock } = await import("../test/mockCore.js");
@@ -70,7 +61,12 @@ function task(id: string, column: string, extra: Partial<Task> = {}): Task {
* `maxConcurrent` is generous on purpose — admission takes `min(projectRoom, worktreeRoom)`, so an
* exhausted agent count would mask the ledger this file is about.
*/
function createStore(tasks: Task[], recorded: RecordedEvent[], maxWorktrees: number): TaskStore {
function createStore(
tasks: Task[],
recorded: RecordedEvent[],
maxWorktrees: number,
worktreeLimitEnabled: boolean = true,
): TaskStore {
return {
getTask: vi.fn(async (id: string) => {
const found = tasks.find((candidate) => candidate.id === id);
@@ -80,11 +76,11 @@ function createStore(tasks: Task[], recorded: RecordedEvent[], maxWorktrees: num
getSettings: vi.fn().mockResolvedValue({
maxConcurrent: 20,
maxWorktrees,
worktreeLimitEnabled,
pollIntervalMs: 600_000,
groupOverlappingFiles: false,
autoMerge: true,
} as Settings),
/* The only store read `resolveProjectColumnsForRoles` makes. */
listWorkflowDefinitions: vi.fn(async () => [{ ir: RENAMED_IR }]),
recordRunAuditEvent: vi.fn(async (event: { mutationType: string; target: string; metadata?: Record<string, unknown> }) => {
recorded.push({ type: event.mutationType, target: event.target, metadata: event.metadata });
@@ -116,14 +112,25 @@ function createStore(tasks: Task[], recorded: RecordedEvent[], maxWorktrees: num
} as unknown as TaskStore;
}
async function pollOnce(store: TaskStore): Promise<void> {
async function pollOnce(store: TaskStore): Promise<ReturnType<typeof vi.fn>> {
const processor = new TriageProcessor(store, "/tmp/fn-admission-renamed-root", {});
/*
FNXC:WorktreeCapacity 2026-08-01-04:38:
Admission is the contract under test; replace the planner body with a narrow seam so a passing
admission cannot launch a real provider subprocess from this unit test.
*/
const specifyTask = vi.fn(async () => undefined);
(processor as unknown as { specifyTask: (task: Task) => Promise<void> }).specifyTask = specifyTask;
/* poll() no-ops unless running; driving one pass directly keeps this time-independent. */
(processor as unknown as { running: boolean }).running = true;
await (processor as unknown as { poll: () => Promise<void> }).poll();
/* The audit write is fire-and-forget. */
await new Promise((resolve) => setImmediate(resolve));
processor.stop();
for (const [admitted] of specifyTask.mock.calls) {
projectAdmissionCoordinator.releaseReservation((admitted as Task).id);
}
return specifyTask;
}
describe("planning admission's worktree ledger on a renamed board", () => {
@@ -141,18 +148,60 @@ describe("planning admission's worktree ledger on a renamed board", () => {
task("FN-SHIPPED", "shipped", { worktree: "/tmp/wt-shipped" }),
], recorded, 1);
await pollOnce(store);
const specifyTask = await pollOnce(store);
/* Against the literals `shipped` is counted, worktreeRoom is 0, and admission is withheld. */
const throttle = recorded.filter((event) => event.type === "task:plan-admission-throttled");
expect(throttle, "a finished card must not consume a worktree slot").toHaveLength(0);
expect(specifyTask).toHaveBeenCalledOnce();
expect(specifyTask).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-WAITING" }));
});
it("does not count an inactive retained worktree against planning capacity", async () => {
const store = createStore([
task("FN-WAITING", "drafting"),
task("FN-PARKED", "drafting", {
worktree: "/tmp/wt-parked",
status: "queued",
paused: true,
}),
], recorded, 1);
const specifyTask = await pollOnce(store);
const throttle = recorded.filter((event) => event.type === "task:plan-admission-throttled");
expect(throttle, "an inactive retained directory must not consume a live-task slot").toHaveLength(0);
expect(specifyTask).toHaveBeenCalledOnce();
expect(specifyTask).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-WAITING" }));
});
it("keeps the planning worktree gate inert when the limit is disabled", async () => {
const store = createStore([
task("FN-WAITING", "drafting"),
], recorded, 0, false);
const specifyTask = await pollOnce(store);
expect(recorded.filter((event) => event.type === "task:plan-admission-throttled")).toHaveLength(0);
expect(specifyTask).toHaveBeenCalledOnce();
});
it("admits only one planner when one active-task worktree slot remains", async () => {
const store = createStore([
task("FN-WAITING-1", "drafting"),
task("FN-WAITING-2", "drafting"),
], recorded, 1);
const specifyTask = await pollOnce(store);
expect(specifyTask).toHaveBeenCalledOnce();
expect(specifyTask).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-WAITING-1" }));
});
/*
THE PAIRED POSITIVE. The case above asserts an ABSENCE, which a processor that never throttled at
all would also satisfy — including one broken into ignoring `maxWorktrees` entirely. This pins that
the SAME renamed board still withholds when a live card genuinely holds the last worktree, so the
first case is measuring the lane role rather than a dead gate.
THE PAIRED POSITIVE. The absence cases would also pass if maxWorktrees were ignored entirely. This
pins that the same renamed board still withholds when a task is genuinely live, so the inactive
cases measure liveness rather than a dead gate.
*/
it("still withholds when a card in a live RENAMED lane holds the last worktree", async () => {
const store = createStore([
@@ -160,9 +209,10 @@ describe("planning admission's worktree ledger on a renamed board", () => {
task("FN-LIVE", "building", { worktree: "/tmp/wt-live" }),
], recorded, 1);
await pollOnce(store);
const specifyTask = await pollOnce(store);
const throttle = recorded.filter((event) => event.type === "task:plan-admission-throttled");
expect(throttle.length, "a live worktree holder must still consume its slot").toBeGreaterThan(0);
expect(throttle.length, "a live task must still consume its worktree-capacity slot").toBeGreaterThan(0);
expect(specifyTask).not.toHaveBeenCalled();
});
});

View File

@@ -2,6 +2,8 @@ import {
compareTaskIdNumeric,
countRunningAgentTasks,
enrichRunningAgentTaskShape,
isRunningAgentTask,
resolveWorktreeCapacityLimit,
resolveWorkflowIrForTask,
type Task,
type WorkflowIrResolverStore,
@@ -17,6 +19,23 @@ export const PRIORITY_EXECUTE = 1;
/** Priority level for specification/triage agents — served last (default). */
export const PRIORITY_SPECIFY = 0;
/*
FNXC:WorktreeCapacity 2026-08-01-04:38:
Agent concurrency and worktree capacity count the same canonical live-task
population. Collapse them to one project admission ceiling so planning, execute,
and merge cannot each observe and claim the final worktree slot independently.
*/
export function resolveActiveTaskCapacityLimit(params: {
maxConcurrent: number;
maxWorktrees: number;
worktreeLimitEnabled?: boolean;
}): number {
const maxWorktrees = resolveWorktreeCapacityLimit(params);
return maxWorktrees === null
? params.maxConcurrent
: Math.min(params.maxConcurrent, maxWorktrees);
}
/** A task waiting to enter one of the top-level agent lanes. */
export interface AdmissionCandidate {
taskId: string;
@@ -90,6 +109,51 @@ export class ProjectAdmissionCoordinator {
return this.reservations.get(projectId)?.size ?? 0;
}
private async occupiedCount(params: {
projectId: string;
claimed: () => Promise<number> | number;
claimedTaskIds?: () => Promise<Iterable<string>> | Iterable<string>;
}): Promise<number> {
const claimed = await params.claimed();
if (!params.claimedTaskIds) return claimed + this.reservationCount(params.projectId);
const claimedIds = new Set(await params.claimedTaskIds());
let pendingReservations = 0;
for (const taskId of this.reservations.get(params.projectId) ?? []) {
if (!claimedIds.has(taskId)) pendingReservations += 1;
}
return claimed + pendingReservations;
}
/** Atomically reserve one live-task slot for a lane that performs its own dispatch sweep. */
async reserveIfAvailable(params: {
projectId: string;
taskId: string;
maxConcurrent: number;
claimed: () => Promise<number> | number;
claimedTaskIds?: () => Promise<Iterable<string>> | Iterable<string>;
}): Promise<boolean> {
const existing = this.draining.get(params.projectId);
if (existing) await existing;
let reserved = false;
const drain = (async () => {
const current = this.reservations.get(params.projectId);
if (current?.has(params.taskId)) {
reserved = true;
return;
}
if (await this.occupiedCount(params) >= params.maxConcurrent) return;
this.reserve(params.projectId, params.taskId);
reserved = true;
})();
this.draining.set(params.projectId, drain);
try {
await drain;
return reserved;
} finally {
if (this.draining.get(params.projectId) === drain) this.draining.delete(params.projectId);
}
}
/** Register a lane's refresh source. Re-registering replaces its prior source. */
registerProvider(providerId: string, provider: AdmissionProvider): () => void {
const projectProviders = this.providers.get(provider.projectId) ?? new Map<string, AdmissionProvider>();
@@ -106,6 +170,8 @@ export class ProjectAdmissionCoordinator {
projectId: string;
maxConcurrent: number;
claimed: () => Promise<number> | number;
/** Canonically live task ids, used to de-duplicate reservations after persistence catches up. */
claimedTaskIds?: () => Promise<Iterable<string>> | Iterable<string>;
/** One-shot source for callers that do not hold a durable lane registration. */
refresh?: () => Promise<AdmissionCandidate[]>;
semaphore?: Pick<AgentSemaphore, "tryAcquire" | "release">;
@@ -124,7 +190,7 @@ export class ProjectAdmissionCoordinator {
// in-memory handoffs to count until they either become live or are dropped.
// Persisted task rows lag a fire-and-forget lane start, so omitting these
// reservations lets a second coordinator pass over-admit one project.
if (candidates.length === 0 || (await params.claimed()) + this.reservationCount(params.projectId) >= params.maxConcurrent) return;
if (candidates.length === 0 || await this.occupiedCount(params) >= params.maxConcurrent) return;
// Older test/runtime semaphore wrappers predate tryAcquire. They still
// exercise project admission, while production semaphores atomically take
// the host slot here.
@@ -153,10 +219,10 @@ export class ProjectAdmissionCoordinator {
// The host semaphore is exhausted for everyone, not just this candidate:
// trying younger candidates cannot succeed, so stop rather than spin.
if (!acquiredHostSlot) return;
// Compatibility-only semaphore shims cannot hold a reservation. Their
// lane tests provide claimed() synchronously, while real host semaphores
// use this durable marker until take/drop below.
if (hasReservableHostSlot) this.reserve(params.projectId, winner.taskId);
// The project reservation is independent of the optional host semaphore.
// It bridges every lane's dispatch-to-persist gap, including runtimes
// where the cross-project semaphore is intentionally absent.
this.reserve(params.projectId, winner.taskId);
/*
FNXC:ConcurrencyAdmission 2026-07-26-10:35:
Unwind EXACTLY what this attempt took. Two ways a naive `semaphore.release()` corrupts
@@ -175,20 +241,9 @@ export class ProjectAdmissionCoordinator {
the explicit unwind.
*/
const releaseAttempt = () => {
if (!hasReservableHostSlot) return;
/*
FNXC:CapacityModel 2026-07-29-13:20: the host-slot release is now
UNCONDITIONAL across both branches. The pre-held branch used to release the
caller's semaphore through dropPreHeldExecutorSlot's second parameter;
with that parameter deleted it would otherwise unwind the registration and
the reservation while LEAKING the host slot this attempt acquired.
*/
if (hasPreHeldExecutorSlot(winner.taskId)) {
dropPreHeldExecutorSlot(winner.taskId);
} else {
this.releaseReservation(winner.taskId);
}
params.semaphore?.release();
dropPreHeldExecutorSlot(winner.taskId);
this.releaseReservation(winner.taskId);
if (hasReservableHostSlot) params.semaphore?.release();
};
try {
winner.reserve?.();
@@ -280,28 +335,36 @@ FNXC:GlobalConcurrencyControls 2026-07-14-18:30:
Operators reported live running-agent counts above the global concurrency cap (e.g. 5 running with cap 4). Live utilization counts every top-level slot holder (in-progress, planning triage, active in-review), but the scheduler only preflighted capacity and acquired the shared semaphore later inside the executor — so a card could sit in-progress (and count as running) while triage still saw free semaphore slots and filled the rest of the cap. Pre-held executor slots close that gap: tryAcquire before todo→in-progress, keep the slot until the executor/graph run claims and releases it, and admit triage against max(semaphore.activeCount, live running count).
FNXC:GlobalConcurrencyControls 2026-07-15-03:50:
Hard invariant: registerPreHeldExecutorSlot may only run immediately after a successful semaphore.tryAcquire() for that same task, and every registration must later be either take()d (caller releases the semaphore) or drop()d (releases the semaphore). The Set is process-local soft state decoupled from activeCount except via this discipline — acquire-without-register or register-without-acquire desyncs capacity accounting.
The semaphore claim and project-admission handoff are tracked separately. A
runtime without the optional host semaphore still needs the admission handoff
to bridge selection until its task row becomes canonically live.
*/
const preHeldExecutorSlots = new Set<string>();
const preHeldAdmissionReservations = new Set<string>();
/**
* Register a semaphore slot that was **just** acquired via `tryAcquire` for a task about to enter in-progress.
* Must not be called without a matching prior acquire; pair with take() or drop().
* Register a project-admission handoff and, when present, its host semaphore slot.
*/
export function registerPreHeldExecutorSlot(taskId: string): void {
preHeldExecutorSlots.add(taskId);
export function registerPreHeldExecutorSlot(taskId: string, semaphoreHeld = true): void {
preHeldAdmissionReservations.add(taskId);
if (semaphoreHeld) preHeldExecutorSlots.add(taskId);
}
/**
* Transfer ownership of a pre-held executor slot to the caller.
* Returns true when a slot was registered; the caller MUST release the underlying semaphore in its finally path.
*/
export function takePreHeldExecutorSlot(taskId: string): boolean {
export function takePreHeldExecutorSlot(taskId: string, retainAdmissionReservation = false): boolean {
const taken = preHeldExecutorSlots.delete(taskId);
if (taken) projectAdmissionCoordinator.releaseReservation(taskId);
if (!retainAdmissionReservation) releasePreHeldAdmissionReservation(taskId);
return taken;
}
/** Release only the project-admission bridge while retaining any host semaphore claim. */
export function releasePreHeldAdmissionReservation(taskId: string): void {
if (preHeldAdmissionReservations.delete(taskId)) projectAdmissionCoordinator.releaseReservation(taskId);
}
/*
FNXC:CapacityModel 2026-07-29-13:20 (drop the cross-project cap — pre-held slots):
Drop a pre-held slot without transferring ownership (failed reserve / cancelled
@@ -333,9 +396,9 @@ export function dropPreHeldExecutorSlot(taskId: string): boolean {
slot this call never held, INFLATING capacity — the opposite of the leak the
cleanup was guarding against.
*/
if (!preHeldExecutorSlots.delete(taskId)) return false;
projectAdmissionCoordinator.releaseReservation(taskId);
return true;
const heldSemaphore = preHeldExecutorSlots.delete(taskId);
if (preHeldAdmissionReservations.delete(taskId)) projectAdmissionCoordinator.releaseReservation(taskId);
return heldSemaphore;
}
/** Test/helper: whether a task currently has an unclaimed pre-held executor slot. */
@@ -346,6 +409,7 @@ export function hasPreHeldExecutorSlot(taskId: string): boolean {
/** Test helper: clear all pre-held registrations without releasing semaphore slots. */
export function clearPreHeldExecutorSlotsForTests(): void {
preHeldExecutorSlots.clear();
preHeldAdmissionReservations.clear();
}
/**
@@ -361,13 +425,30 @@ export function persistedTopLevelAgentSlots(tasks: Task[]): number {
* predicate; raw `persistedTopLevelAgentSlots` remains for pre-enriched tests
* and callers that cannot resolve an IR.
*/
export async function persistedTopLevelAgentSlotsFromStore(store: WorkflowIrResolverStore, tasks: Task[]): Promise<number> {
/*
FNXC:WorktreeCapacity 2026-08-01-04:38:
Expose the ids behind the canonical live-task count so capacity diagnostics and arithmetic use the
same enriched predicate as the dashboard. Retained worktree metadata is deliberately not an input.
*/
async function enrichedTopLevelAgentTasksFromStore(store: WorkflowIrResolverStore, tasks: Task[]) {
const irCache = new Map();
const enriched = await Promise.all(tasks.map(async (task) => {
return Promise.all(tasks.map(async (task) => {
const ir = await resolveWorkflowIrForTask(store, task.id, irCache);
return enrichRunningAgentTaskShape(task, ir);
}));
return countRunningAgentTasks(enriched);
}
export async function persistedTopLevelAgentTaskIdsFromStore(store: WorkflowIrResolverStore, tasks: Task[]): Promise<string[]> {
const enriched = await enrichedTopLevelAgentTasksFromStore(store, tasks);
const ids: string[] = [];
for (const task of enriched) {
if (isRunningAgentTask(task)) ids.push(task.id);
}
return ids;
}
export async function persistedTopLevelAgentSlotsFromStore(store: WorkflowIrResolverStore, tasks: Task[]): Promise<number> {
return countRunningAgentTasks(await enrichedTopLevelAgentTasksFromStore(store, tasks));
}
/**

View File

@@ -3520,7 +3520,10 @@ export class TaskExecutor {
*/
private async runWithExecutorSemaphore<T>(taskId: string, work: () => Promise<T>): Promise<T> {
const sem = this.options.semaphore;
if (!sem) return work();
if (!sem) {
takePreHeldExecutorSlot(taskId);
return work();
}
if (this.outerConcurrencyClaims.has(taskId)) {
return work();
}

View File

@@ -75,7 +75,11 @@ import type { RoutineRunner } from "./routine-runner.js";
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, projectAdmissionCoordinator } from "./concurrency.js";
import {
computeTopLevelConcurrencyClaimedFromStore,
projectAdmissionCoordinator,
resolveActiveTaskCapacityLimit,
} from "./concurrency.js";
import { canStartNextMergeBody } from "./merge-reclaim-policy.js";
import {
registerProjectVerificationLimit,
@@ -3888,6 +3892,7 @@ export class ProjectEngine {
}
}
let selected = false;
const admissionSettings = await store.getSettings();
/*
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 —
@@ -3908,7 +3913,11 @@ export class ProjectEngine {
*/
await projectAdmissionCoordinator.admitOldest({
projectId: cwd,
maxConcurrent: (await store.getSettings()).maxConcurrent ?? 2,
maxConcurrent: resolveActiveTaskCapacityLimit({
maxConcurrent: admissionSettings.maxConcurrent ?? 2,
maxWorktrees: admissionSettings.maxWorktrees ?? 4,
worktreeLimitEnabled: admissionSettings.worktreeLimitEnabled,
}),
claimed: async () => computeTopLevelConcurrencyClaimedFromStore({
store,
tasks: await store.listTasks({ slim: true, includeArchived: false }),

View File

@@ -16,12 +16,13 @@ import { existsSync } from "node:fs";
import { readFile } from "node:fs/promises";
import { join } from "node:path";
import {
computeTopLevelConcurrencyClaimedFromStore,
dropPreHeldExecutorSlot,
hasPreHeldExecutorSlot,
projectAdmissionCoordinator,
persistedTopLevelAgentTaskIdsFromStore,
recoverIdleSemaphoreLeakCandidate,
registerPreHeldExecutorSlot,
resolveActiveTaskCapacityLimit,
type AgentSemaphore,
} from "./concurrency.js";
import { planTaskWorktreePath, resolveTaskWorkingBranch } from "./worktree-names.js";
@@ -41,7 +42,7 @@ import { BacklogPressureReporter } from "./backlog-pressure-reporter.js";
import { UnlinkedMissionsAdvisoryReporter } from "./unlinked-missions-advisory-reporter.js";
import { createRunAuditor, generateSyntheticRunId } from "./run-audit.js";
import type { TaskMoveLanes } from "@fusion/core";
import { resolveProjectColumnsForRoles, resolveWorkflowIrForTask, resolveWorkflowIrById, resolveColumnFlags, resolveWorktreeCapacityLimit, resolveLifecycleColumns, isWipColumnRole, isReviewColumnRole, isCompleteColumnRole, isTerminalColumnRole, columnsWithFlag } from "@fusion/core";
import { resolveProjectColumnsForRoles, resolveWorkflowIrForTask, resolveWorkflowIrById, resolveColumnFlags, resolveWorktreeCapacityLimit, resolveLifecycleColumns, isWipColumnRole, isReviewColumnRole, isCompleteColumnRole, columnsWithFlag } from "@fusion/core";
import type { ColumnRoleTraitFlags } from "@fusion/core";
import type { WorkflowIr, WorkflowIrV2 } from "@fusion/core";
import { runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./hold-release.js";
@@ -701,8 +702,7 @@ function computeConcurrencyGateDiagnostic(params: {
* semaphore gate uses this instead of only in-progress agentSlots.
*/
topLevelClaimedSlots?: number;
/** FNXC:WorkflowScheduling 2026-07-31-23:50: every live worktree holder (wip + planning/review
* lanes), so the maxWorktrees holders diagnostic names who actually occupies the slots. */
/** FNXC:WorktreeCapacity 2026-08-01-04:38: active task ids occupying worktree-capacity slots. */
worktreeHolderTaskIds?: string[];
/** U6: additive per-column capacity gates (flag-ON only). Omitted → the legacy
* three-gate report is byte-identical. */
@@ -890,77 +890,7 @@ export interface SchedulerOptions {
* itself is also refreshed: if `pollIntervalMs` differs from the active
* timer, the `setInterval` is transparently restarted.
*/
/*
FNXC:WorktreeCapacity 2026-08-01-00:45:
THE WORKTREE-CAPACITY HOLDER SET, behind a seam so its arithmetic can be pinned.
Two live defects came out of these few lines and neither had behavioural coverage. Measured by the
author of #3262: blinding the terminal predicate to `false` leaves all 22 scheduler suites green.
UNDER-COUNT admits work over the cap. The commit that added this gate reports maxWorktrees=4 with
four planning sessions holding worktrees and a fifth admitted, because the ledger counted WIP cards
only and never learned to count planners.
OVER-COUNT self-deadlocks. A planned Ready card RETAINS its planning worktree for execution reuse,
so counting it as a holder blocked its own release: 2 wip + 3 idle-held = 5/4 and the first unpause
released 2 of 4 slots' worth of work.
The two errors are not symmetric — under-counting breaks the cap, over-counting only starves dispatch
— so both directions are pinned rather than just the one that reads as "safe".
Extracted rather than tested in place because the call site is inside `schedule()`, which a unit test
has no business standing up (same reasoning as `scheduler-load-lane-union.test.ts`).
*/
export function nonWipWorktreeHolderIdsOf(
tasks: readonly Task[],
wipTaskIds: readonly string[],
isTerminalColumnTask: (task: Task) => boolean,
): string[] {
const wipTaskIdSet = new Set(wipTaskIds);
return tasks
.filter((task) => !wipTaskIdSet.has(task.id)
&& !isTerminalColumnTask(task)
&& typeof task.worktree === "string" && task.worktree.length > 0)
.map((task) => task.id);
}
/**
* FNXC:WorktreeCapacity 2026-08-01-04:07:
* Resolve whether worktree capacity applies to this candidate. A task that already holds a
* worktree allocates no additional tree when it dispatches, so the worktree dimension must not
* gate that transfer. Agent and per-column capacity continue to apply independently.
*/
export function resolveCandidateWorktreeCapacityLimit(
maxWorktrees: number | null,
candidateHoldsWorktree: boolean,
): number | null {
return candidateHoldsWorktree ? null : maxWorktrees;
}
/**
* FNXC:WorktreeCapacity 2026-08-01-01:05:
* The ledger MUTATIONS, named so they can be pinned alongside the totals they modify.
*
* `reserveWorktreeOnDispatch` is the ledger counterpart of
* `resolveCandidateWorktreeCapacityLimit`: a candidate that already holds a worktree reuses it, so
* dispatch TRANSFERS the slot rather than adding one. The two must agree — bypassing allocation
* capacity for a transfer but incrementing anyway would leak a slot per dispatch until the cap
* wedged.
*
* `releaseWorktreeReservation` gives back only a slot that this candidate added.
* `releaseReservedSlot` carries the floor; its `Math.max(0, …)` stops a double-release from handing
* out capacity that does not exist — a negative reserved count reads as free slots to every later
* comparison in the loop.
*/
export function reserveWorktreeOnDispatch(reservedWorktreeSlots: number, candidateHoldsWorktree: boolean): number {
return candidateHoldsWorktree ? reservedWorktreeSlots : reservedWorktreeSlots + 1;
}
export function releaseWorktreeReservation(reservedWorktreeSlots: number, candidateHoldsWorktree: boolean): number {
return candidateHoldsWorktree ? reservedWorktreeSlots : releaseReservedSlot(reservedWorktreeSlots);
}
export function releaseReservedSlot(reservedSlots: number): number {
function releaseReservedSlot(reservedSlots: number): number {
return Math.max(0, reservedSlots - 1);
}
@@ -1049,7 +979,7 @@ export class Scheduler {
taskId: task.id,
projectId,
createdAt: task.createdAt,
reserve: () => { if (this.options.semaphore) registerPreHeldExecutorSlot(task.id); },
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
start: async () => {
this.coordinatorReadyTasks.delete(task.id);
this.coordinatorAdmittedTaskIds.add(task.id);
@@ -2223,6 +2153,11 @@ export class Scheduler {
worktreeLimitEnabled: settings.worktreeLimitEnabled,
});
const maxConcurrent = settings.maxConcurrent ?? this.options.maxConcurrent ?? 2;
const activeTaskLimit = resolveActiveTaskCapacityLimit({
maxConcurrent,
maxWorktrees: settings.maxWorktrees ?? this.options.maxWorktrees ?? 4,
worktreeLimitEnabled: settings.worktreeLimitEnabled,
});
/*
FNXC:WorkflowScheduling 2026-07-19-02:35 (U4/KTD-9):
Count active WIP reservations by the `wip` trait, not the literal
@@ -2297,65 +2232,16 @@ export class Scheduler {
isReviewColumnRole(columnFlagsForTask(task), task.column);
const wipTaskIds = tasks.filter(isWipColumnTask).map((task) => task.id);
/*
FNXC:WorkflowScheduling 2026-07-31-23:50 (maxWorktrees counted only WIP — live board breach):
Under plan-in-place EVERY lane's live card holds a real worktree — planning runs in the task
worktree (triage.ts) and review/merge keeps it — but this ledger counted WIP cards only. The
protection that used to catch the difference was the GLOBAL SEMAPHORE gate, whose FNXC below
says exactly this ("must include every live top-level agent holder (planning triage and active
in-review), otherwise the hold/release sweep can admit an executor on top of a full planner
fleet"); the two-number capacity model deleted the semaphore, and this gate never learned to
count planners. Observed live: maxWorktrees=4, four planning sessions each holding a worktree,
and a replan dispatch admitted as the FIFTH worktree because the gate read used=0/4.
Count every non-terminal task that HOLDS a worktree (`task.worktree` set) in addition to WIP
membership — wip cards without a worktree yet still reserve (they are about to acquire), and
terminal lanes are excluded because their retained worktrees are cleanup-owned, not capacity.
FNXC:WorktreeCapacity 2026-08-01-04:38 (inactive retained-worktree capacity inversion):
Worktree capacity is a LIVE-TASK budget, not a count of directories retained on disk. The
dashboard showed seven active tasks against maxWorktrees=9, but two dependency-blocked queued
cards retained worktree paths. Counting those inactive paths filled the ledger and prevented
the dependency-free roots from starting. Use the canonical enriched running-agent predicate
shared with the board; planning, WIP, and active review count, while queued/paused/terminal
tasks do not. Same-sweep reservations below keep newly released tasks visible immediately.
*/
/*
FNXC:WorkflowResolvedColumns 2026-07-31-20:55 (u12 — the ratchet caught this, correctly):
DELIBERATE-LITERAL — the unresolvable-workflow default. The trait path above is the real answer;
this arm is reached only when `columnFlagsForTask` returns undefined, i.e. the task's workflow
could not be read at all, and it then gives the same answer the pre-trait code gave.
Recorded rather than converted because there is nothing to convert TO: a task with no readable
workflow has no resolved lane, and treating it as non-terminal would count a finished card's
retained worktree against live capacity — the opposite of what the surrounding fix does.
Marker sits in the DECLARATION's leading comments, not inline: markers are read from a node's
leading comments, so a mid-expression one attaches to the wrong node and is silently ignored.
*/
/*
FNXC:WorkflowResolvedColumns 2026-08-01-10:05 (fleet — correcting "nothing to convert TO"):
There is: `isTerminalColumnRole` in core is this predicate, term for term. With flags it reads
`flags.complete === true || flags.archived === true`; without them it falls back to the legacy
complete/archived ids — the same two arms written out below, verified against `column-roles.ts`.
Its own doc-comment names this exact case: it exists "because the pattern
`column !== \"done\" && column !== \"archived\"` is the single most repeated shape in the backlog"
and "keeps callers from re-deriving it and from accidentally dropping one half."
The DELIBERATE reasoning above is otherwise correct and kept: the undefined-flags arm is a live
path and treating an unreadable workflow as non-terminal would count a finished card's retained
worktree against live capacity. That argument is about the FALLBACK's existence, not about where
the predicate lives — and the shared helper carries the identical fallback.
Second time in this file: `isWipColumnTask` two lines up records that it was itself once "a
hand-rolled copy of `isWipColumnRole`".
*/
const isTerminalColumnTask = (task: Task): boolean =>
isTerminalColumnRole(columnFlagsForTask(task), task.column);
const nonWipWorktreeHolderIds = nonWipWorktreeHolderIdsOf(tasks, wipTaskIds, isTerminalColumnTask);
/*
FNXC:WorkflowScheduling 2026-07-31-01:05 (self-deadlock in the widened ledger, observed live):
A planned Ready card RETAINS its planning worktree for execution reuse, so counting it as a
holder must not block ITS OWN release — on release the slot TRANSFERS (the card executes in
the same worktree), it does not add. Without this exclusion the first unpause released only
2 of 4 slots' worth of work: the two remaining Ready cards were gated out by the very
worktrees they would reuse (2 wip + 3 idle-held = 5/4). Candidates in this set subtract
their own slot from the gate and skip the dispatch increment.
*/
const nonWipWorktreeHolderIdSet = new Set(nonWipWorktreeHolderIds);
let reservedWorktreeSlots = wipTaskIds.length + nonWipWorktreeHolderIds.length;
const activeWorktreeTaskIds = await persistedTopLevelAgentTaskIdsFromStore(this.store, tasks);
let reservedWorktreeSlots = activeWorktreeTaskIds.length;
let reservedConcurrentSlots = wipTaskIds.length;
const inProgressTaskIds = wipTaskIds;
const dispatchPrepByTaskId = new Map<string, {
@@ -2962,28 +2848,15 @@ export class Scheduler {
}
}
const topLevelClaimedSlots = await computeTopLevelConcurrencyClaimedFromStore({
store: this.store,
tasks,
});
const candidateHoldsWorktree = nonWipWorktreeHolderIdSet.has(task.id);
/*
FNXC:WorktreeCapacity 2026-08-01-03:55 (live over-cap transfer deadlock):
A retained-worktree candidate consumes zero NEW worktree slots. Subtracting only its own
slot still deadlocked an already-over-cap durable ledger (12 - 1 remained 11 against 9),
and even a ledger exactly at cap produced zero slack and blocked the transfer. Remove only
the worktree dimension for this candidate; maxConcurrent, semaphore, and column capacity
still decide whether another agent may run. Candidates without a tree keep the full gate.
*/
const concurrencyDiagnostic = computeConcurrencyGateDiagnostic({
agentSlots: reservedConcurrentSlots,
maxConcurrent,
activeWorktrees: reservedWorktreeSlots,
maxWorktrees: resolveCandidateWorktreeCapacityLimit(maxWorktrees, candidateHoldsWorktree),
worktreeHolderTaskIds: [...inProgressTaskIds, ...nonWipWorktreeHolderIds],
maxWorktrees,
worktreeHolderTaskIds: [...activeWorktreeTaskIds, ...dispatchPrepByTaskId.keys()],
semaphore: this.options.semaphore,
inProgressTaskIds,
topLevelClaimedSlots,
topLevelClaimedSlots: reservedWorktreeSlots,
});
/*
FNXC:WorkflowScheduling 2026-06-23-20:58:
@@ -3003,9 +2876,42 @@ export class Scheduler {
return null;
}
/*
FNXC:WorktreeCapacity 2026-08-01-04:38:
Serialize the workflow scheduler's direct hold release with planning and merge admission.
All three lanes count the same active population, so independent snapshots must not each
claim the final slot. The reservation remains until the executor observes the persisted
WIP row and takes the handoff.
*/
const projectSlotReserved = await projectAdmissionCoordinator.reserveIfAvailable({
projectId: this.store.getRootDir(),
taskId: task.id,
maxConcurrent: activeTaskLimit,
claimed: () => activeWorktreeTaskIds.length,
});
if (!projectSlotReserved) {
if (reservedScope) {
activeScopes.delete(task.id);
activeScopeColumns.delete(task.id);
}
const bindingGate: ConcurrencyGateName = maxWorktrees !== null && maxWorktrees <= maxConcurrent
? "maxWorktrees"
: "maxConcurrent";
const reason = formatConcurrencyLimitReason({
...concurrencyDiagnostic,
available: 0,
bindingGates: [...new Set([...concurrencyDiagnostic.bindingGates, bindingGate])],
});
await this.store.updateTask(task.id, { status: "queued" });
await this.logDispatchQueuedReason(task.id, reason);
return null;
}
const sem = this.options.semaphore;
const coordinatorReserved = hasPreHeldExecutorSlot(task.id);
if (sem && !coordinatorReserved && !sem.tryAcquire()) {
const hostSlotReserved = hasPreHeldExecutorSlot(task.id);
registerPreHeldExecutorSlot(task.id, hostSlotReserved);
if (sem && !hostSlotReserved && !sem.tryAcquire()) {
dropPreHeldExecutorSlot(task.id);
if (reservedScope) {
activeScopes.delete(task.id);
activeScopeColumns.delete(task.id);
@@ -3019,7 +2925,7 @@ export class Scheduler {
await this.logDispatchQueuedReason(task.id, reason, formatConcurrencyLimitMemoKey(concurrencyDiagnostic));
return null;
}
if (sem && !coordinatorReserved) {
if (sem && !hostSlotReserved) {
registerPreHeldExecutorSlot(task.id);
}
@@ -3063,8 +2969,7 @@ export class Scheduler {
task: freshTask,
});
// Transfer, not addition, for a candidate that already holds its worktree.
reservedWorktreeSlots = reserveWorktreeOnDispatch(reservedWorktreeSlots, candidateHoldsWorktree);
reservedWorktreeSlots += 1;
reservedConcurrentSlots += 1;
let released = false;
return {
@@ -3075,7 +2980,7 @@ export class Scheduler {
activeScopes.delete(task.id);
activeScopeColumns.delete(task.id);
}
reservedWorktreeSlots = releaseWorktreeReservation(reservedWorktreeSlots, candidateHoldsWorktree);
reservedWorktreeSlots = releaseReservedSlot(reservedWorktreeSlots);
reservedConcurrentSlots = releaseReservedSlot(reservedConcurrentSlots);
dispatchPrepByTaskId.delete(task.id);
if (dropPreHeldExecutorSlot(task.id)) sem?.release();
@@ -3089,6 +2994,9 @@ export class Scheduler {
this.planWorktreePath(task, settings.worktreeNaming, reservedNames, settings),
});
for (const taskId of result.released) {
// The authoritative move has made this task visible to canonical live-task
// counting; the transient project reservation no longer needs to bridge it.
projectAdmissionCoordinator.releaseReservation(taskId);
const prep = dispatchPrepByTaskId.get(taskId);
if (!prep) continue;
/*

View File

@@ -37,6 +37,7 @@ import {
resolveLifecycleColumns,
resolveWorkflowIrForTaskWithProvenance,
resolveProjectColumnsForRoles,
resolveWorktreeCapacityLimit,
workflowHasColumn,
getStepParser,
computePlanApprovalFingerprint,
@@ -137,8 +138,11 @@ import {
PRIORITY_SPECIFY,
computeTopLevelConcurrencyClaimedFromStore,
dropPreHeldExecutorSlot,
persistedTopLevelAgentTaskIdsFromStore,
projectAdmissionCoordinator,
registerPreHeldExecutorSlot,
releasePreHeldAdmissionReservation,
resolveActiveTaskCapacityLimit,
takePreHeldExecutorSlot,
recoverIdleSemaphoreLeakCandidate,
type AgentSemaphore,
@@ -581,7 +585,7 @@ export class TriageProcessor {
);
return tasks.filter((task) => !this.coordinatorAdmittedTaskIds.has(task.id)).map((task) => ({
taskId: task.id, projectId: this.rootDir, createdAt: task.createdAt,
reserve: () => { if (this.options.semaphore) registerPreHeldExecutorSlot(task.id); },
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
start: async () => {
this.coordinatorAdmittedTaskIds.add(task.id);
void this.specifyTask(task);
@@ -2128,43 +2132,24 @@ export class TriageProcessor {
*/
const projectRoom = Math.max(0, maxConcurrent - claimed);
/*
FNXC:CapacityModel 2026-08-01-02:15 (live breach — 8 planners on a maxWorktrees=4 board):
Planning admission gated ONLY on the agent count. Every planner acquires a REAL worktree
(plan-in-place), so admission must also respect the worktree budget — the same two-number
model the scheduler's dispatch gate enforces. This gap was invisible for days because the
merge-inside-admission-drain bug (00769fad7c) froze admission during every merge and
accidentally throttled planners; unfreezing it exposed 8 concurrent planning sessions and 11
worktrees on a 4-slot board within one restart.
The ledger mirrors the scheduler's: every non-terminal task HOLDING a worktree occupies a
slot, and a candidate that already holds one (replan re-entry) transfers rather than adds —
only worktree-less candidates consume remaining slots, so a held tree never blocks its own
resume. maxWorktrees is absent/null when worktrees are off; admission then falls back to the
agent gate alone.
FNXC:WorktreeCapacity 2026-08-01-04:38:
Worktree slots follow the canonical LIVE-TASK count (`claimed`), not retained directory
metadata. A queued, paused, dependency-blocked, or terminal card may keep a worktree path for
reuse/cleanup without consuming admission capacity. Every newly admitted planner becomes live
and spends one slot below, even when it reuses an existing directory.
*/
const maxWorktrees = (settings as { maxWorktrees?: number | null }).maxWorktrees ?? 4;
let worktreeRoom = Number.POSITIVE_INFINITY;
if (typeof maxWorktrees === "number" && Number.isFinite(maxWorktrees)) {
/*
FNXC:WorkflowResolvedColumns 2026-08-01-03:05:
TERMINAL IS A ROLE HERE, NOT A NAME. The ledger above excludes terminal lanes because their
retained worktrees are cleanup-owned rather than capacity; against the literals `done` and
`archived` a RENAMED board matched neither, so every finished card still counted as a live
worktree holder. The gate then reads worktreeRoom=0 on a board with free slots and planning
admission stalls permanently — the mirror of the breach this commit set out to fix, and
strictly worse, because a stall is silent where a breach is at least visible as 8 planners.
PROJECT-level, matching this file's existing use at `sweepStalePlanningStatuses`: the ledger
spans every card on the board, so there is no single task to resolve against.
`resolveProjectColumnsForRoles` is legacy-seeded, so a default board still excludes exactly
`done` and `archived` and this conversion is byte-identical there.
*/
const terminalColumns = await resolveProjectColumnsForRoles(this.store, ["complete", "archived"]);
const heldWorktrees = allTasks.filter((t) =>
!terminalColumns.has(t.column)
&& typeof t.worktree === "string" && t.worktree.length > 0).length;
worktreeRoom = Math.max(0, maxWorktrees - heldWorktrees);
}
const maxWorktrees = resolveWorktreeCapacityLimit({
maxWorktrees: settings.maxWorktrees ?? 4,
worktreeLimitEnabled: settings.worktreeLimitEnabled,
});
const activeTaskLimit = resolveActiveTaskCapacityLimit({
maxConcurrent,
maxWorktrees: settings.maxWorktrees ?? 4,
worktreeLimitEnabled: settings.worktreeLimitEnabled,
});
const worktreeRoom = maxWorktrees === null
? Number.POSITIVE_INFINITY
: Math.max(0, maxWorktrees - claimed);
const maxToStart = Math.min(projectRoom, worktreeRoom);
if (maxToStart <= 0 && triageTasks.length > 0) {
@@ -2276,36 +2261,33 @@ export class TriageProcessor {
// the planner's synchronous processing claim until after this poll returns.
const admittedThisPoll = new Set<string>();
/*
FNXC:CapacityModel 2026-08-01-02:20: the transfer rule from the scheduler ledger, applied at
admission — a candidate that already HOLDS a worktree (replan re-entry) reuses it, so it
spends only agent room, never a fresh worktree slot. Budgets decrement per admission below.
FNXC:WorktreeCapacity 2026-08-01-04:38:
Each admitted planner becomes an active task and therefore spends one worktree-capacity slot.
Whether its directory is newly created or retained is deliberately irrelevant.
*/
let agentBudget = projectRoom;
let worktreeBudget = worktreeRoom;
for (let i = 0; i < triageTasks.length; i++) {
if (agentBudget <= 0) break;
const candidate = triageTasks[i];
const candidateHoldsWorktree = typeof candidate.worktree === "string" && candidate.worktree.length > 0;
if (!candidateHoldsWorktree && worktreeBudget <= 0) continue;
if (agentBudget <= 0 || worktreeBudget <= 0) break;
agentBudget -= 1;
if (!candidateHoldsWorktree) worktreeBudget -= 1;
worktreeBudget -= 1;
let freshClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
const getFreshClaimSnapshot = () => freshClaimSnapshot ??= (async () => {
const fresh = await this.store.listTasks({ slim: true, includeArchived: false });
let pending = 0;
for (const id of this.processing) {
const row = fresh.find((task) => task.id === id);
if (!row || row.status !== "planning") pending++;
}
const ids = await persistedTopLevelAgentTaskIdsFromStore(this.store, fresh);
return { count: ids.length + pending, ids: [...new Set([...ids, ...this.processing])] };
})();
await projectAdmissionCoordinator.admitOldest({
// rootDir is the stable per-project identity held by this processor.
projectId: this.rootDir,
maxConcurrent,
claimed: async () => {
const fresh = await this.store.listTasks({ slim: true, includeArchived: false });
let pending = 0;
for (const id of this.processing) {
const row = fresh.find((task) => task.id === id);
if (!row || row.status !== "planning") pending++;
}
return computeTopLevelConcurrencyClaimedFromStore({
store: this.store,
tasks: fresh,
pendingSpecifyCount: pending,
});
},
maxConcurrent: activeTaskLimit,
claimed: async () => (await getFreshClaimSnapshot()).count,
claimedTaskIds: async () => (await getFreshClaimSnapshot()).ids,
semaphore: this.options.semaphore,
refresh: async () => triageTasks
.filter((task) => !admittedThisPoll.has(task.id) && !this.coordinatorAdmittedTaskIds.has(task.id) && !this.processing.has(task.id) && !this.hasLivePlanningWork(task.id))
@@ -2316,7 +2298,7 @@ export class TriageProcessor {
// FNXC:ConcurrencyAdmission 2026-08-05-10:00: the planner must
// own the coordinator's real host reservation before it starts;
// deferring to semaphore.run would reintroduce priority overtaking.
reserve: () => { if (this.options.semaphore) registerPreHeldExecutorSlot(task.id); },
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
start: async () => {
admittedThisPoll.add(task.id);
this.coordinatorAdmittedTaskIds.add(task.id);
@@ -2521,7 +2503,13 @@ export class TriageProcessor {
updatePlanningStateIfStillCurrent); this one must be visible — the FN-7977 steps>0
wedge stalled the whole planner for hours precisely because it logged nothing.
*/
if (!await this.updatePlanningStateIfStillCurrent(task, { status: "planning" })) {
let planningClaimed = false;
try {
planningClaimed = await this.updatePlanningStateIfStillCurrent(task, { status: "planning" });
} finally {
releasePreHeldAdmissionReservation(task.id);
}
if (!planningClaimed) {
planLog.warn(
`${task.id}: planning claim skipped — live row is no longer in the planning stage; `
+ "it will be re-claimed on the next poll",
@@ -3256,7 +3244,8 @@ export class TriageProcessor {
},
});
if (this.options.semaphore && takePreHeldExecutorSlot(task.id)) {
const heldHostSlot = takePreHeldExecutorSlot(task.id, true);
if (this.options.semaphore && heldHostSlot) {
// Coordinator already owns this top-level slot; run directly so it
// cannot join the priority queue after age-based admission.
try {