diff --git a/.changeset/fix-active-worktree-slot-accounting.md b/.changeset/fix-active-worktree-slot-accounting.md new file mode 100644 index 0000000000..157031b49c --- /dev/null +++ b/.changeset/fix-active-worktree-slot-accounting.md @@ -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. diff --git a/packages/engine/src/__tests__/concurrency.test.ts b/packages/engine/src/__tests__/concurrency.test.ts index b5fdad1b61..fc2acd8705 100644 --- a/packages/engine/src/__tests__/concurrency.test.ts +++ b/packages/engine/src/__tests__/concurrency.test.ts @@ -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[] = []; diff --git a/packages/engine/src/__tests__/scheduler-terminal-capacity-role.test.ts b/packages/engine/src/__tests__/scheduler-terminal-capacity-role.test.ts index 150ba63785..ae4dfde998 100644 --- a/packages/engine/src/__tests__/scheduler-terminal-capacity-role.test.ts +++ b/packages/engine/src/__tests__/scheduler-terminal-capacity-role.test.ts @@ -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 => ({ + 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); }); }); diff --git a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts index c0daead2f0..dddce4448e 100644 --- a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts +++ b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts @@ -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 () => { diff --git a/packages/engine/src/__tests__/scheduler-worktree-capacity-arithmetic.test.ts b/packages/engine/src/__tests__/scheduler-worktree-capacity-arithmetic.test.ts deleted file mode 100644 index b4797ebfb7..0000000000 --- a/packages/engine/src/__tests__/scheduler-worktree-capacity-arithmetic.test.ts +++ /dev/null @@ -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 => ({ - 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); - }); -}); diff --git a/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts b/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts index 055be69691..be4a463c92 100644 --- a/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts +++ b/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts @@ -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 { * `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 }) => { 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 { +async function pollOnce(store: TaskStore): Promise> { 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 }).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 }).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(); }); }); diff --git a/packages/engine/src/concurrency.ts b/packages/engine/src/concurrency.ts index def9035c8f..4702732043 100644 --- a/packages/engine/src/concurrency.ts +++ b/packages/engine/src/concurrency.ts @@ -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; + claimedTaskIds?: () => Promise> | Iterable; + }): Promise { + 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; + claimedTaskIds?: () => Promise> | Iterable; + }): Promise { + 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(); @@ -106,6 +170,8 @@ export class ProjectAdmissionCoordinator { projectId: string; maxConcurrent: number; claimed: () => Promise | number; + /** Canonically live task ids, used to de-duplicate reservations after persistence catches up. */ + claimedTaskIds?: () => Promise> | Iterable; /** One-shot source for callers that do not hold a durable lane registration. */ refresh?: () => Promise; semaphore?: Pick; @@ -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(); +const preHeldAdmissionReservations = new Set(); /** - * 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 { +/* +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 { + 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 { + return countRunningAgentTasks(await enrichedTopLevelAgentTasksFromStore(store, tasks)); } /** diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index 03518be443..abe7db6870 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -3520,7 +3520,10 @@ export class TaskExecutor { */ private async runWithExecutorSemaphore(taskId: string, work: () => Promise): Promise { const sem = this.options.semaphore; - if (!sem) return work(); + if (!sem) { + takePreHeldExecutorSlot(taskId); + return work(); + } if (this.outerConcurrencyClaims.has(taskId)) { return work(); } diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index 17924952c7..030267dc87 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -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 }), diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts index cc783c9207..9ec35dcf70 100644 --- a/packages/engine/src/scheduler.ts +++ b/packages/engine/src/scheduler.ts @@ -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 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; /* diff --git a/packages/engine/src/triage.ts b/packages/engine/src/triage.ts index 31c3d6c141..b67c7f54ab 100644 --- a/packages/engine/src/triage.ts +++ b/packages/engine/src/triage.ts @@ -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(); /* - 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 {