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:
7
.changeset/fix-active-worktree-slot-accounting.md
Normal file
7
.changeset/fix-active-worktree-slot-accounting.md
Normal 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.
|
||||||
@@ -15,6 +15,7 @@ import {
|
|||||||
persistedTopLevelAgentSlots,
|
persistedTopLevelAgentSlots,
|
||||||
recoverIdleSemaphoreLeakCandidate,
|
recoverIdleSemaphoreLeakCandidate,
|
||||||
registerPreHeldExecutorSlot,
|
registerPreHeldExecutorSlot,
|
||||||
|
resolveActiveTaskCapacityLimit,
|
||||||
takePreHeldExecutorSlot,
|
takePreHeldExecutorSlot,
|
||||||
} from "../concurrency.js";
|
} from "../concurrency.js";
|
||||||
|
|
||||||
@@ -1067,6 +1068,47 @@ describe("AgentSemaphore resilience (FN-978)", () => {
|
|||||||
|
|
||||||
|
|
||||||
describe("ProjectAdmissionCoordinator", () => {
|
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 () => {
|
it("admits the oldest same-project candidate atomically and partitions projects", async () => {
|
||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const started: string[] = [];
|
const started: string[] = [];
|
||||||
|
|||||||
@@ -1,30 +1,14 @@
|
|||||||
// @vitest-environment node
|
// @vitest-environment node
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowResolvedColumns 2026-08-01-09:20 (fleet):
|
FNXC:WorktreeCapacity 2026-08-01-04:38:
|
||||||
|
The worktree-capacity ledger uses the canonical enriched live-task predicate. Renamed workflow roles
|
||||||
THE INVARIANT: the worktree-capacity read excludes terminal lanes by ROLE, on any board.
|
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.
|
||||||
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.
|
|
||||||
*/
|
*/
|
||||||
|
|
||||||
import { describe, expect, it } from "vitest";
|
import { describe, expect, it } from "vitest";
|
||||||
import { isTerminalColumnRole, resolveColumnFlags } from "@fusion/core";
|
import { enrichRunningAgentTaskShape, isRunningAgentTask } from "@fusion/core";
|
||||||
import type { WorkflowIr } from "@fusion/core";
|
import type { Task, WorkflowIr } from "@fusion/core";
|
||||||
|
|
||||||
const RENAMED_IR = {
|
const RENAMED_IR = {
|
||||||
version: "v2", id: "wf-renamed", name: "renamed", nodes: [], edges: [],
|
version: "v2", id: "wf-renamed", name: "renamed", nodes: [], edges: [],
|
||||||
@@ -36,29 +20,36 @@ const RENAMED_IR = {
|
|||||||
],
|
],
|
||||||
} as unknown as WorkflowIr;
|
} as unknown as WorkflowIr;
|
||||||
|
|
||||||
const flagsFor = (columnId: string) => resolveColumnFlags(
|
const task = (id: string, column: string, overrides: Partial<Task> = {}): Task => ({
|
||||||
RENAMED_IR.columns.find((c) => c.id === columnId) as never,
|
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", () => {
|
const isLive = (candidate: Task): boolean =>
|
||||||
it("excludes BOTH renamed terminal lanes, not just the complete one", () => {
|
isRunningAgentTask(enrichRunningAgentTaskShape(candidate, RENAMED_IR));
|
||||||
expect(isTerminalColumnRole(flagsFor("shipped"), "shipped")).toBe(true);
|
|
||||||
expect(isTerminalColumnRole(flagsFor("attic"), "attic")).toBe(true);
|
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", () => {
|
it("ignores an inactive retained directory", () => {
|
||||||
expect(isTerminalColumnRole(flagsFor("building"), "building")).toBe(false);
|
expect(isLive(task("FN-QUEUED", "backlog", {
|
||||||
expect(isTerminalColumnRole(flagsFor("backlog"), "backlog")).toBe(false);
|
status: "queued",
|
||||||
|
worktree: "/tmp/project/.worktrees/queued",
|
||||||
|
}))).toBe(false);
|
||||||
});
|
});
|
||||||
|
|
||||||
/*
|
it("ignores complete and archived cards even with stale live-looking metadata", () => {
|
||||||
The degraded arm is why the shared helper is used rather than the hand-rolled version: an
|
expect(isLive(task("FN-DONE", "shipped", { status: "planning" }))).toBe(false);
|
||||||
unresolvable column must keep answering for the legacy ids, or a board mid-migration starts counting
|
expect(isLive(task("FN-ARCHIVED", "attic", { sessionFile: "/tmp/stale-session" }))).toBe(false);
|
||||||
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);
|
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -460,7 +460,7 @@ describe("Scheduler workflow cutover", () => {
|
|||||||
expect(ready.column).toBe("todo");
|
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 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 ready = task({ id: "FN-200", status: "queued", worktree: "/tmp/project/.worktrees/fn-200" });
|
||||||
const store = storeWith([active, ready], { maxConcurrent: 4, maxWorktrees: 1 });
|
const store = storeWith([active, ready], { maxConcurrent: 4, maxWorktrees: 1 });
|
||||||
@@ -470,8 +470,50 @@ describe("Scheduler workflow cutover", () => {
|
|||||||
|
|
||||||
await scheduler.schedule();
|
await scheduler.schedule();
|
||||||
|
|
||||||
expect(store.moveTaskIf).toHaveBeenCalledWith("FN-200", "in-progress", expect.any(Function), expect.anything());
|
expect(store.moveTaskIf).not.toHaveBeenCalledWith("FN-200", "in-progress", expect.anything(), expect.anything());
|
||||||
expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-200", column: "in-progress" }));
|
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");
|
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 retained = task({ id: "FN-001", status: "queued", worktree: "/tmp/project/.worktrees/fn-001" });
|
||||||
const fresh = task({ id: "FN-002", status: "queued" });
|
const fresh = task({ id: "FN-002", status: "queued" });
|
||||||
const store = storeWith([retained, fresh], { maxConcurrent: 4, maxWorktrees: 1 });
|
const store = storeWith([retained, fresh], { maxConcurrent: 4, maxWorktrees: 1 });
|
||||||
@@ -635,11 +677,11 @@ describe("Scheduler workflow cutover", () => {
|
|||||||
|
|
||||||
await scheduler.schedule();
|
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).toHaveBeenCalledWith("FN-001", "in-progress", expect.any(Function), expect.anything());
|
||||||
expect(store.moveTaskIf).not.toHaveBeenCalledWith("FN-002", "in-progress", expect.anything(), expect.anything());
|
expect(store.moveTaskIf).toHaveBeenCalledWith("FN-002", "in-progress", expect.any(Function), expect.anything());
|
||||||
expect(onSchedule).not.toHaveBeenCalled();
|
expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-002", column: "in-progress" }));
|
||||||
expect(fresh.column).toBe("todo");
|
expect(fresh.column).toBe("in-progress");
|
||||||
});
|
});
|
||||||
|
|
||||||
it("does not release work when the shared semaphore is saturated", async () => {
|
it("does not release work when the shared semaphore is saturated", async () => {
|
||||||
|
|||||||
@@ -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);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,17 +1,11 @@
|
|||||||
/*
|
/*
|
||||||
FNXC:WorkflowResolvedColumns 2026-08-01-03:35:
|
FNXC:WorktreeCapacity 2026-08-01-04:38:
|
||||||
PLANNING ADMISSION'S WORKTREE LEDGER EXCLUDES TERMINAL LANES BY ROLE, NOT BY NAME.
|
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
|
The worktree cap limits simultaneous task activity. Planning, WIP, and live review sessions count;
|
||||||
8 planners against a 4-slot budget — and counted holders with
|
queued, paused, dependency-blocked, and terminal cards do not count even when their worktree path is
|
||||||
`t.column !== "done" && t.column !== "archived"`.
|
retained for reuse or cleanup. The canonical running-agent predicate resolves renamed workflow roles
|
||||||
|
before making that decision.
|
||||||
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.
|
|
||||||
|
|
||||||
WHY THIS FILE RATHER THAN A CASE IN `triage-plan-admission-throttle-audit.test.ts`: that suite's
|
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
|
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
|
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.
|
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 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
|
the withhold path DOES, and it names the binding gate. A boolean "did anything start" would also be
|
||||||
satisfied by an unrelated early return.
|
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 { describe, expect, it, vi, beforeEach } from "vitest";
|
||||||
import type { Settings, Task, TaskStore } from "@fusion/core";
|
import type { Settings, Task, TaskStore } from "@fusion/core";
|
||||||
import { TriageProcessor } from "../triage.js";
|
import { TriageProcessor } from "../triage.js";
|
||||||
|
import { projectAdmissionCoordinator } from "../concurrency.js";
|
||||||
|
|
||||||
vi.mock("@fusion/core", async (importOriginal) => {
|
vi.mock("@fusion/core", async (importOriginal) => {
|
||||||
const { createEngineCoreMock } = await import("../test/mockCore.js");
|
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
|
* `maxConcurrent` is generous on purpose — admission takes `min(projectRoom, worktreeRoom)`, so an
|
||||||
* exhausted agent count would mask the ledger this file is about.
|
* 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 {
|
return {
|
||||||
getTask: vi.fn(async (id: string) => {
|
getTask: vi.fn(async (id: string) => {
|
||||||
const found = tasks.find((candidate) => candidate.id === id);
|
const found = tasks.find((candidate) => candidate.id === id);
|
||||||
@@ -80,11 +76,11 @@ function createStore(tasks: Task[], recorded: RecordedEvent[], maxWorktrees: num
|
|||||||
getSettings: vi.fn().mockResolvedValue({
|
getSettings: vi.fn().mockResolvedValue({
|
||||||
maxConcurrent: 20,
|
maxConcurrent: 20,
|
||||||
maxWorktrees,
|
maxWorktrees,
|
||||||
|
worktreeLimitEnabled,
|
||||||
pollIntervalMs: 600_000,
|
pollIntervalMs: 600_000,
|
||||||
groupOverlappingFiles: false,
|
groupOverlappingFiles: false,
|
||||||
autoMerge: true,
|
autoMerge: true,
|
||||||
} as Settings),
|
} as Settings),
|
||||||
/* The only store read `resolveProjectColumnsForRoles` makes. */
|
|
||||||
listWorkflowDefinitions: vi.fn(async () => [{ ir: RENAMED_IR }]),
|
listWorkflowDefinitions: vi.fn(async () => [{ ir: RENAMED_IR }]),
|
||||||
recordRunAuditEvent: vi.fn(async (event: { mutationType: string; target: string; metadata?: Record<string, unknown> }) => {
|
recordRunAuditEvent: vi.fn(async (event: { mutationType: string; target: string; metadata?: Record<string, unknown> }) => {
|
||||||
recorded.push({ type: event.mutationType, target: event.target, metadata: event.metadata });
|
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;
|
} 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", {});
|
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. */
|
/* poll() no-ops unless running; driving one pass directly keeps this time-independent. */
|
||||||
(processor as unknown as { running: boolean }).running = true;
|
(processor as unknown as { running: boolean }).running = true;
|
||||||
await (processor as unknown as { poll: () => Promise<void> }).poll();
|
await (processor as unknown as { poll: () => Promise<void> }).poll();
|
||||||
/* The audit write is fire-and-forget. */
|
/* The audit write is fire-and-forget. */
|
||||||
await new Promise((resolve) => setImmediate(resolve));
|
await new Promise((resolve) => setImmediate(resolve));
|
||||||
processor.stop();
|
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", () => {
|
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" }),
|
task("FN-SHIPPED", "shipped", { worktree: "/tmp/wt-shipped" }),
|
||||||
], recorded, 1);
|
], recorded, 1);
|
||||||
|
|
||||||
await pollOnce(store);
|
const specifyTask = await pollOnce(store);
|
||||||
|
|
||||||
/* Against the literals `shipped` is counted, worktreeRoom is 0, and admission is withheld. */
|
/* Against the literals `shipped` is counted, worktreeRoom is 0, and admission is withheld. */
|
||||||
const throttle = recorded.filter((event) => event.type === "task:plan-admission-throttled");
|
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(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
|
THE PAIRED POSITIVE. The absence cases would also pass if maxWorktrees were ignored entirely. This
|
||||||
all would also satisfy — including one broken into ignoring `maxWorktrees` entirely. This pins that
|
pins that the same renamed board still withholds when a task is genuinely live, so the inactive
|
||||||
the SAME renamed board still withholds when a live card genuinely holds the last worktree, so the
|
cases measure liveness rather than a dead gate.
|
||||||
first case is measuring the lane role rather than a dead gate.
|
|
||||||
*/
|
*/
|
||||||
it("still withholds when a card in a live RENAMED lane holds the last worktree", async () => {
|
it("still withholds when a card in a live RENAMED lane holds the last worktree", async () => {
|
||||||
const store = createStore([
|
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" }),
|
task("FN-LIVE", "building", { worktree: "/tmp/wt-live" }),
|
||||||
], recorded, 1);
|
], recorded, 1);
|
||||||
|
|
||||||
await pollOnce(store);
|
const specifyTask = await pollOnce(store);
|
||||||
|
|
||||||
const throttle = recorded.filter((event) => event.type === "task:plan-admission-throttled");
|
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();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ import {
|
|||||||
compareTaskIdNumeric,
|
compareTaskIdNumeric,
|
||||||
countRunningAgentTasks,
|
countRunningAgentTasks,
|
||||||
enrichRunningAgentTaskShape,
|
enrichRunningAgentTaskShape,
|
||||||
|
isRunningAgentTask,
|
||||||
|
resolveWorktreeCapacityLimit,
|
||||||
resolveWorkflowIrForTask,
|
resolveWorkflowIrForTask,
|
||||||
type Task,
|
type Task,
|
||||||
type WorkflowIrResolverStore,
|
type WorkflowIrResolverStore,
|
||||||
@@ -17,6 +19,23 @@ export const PRIORITY_EXECUTE = 1;
|
|||||||
/** Priority level for specification/triage agents — served last (default). */
|
/** Priority level for specification/triage agents — served last (default). */
|
||||||
export const PRIORITY_SPECIFY = 0;
|
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. */
|
/** A task waiting to enter one of the top-level agent lanes. */
|
||||||
export interface AdmissionCandidate {
|
export interface AdmissionCandidate {
|
||||||
taskId: string;
|
taskId: string;
|
||||||
@@ -90,6 +109,51 @@ export class ProjectAdmissionCoordinator {
|
|||||||
return this.reservations.get(projectId)?.size ?? 0;
|
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. */
|
/** Register a lane's refresh source. Re-registering replaces its prior source. */
|
||||||
registerProvider(providerId: string, provider: AdmissionProvider): () => void {
|
registerProvider(providerId: string, provider: AdmissionProvider): () => void {
|
||||||
const projectProviders = this.providers.get(provider.projectId) ?? new Map<string, AdmissionProvider>();
|
const projectProviders = this.providers.get(provider.projectId) ?? new Map<string, AdmissionProvider>();
|
||||||
@@ -106,6 +170,8 @@ export class ProjectAdmissionCoordinator {
|
|||||||
projectId: string;
|
projectId: string;
|
||||||
maxConcurrent: number;
|
maxConcurrent: number;
|
||||||
claimed: () => Promise<number> | 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. */
|
/** One-shot source for callers that do not hold a durable lane registration. */
|
||||||
refresh?: () => Promise<AdmissionCandidate[]>;
|
refresh?: () => Promise<AdmissionCandidate[]>;
|
||||||
semaphore?: Pick<AgentSemaphore, "tryAcquire" | "release">;
|
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.
|
// 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
|
// Persisted task rows lag a fire-and-forget lane start, so omitting these
|
||||||
// reservations lets a second coordinator pass over-admit one project.
|
// 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
|
// Older test/runtime semaphore wrappers predate tryAcquire. They still
|
||||||
// exercise project admission, while production semaphores atomically take
|
// exercise project admission, while production semaphores atomically take
|
||||||
// the host slot here.
|
// the host slot here.
|
||||||
@@ -153,10 +219,10 @@ export class ProjectAdmissionCoordinator {
|
|||||||
// The host semaphore is exhausted for everyone, not just this candidate:
|
// The host semaphore is exhausted for everyone, not just this candidate:
|
||||||
// trying younger candidates cannot succeed, so stop rather than spin.
|
// trying younger candidates cannot succeed, so stop rather than spin.
|
||||||
if (!acquiredHostSlot) return;
|
if (!acquiredHostSlot) return;
|
||||||
// Compatibility-only semaphore shims cannot hold a reservation. Their
|
// The project reservation is independent of the optional host semaphore.
|
||||||
// lane tests provide claimed() synchronously, while real host semaphores
|
// It bridges every lane's dispatch-to-persist gap, including runtimes
|
||||||
// use this durable marker until take/drop below.
|
// where the cross-project semaphore is intentionally absent.
|
||||||
if (hasReservableHostSlot) this.reserve(params.projectId, winner.taskId);
|
this.reserve(params.projectId, winner.taskId);
|
||||||
/*
|
/*
|
||||||
FNXC:ConcurrencyAdmission 2026-07-26-10:35:
|
FNXC:ConcurrencyAdmission 2026-07-26-10:35:
|
||||||
Unwind EXACTLY what this attempt took. Two ways a naive `semaphore.release()` corrupts
|
Unwind EXACTLY what this attempt took. Two ways a naive `semaphore.release()` corrupts
|
||||||
@@ -175,20 +241,9 @@ export class ProjectAdmissionCoordinator {
|
|||||||
the explicit unwind.
|
the explicit unwind.
|
||||||
*/
|
*/
|
||||||
const releaseAttempt = () => {
|
const releaseAttempt = () => {
|
||||||
if (!hasReservableHostSlot) return;
|
dropPreHeldExecutorSlot(winner.taskId);
|
||||||
/*
|
this.releaseReservation(winner.taskId);
|
||||||
FNXC:CapacityModel 2026-07-29-13:20: the host-slot release is now
|
if (hasReservableHostSlot) params.semaphore?.release();
|
||||||
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();
|
|
||||||
};
|
};
|
||||||
try {
|
try {
|
||||||
winner.reserve?.();
|
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).
|
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:
|
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 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.
|
* Register a project-admission handoff and, when present, its host semaphore slot.
|
||||||
* Must not be called without a matching prior acquire; pair with take() or drop().
|
|
||||||
*/
|
*/
|
||||||
export function registerPreHeldExecutorSlot(taskId: string): void {
|
export function registerPreHeldExecutorSlot(taskId: string, semaphoreHeld = true): void {
|
||||||
preHeldExecutorSlots.add(taskId);
|
preHeldAdmissionReservations.add(taskId);
|
||||||
|
if (semaphoreHeld) preHeldExecutorSlots.add(taskId);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Transfer ownership of a pre-held executor slot to the caller.
|
* 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.
|
* 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);
|
const taken = preHeldExecutorSlots.delete(taskId);
|
||||||
if (taken) projectAdmissionCoordinator.releaseReservation(taskId);
|
if (!retainAdmissionReservation) releasePreHeldAdmissionReservation(taskId);
|
||||||
return taken;
|
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):
|
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
|
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
|
slot this call never held, INFLATING capacity — the opposite of the leak the
|
||||||
cleanup was guarding against.
|
cleanup was guarding against.
|
||||||
*/
|
*/
|
||||||
if (!preHeldExecutorSlots.delete(taskId)) return false;
|
const heldSemaphore = preHeldExecutorSlots.delete(taskId);
|
||||||
projectAdmissionCoordinator.releaseReservation(taskId);
|
if (preHeldAdmissionReservations.delete(taskId)) projectAdmissionCoordinator.releaseReservation(taskId);
|
||||||
return true;
|
return heldSemaphore;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Test/helper: whether a task currently has an unclaimed pre-held executor slot. */
|
/** 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. */
|
/** Test helper: clear all pre-held registrations without releasing semaphore slots. */
|
||||||
export function clearPreHeldExecutorSlotsForTests(): void {
|
export function clearPreHeldExecutorSlotsForTests(): void {
|
||||||
preHeldExecutorSlots.clear();
|
preHeldExecutorSlots.clear();
|
||||||
|
preHeldAdmissionReservations.clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -361,13 +425,30 @@ export function persistedTopLevelAgentSlots(tasks: Task[]): number {
|
|||||||
* predicate; raw `persistedTopLevelAgentSlots` remains for pre-enriched tests
|
* predicate; raw `persistedTopLevelAgentSlots` remains for pre-enriched tests
|
||||||
* and callers that cannot resolve an IR.
|
* 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 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);
|
const ir = await resolveWorkflowIrForTask(store, task.id, irCache);
|
||||||
return enrichRunningAgentTaskShape(task, ir);
|
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));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -3520,7 +3520,10 @@ export class TaskExecutor {
|
|||||||
*/
|
*/
|
||||||
private async runWithExecutorSemaphore<T>(taskId: string, work: () => Promise<T>): Promise<T> {
|
private async runWithExecutorSemaphore<T>(taskId: string, work: () => Promise<T>): Promise<T> {
|
||||||
const sem = this.options.semaphore;
|
const sem = this.options.semaphore;
|
||||||
if (!sem) return work();
|
if (!sem) {
|
||||||
|
takePreHeldExecutorSlot(taskId);
|
||||||
|
return work();
|
||||||
|
}
|
||||||
if (this.outerConcurrencyClaims.has(taskId)) {
|
if (this.outerConcurrencyClaims.has(taskId)) {
|
||||||
return work();
|
return work();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -75,7 +75,11 @@ import type { RoutineRunner } from "./routine-runner.js";
|
|||||||
import { sweepStaleAutostashes, VerificationError } from "./merger.js";
|
import { sweepStaleAutostashes, VerificationError } from "./merger.js";
|
||||||
import { runAiMerge, landWorkspaceTask, WorkspacePartialLandError, WorkspaceRepoLandBusyError } from "./merger-ai.js";
|
import { runAiMerge, landWorkspaceTask, WorkspacePartialLandError, WorkspaceRepoLandBusyError } from "./merger-ai.js";
|
||||||
import { promoteBranchGroup, type BranchGroupPromotionResult, type CreateGroupPrFn, type SyncGroupPrFn } from "./group-merge-coordinator.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 { canStartNextMergeBody } from "./merge-reclaim-policy.js";
|
||||||
import {
|
import {
|
||||||
registerProjectVerificationLimit,
|
registerProjectVerificationLimit,
|
||||||
@@ -3888,6 +3892,7 @@ export class ProjectEngine {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
let selected = false;
|
let selected = false;
|
||||||
|
const admissionSettings = await store.getSettings();
|
||||||
/*
|
/*
|
||||||
FNXC:ConcurrencyAdmission 2026-08-01-01:50 (ROOT CAUSE — triage admission died during every merge):
|
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 —
|
This lane previously ran `value = await start()` INSIDE its admission `start()` callback —
|
||||||
@@ -3908,7 +3913,11 @@ export class ProjectEngine {
|
|||||||
*/
|
*/
|
||||||
await projectAdmissionCoordinator.admitOldest({
|
await projectAdmissionCoordinator.admitOldest({
|
||||||
projectId: cwd,
|
projectId: cwd,
|
||||||
maxConcurrent: (await store.getSettings()).maxConcurrent ?? 2,
|
maxConcurrent: resolveActiveTaskCapacityLimit({
|
||||||
|
maxConcurrent: admissionSettings.maxConcurrent ?? 2,
|
||||||
|
maxWorktrees: admissionSettings.maxWorktrees ?? 4,
|
||||||
|
worktreeLimitEnabled: admissionSettings.worktreeLimitEnabled,
|
||||||
|
}),
|
||||||
claimed: async () => computeTopLevelConcurrencyClaimedFromStore({
|
claimed: async () => computeTopLevelConcurrencyClaimedFromStore({
|
||||||
store,
|
store,
|
||||||
tasks: await store.listTasks({ slim: true, includeArchived: false }),
|
tasks: await store.listTasks({ slim: true, includeArchived: false }),
|
||||||
|
|||||||
@@ -16,12 +16,13 @@ import { existsSync } from "node:fs";
|
|||||||
import { readFile } from "node:fs/promises";
|
import { readFile } from "node:fs/promises";
|
||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import {
|
import {
|
||||||
computeTopLevelConcurrencyClaimedFromStore,
|
|
||||||
dropPreHeldExecutorSlot,
|
dropPreHeldExecutorSlot,
|
||||||
hasPreHeldExecutorSlot,
|
hasPreHeldExecutorSlot,
|
||||||
projectAdmissionCoordinator,
|
projectAdmissionCoordinator,
|
||||||
|
persistedTopLevelAgentTaskIdsFromStore,
|
||||||
recoverIdleSemaphoreLeakCandidate,
|
recoverIdleSemaphoreLeakCandidate,
|
||||||
registerPreHeldExecutorSlot,
|
registerPreHeldExecutorSlot,
|
||||||
|
resolveActiveTaskCapacityLimit,
|
||||||
type AgentSemaphore,
|
type AgentSemaphore,
|
||||||
} from "./concurrency.js";
|
} from "./concurrency.js";
|
||||||
import { planTaskWorktreePath, resolveTaskWorkingBranch } from "./worktree-names.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 { UnlinkedMissionsAdvisoryReporter } from "./unlinked-missions-advisory-reporter.js";
|
||||||
import { createRunAuditor, generateSyntheticRunId } from "./run-audit.js";
|
import { createRunAuditor, generateSyntheticRunId } from "./run-audit.js";
|
||||||
import type { TaskMoveLanes } from "@fusion/core";
|
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 { ColumnRoleTraitFlags } from "@fusion/core";
|
||||||
import type { WorkflowIr, WorkflowIrV2 } from "@fusion/core";
|
import type { WorkflowIr, WorkflowIrV2 } from "@fusion/core";
|
||||||
import { runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./hold-release.js";
|
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.
|
* semaphore gate uses this instead of only in-progress agentSlots.
|
||||||
*/
|
*/
|
||||||
topLevelClaimedSlots?: number;
|
topLevelClaimedSlots?: number;
|
||||||
/** FNXC:WorkflowScheduling 2026-07-31-23:50: every live worktree holder (wip + planning/review
|
/** FNXC:WorktreeCapacity 2026-08-01-04:38: active task ids occupying worktree-capacity slots. */
|
||||||
* lanes), so the maxWorktrees holders diagnostic names who actually occupies the slots. */
|
|
||||||
worktreeHolderTaskIds?: string[];
|
worktreeHolderTaskIds?: string[];
|
||||||
/** U6: additive per-column capacity gates (flag-ON only). Omitted → the legacy
|
/** U6: additive per-column capacity gates (flag-ON only). Omitted → the legacy
|
||||||
* three-gate report is byte-identical. */
|
* three-gate report is byte-identical. */
|
||||||
@@ -890,77 +890,7 @@ export interface SchedulerOptions {
|
|||||||
* itself is also refreshed: if `pollIntervalMs` differs from the active
|
* itself is also refreshed: if `pollIntervalMs` differs from the active
|
||||||
* timer, the `setInterval` is transparently restarted.
|
* timer, the `setInterval` is transparently restarted.
|
||||||
*/
|
*/
|
||||||
/*
|
function releaseReservedSlot(reservedSlots: number): number {
|
||||||
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 {
|
|
||||||
return Math.max(0, reservedSlots - 1);
|
return Math.max(0, reservedSlots - 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1049,7 +979,7 @@ export class Scheduler {
|
|||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
projectId,
|
projectId,
|
||||||
createdAt: task.createdAt,
|
createdAt: task.createdAt,
|
||||||
reserve: () => { if (this.options.semaphore) registerPreHeldExecutorSlot(task.id); },
|
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
|
||||||
start: async () => {
|
start: async () => {
|
||||||
this.coordinatorReadyTasks.delete(task.id);
|
this.coordinatorReadyTasks.delete(task.id);
|
||||||
this.coordinatorAdmittedTaskIds.add(task.id);
|
this.coordinatorAdmittedTaskIds.add(task.id);
|
||||||
@@ -2223,6 +2153,11 @@ export class Scheduler {
|
|||||||
worktreeLimitEnabled: settings.worktreeLimitEnabled,
|
worktreeLimitEnabled: settings.worktreeLimitEnabled,
|
||||||
});
|
});
|
||||||
const maxConcurrent = settings.maxConcurrent ?? this.options.maxConcurrent ?? 2;
|
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):
|
FNXC:WorkflowScheduling 2026-07-19-02:35 (U4/KTD-9):
|
||||||
Count active WIP reservations by the `wip` trait, not the literal
|
Count active WIP reservations by the `wip` trait, not the literal
|
||||||
@@ -2297,65 +2232,16 @@ export class Scheduler {
|
|||||||
isReviewColumnRole(columnFlagsForTask(task), task.column);
|
isReviewColumnRole(columnFlagsForTask(task), task.column);
|
||||||
const wipTaskIds = tasks.filter(isWipColumnTask).map((task) => task.id);
|
const wipTaskIds = tasks.filter(isWipColumnTask).map((task) => task.id);
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowScheduling 2026-07-31-23:50 (maxWorktrees counted only WIP — live board breach):
|
FNXC:WorktreeCapacity 2026-08-01-04:38 (inactive retained-worktree capacity inversion):
|
||||||
Under plan-in-place EVERY lane's live card holds a real worktree — planning runs in the task
|
Worktree capacity is a LIVE-TASK budget, not a count of directories retained on disk. The
|
||||||
worktree (triage.ts) and review/merge keeps it — but this ledger counted WIP cards only. The
|
dashboard showed seven active tasks against maxWorktrees=9, but two dependency-blocked queued
|
||||||
protection that used to catch the difference was the GLOBAL SEMAPHORE gate, whose FNXC below
|
cards retained worktree paths. Counting those inactive paths filled the ledger and prevented
|
||||||
says exactly this ("must include every live top-level agent holder (planning triage and active
|
the dependency-free roots from starting. Use the canonical enriched running-agent predicate
|
||||||
in-review), otherwise the hold/release sweep can admit an executor on top of a full planner
|
shared with the board; planning, WIP, and active review count, while queued/paused/terminal
|
||||||
fleet"); the two-number capacity model deleted the semaphore, and this gate never learned to
|
tasks do not. Same-sweep reservations below keep newly released tasks visible immediately.
|
||||||
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.
|
|
||||||
*/
|
*/
|
||||||
/*
|
const activeWorktreeTaskIds = await persistedTopLevelAgentTaskIdsFromStore(this.store, tasks);
|
||||||
FNXC:WorkflowResolvedColumns 2026-07-31-20:55 (u12 — the ratchet caught this, correctly):
|
let reservedWorktreeSlots = activeWorktreeTaskIds.length;
|
||||||
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;
|
|
||||||
let reservedConcurrentSlots = wipTaskIds.length;
|
let reservedConcurrentSlots = wipTaskIds.length;
|
||||||
const inProgressTaskIds = wipTaskIds;
|
const inProgressTaskIds = wipTaskIds;
|
||||||
const dispatchPrepByTaskId = new Map<string, {
|
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({
|
const concurrencyDiagnostic = computeConcurrencyGateDiagnostic({
|
||||||
agentSlots: reservedConcurrentSlots,
|
agentSlots: reservedConcurrentSlots,
|
||||||
maxConcurrent,
|
maxConcurrent,
|
||||||
activeWorktrees: reservedWorktreeSlots,
|
activeWorktrees: reservedWorktreeSlots,
|
||||||
maxWorktrees: resolveCandidateWorktreeCapacityLimit(maxWorktrees, candidateHoldsWorktree),
|
maxWorktrees,
|
||||||
worktreeHolderTaskIds: [...inProgressTaskIds, ...nonWipWorktreeHolderIds],
|
worktreeHolderTaskIds: [...activeWorktreeTaskIds, ...dispatchPrepByTaskId.keys()],
|
||||||
semaphore: this.options.semaphore,
|
semaphore: this.options.semaphore,
|
||||||
inProgressTaskIds,
|
inProgressTaskIds,
|
||||||
topLevelClaimedSlots,
|
topLevelClaimedSlots: reservedWorktreeSlots,
|
||||||
});
|
});
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowScheduling 2026-06-23-20:58:
|
FNXC:WorkflowScheduling 2026-06-23-20:58:
|
||||||
@@ -3003,9 +2876,42 @@ export class Scheduler {
|
|||||||
return null;
|
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 sem = this.options.semaphore;
|
||||||
const coordinatorReserved = hasPreHeldExecutorSlot(task.id);
|
const hostSlotReserved = hasPreHeldExecutorSlot(task.id);
|
||||||
if (sem && !coordinatorReserved && !sem.tryAcquire()) {
|
registerPreHeldExecutorSlot(task.id, hostSlotReserved);
|
||||||
|
if (sem && !hostSlotReserved && !sem.tryAcquire()) {
|
||||||
|
dropPreHeldExecutorSlot(task.id);
|
||||||
if (reservedScope) {
|
if (reservedScope) {
|
||||||
activeScopes.delete(task.id);
|
activeScopes.delete(task.id);
|
||||||
activeScopeColumns.delete(task.id);
|
activeScopeColumns.delete(task.id);
|
||||||
@@ -3019,7 +2925,7 @@ export class Scheduler {
|
|||||||
await this.logDispatchQueuedReason(task.id, reason, formatConcurrencyLimitMemoKey(concurrencyDiagnostic));
|
await this.logDispatchQueuedReason(task.id, reason, formatConcurrencyLimitMemoKey(concurrencyDiagnostic));
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
if (sem && !coordinatorReserved) {
|
if (sem && !hostSlotReserved) {
|
||||||
registerPreHeldExecutorSlot(task.id);
|
registerPreHeldExecutorSlot(task.id);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3063,8 +2969,7 @@ export class Scheduler {
|
|||||||
task: freshTask,
|
task: freshTask,
|
||||||
});
|
});
|
||||||
|
|
||||||
// Transfer, not addition, for a candidate that already holds its worktree.
|
reservedWorktreeSlots += 1;
|
||||||
reservedWorktreeSlots = reserveWorktreeOnDispatch(reservedWorktreeSlots, candidateHoldsWorktree);
|
|
||||||
reservedConcurrentSlots += 1;
|
reservedConcurrentSlots += 1;
|
||||||
let released = false;
|
let released = false;
|
||||||
return {
|
return {
|
||||||
@@ -3075,7 +2980,7 @@ export class Scheduler {
|
|||||||
activeScopes.delete(task.id);
|
activeScopes.delete(task.id);
|
||||||
activeScopeColumns.delete(task.id);
|
activeScopeColumns.delete(task.id);
|
||||||
}
|
}
|
||||||
reservedWorktreeSlots = releaseWorktreeReservation(reservedWorktreeSlots, candidateHoldsWorktree);
|
reservedWorktreeSlots = releaseReservedSlot(reservedWorktreeSlots);
|
||||||
reservedConcurrentSlots = releaseReservedSlot(reservedConcurrentSlots);
|
reservedConcurrentSlots = releaseReservedSlot(reservedConcurrentSlots);
|
||||||
dispatchPrepByTaskId.delete(task.id);
|
dispatchPrepByTaskId.delete(task.id);
|
||||||
if (dropPreHeldExecutorSlot(task.id)) sem?.release();
|
if (dropPreHeldExecutorSlot(task.id)) sem?.release();
|
||||||
@@ -3089,6 +2994,9 @@ export class Scheduler {
|
|||||||
this.planWorktreePath(task, settings.worktreeNaming, reservedNames, settings),
|
this.planWorktreePath(task, settings.worktreeNaming, reservedNames, settings),
|
||||||
});
|
});
|
||||||
for (const taskId of result.released) {
|
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);
|
const prep = dispatchPrepByTaskId.get(taskId);
|
||||||
if (!prep) continue;
|
if (!prep) continue;
|
||||||
/*
|
/*
|
||||||
|
|||||||
@@ -37,6 +37,7 @@ import {
|
|||||||
resolveLifecycleColumns,
|
resolveLifecycleColumns,
|
||||||
resolveWorkflowIrForTaskWithProvenance,
|
resolveWorkflowIrForTaskWithProvenance,
|
||||||
resolveProjectColumnsForRoles,
|
resolveProjectColumnsForRoles,
|
||||||
|
resolveWorktreeCapacityLimit,
|
||||||
workflowHasColumn,
|
workflowHasColumn,
|
||||||
getStepParser,
|
getStepParser,
|
||||||
computePlanApprovalFingerprint,
|
computePlanApprovalFingerprint,
|
||||||
@@ -137,8 +138,11 @@ import {
|
|||||||
PRIORITY_SPECIFY,
|
PRIORITY_SPECIFY,
|
||||||
computeTopLevelConcurrencyClaimedFromStore,
|
computeTopLevelConcurrencyClaimedFromStore,
|
||||||
dropPreHeldExecutorSlot,
|
dropPreHeldExecutorSlot,
|
||||||
|
persistedTopLevelAgentTaskIdsFromStore,
|
||||||
projectAdmissionCoordinator,
|
projectAdmissionCoordinator,
|
||||||
registerPreHeldExecutorSlot,
|
registerPreHeldExecutorSlot,
|
||||||
|
releasePreHeldAdmissionReservation,
|
||||||
|
resolveActiveTaskCapacityLimit,
|
||||||
takePreHeldExecutorSlot,
|
takePreHeldExecutorSlot,
|
||||||
recoverIdleSemaphoreLeakCandidate,
|
recoverIdleSemaphoreLeakCandidate,
|
||||||
type AgentSemaphore,
|
type AgentSemaphore,
|
||||||
@@ -581,7 +585,7 @@ export class TriageProcessor {
|
|||||||
);
|
);
|
||||||
return tasks.filter((task) => !this.coordinatorAdmittedTaskIds.has(task.id)).map((task) => ({
|
return tasks.filter((task) => !this.coordinatorAdmittedTaskIds.has(task.id)).map((task) => ({
|
||||||
taskId: task.id, projectId: this.rootDir, createdAt: task.createdAt,
|
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 () => {
|
start: async () => {
|
||||||
this.coordinatorAdmittedTaskIds.add(task.id);
|
this.coordinatorAdmittedTaskIds.add(task.id);
|
||||||
void this.specifyTask(task);
|
void this.specifyTask(task);
|
||||||
@@ -2128,43 +2132,24 @@ export class TriageProcessor {
|
|||||||
*/
|
*/
|
||||||
const projectRoom = Math.max(0, maxConcurrent - claimed);
|
const projectRoom = Math.max(0, maxConcurrent - claimed);
|
||||||
/*
|
/*
|
||||||
FNXC:CapacityModel 2026-08-01-02:15 (live breach — 8 planners on a maxWorktrees=4 board):
|
FNXC:WorktreeCapacity 2026-08-01-04:38:
|
||||||
Planning admission gated ONLY on the agent count. Every planner acquires a REAL worktree
|
Worktree slots follow the canonical LIVE-TASK count (`claimed`), not retained directory
|
||||||
(plan-in-place), so admission must also respect the worktree budget — the same two-number
|
metadata. A queued, paused, dependency-blocked, or terminal card may keep a worktree path for
|
||||||
model the scheduler's dispatch gate enforces. This gap was invisible for days because the
|
reuse/cleanup without consuming admission capacity. Every newly admitted planner becomes live
|
||||||
merge-inside-admission-drain bug (00769fad7c) froze admission during every merge and
|
and spends one slot below, even when it reuses an existing directory.
|
||||||
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.
|
|
||||||
*/
|
*/
|
||||||
const maxWorktrees = (settings as { maxWorktrees?: number | null }).maxWorktrees ?? 4;
|
const maxWorktrees = resolveWorktreeCapacityLimit({
|
||||||
let worktreeRoom = Number.POSITIVE_INFINITY;
|
maxWorktrees: settings.maxWorktrees ?? 4,
|
||||||
if (typeof maxWorktrees === "number" && Number.isFinite(maxWorktrees)) {
|
worktreeLimitEnabled: settings.worktreeLimitEnabled,
|
||||||
/*
|
});
|
||||||
FNXC:WorkflowResolvedColumns 2026-08-01-03:05:
|
const activeTaskLimit = resolveActiveTaskCapacityLimit({
|
||||||
TERMINAL IS A ROLE HERE, NOT A NAME. The ledger above excludes terminal lanes because their
|
maxConcurrent,
|
||||||
retained worktrees are cleanup-owned rather than capacity; against the literals `done` and
|
maxWorktrees: settings.maxWorktrees ?? 4,
|
||||||
`archived` a RENAMED board matched neither, so every finished card still counted as a live
|
worktreeLimitEnabled: settings.worktreeLimitEnabled,
|
||||||
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
|
const worktreeRoom = maxWorktrees === null
|
||||||
strictly worse, because a stall is silent where a breach is at least visible as 8 planners.
|
? Number.POSITIVE_INFINITY
|
||||||
|
: Math.max(0, maxWorktrees - claimed);
|
||||||
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 maxToStart = Math.min(projectRoom, worktreeRoom);
|
const maxToStart = Math.min(projectRoom, worktreeRoom);
|
||||||
|
|
||||||
if (maxToStart <= 0 && triageTasks.length > 0) {
|
if (maxToStart <= 0 && triageTasks.length > 0) {
|
||||||
@@ -2276,36 +2261,33 @@ export class TriageProcessor {
|
|||||||
// the planner's synchronous processing claim until after this poll returns.
|
// the planner's synchronous processing claim until after this poll returns.
|
||||||
const admittedThisPoll = new Set<string>();
|
const admittedThisPoll = new Set<string>();
|
||||||
/*
|
/*
|
||||||
FNXC:CapacityModel 2026-08-01-02:20: the transfer rule from the scheduler ledger, applied at
|
FNXC:WorktreeCapacity 2026-08-01-04:38:
|
||||||
admission — a candidate that already HOLDS a worktree (replan re-entry) reuses it, so it
|
Each admitted planner becomes an active task and therefore spends one worktree-capacity slot.
|
||||||
spends only agent room, never a fresh worktree slot. Budgets decrement per admission below.
|
Whether its directory is newly created or retained is deliberately irrelevant.
|
||||||
*/
|
*/
|
||||||
let agentBudget = projectRoom;
|
let agentBudget = projectRoom;
|
||||||
let worktreeBudget = worktreeRoom;
|
let worktreeBudget = worktreeRoom;
|
||||||
for (let i = 0; i < triageTasks.length; i++) {
|
for (let i = 0; i < triageTasks.length; i++) {
|
||||||
if (agentBudget <= 0) break;
|
if (agentBudget <= 0 || worktreeBudget <= 0) break;
|
||||||
const candidate = triageTasks[i];
|
|
||||||
const candidateHoldsWorktree = typeof candidate.worktree === "string" && candidate.worktree.length > 0;
|
|
||||||
if (!candidateHoldsWorktree && worktreeBudget <= 0) continue;
|
|
||||||
agentBudget -= 1;
|
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({
|
await projectAdmissionCoordinator.admitOldest({
|
||||||
// rootDir is the stable per-project identity held by this processor.
|
// rootDir is the stable per-project identity held by this processor.
|
||||||
projectId: this.rootDir,
|
projectId: this.rootDir,
|
||||||
maxConcurrent,
|
maxConcurrent: activeTaskLimit,
|
||||||
claimed: async () => {
|
claimed: async () => (await getFreshClaimSnapshot()).count,
|
||||||
const fresh = await this.store.listTasks({ slim: true, includeArchived: false });
|
claimedTaskIds: async () => (await getFreshClaimSnapshot()).ids,
|
||||||
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,
|
|
||||||
});
|
|
||||||
},
|
|
||||||
semaphore: this.options.semaphore,
|
semaphore: this.options.semaphore,
|
||||||
refresh: async () => triageTasks
|
refresh: async () => triageTasks
|
||||||
.filter((task) => !admittedThisPoll.has(task.id) && !this.coordinatorAdmittedTaskIds.has(task.id) && !this.processing.has(task.id) && !this.hasLivePlanningWork(task.id))
|
.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
|
// FNXC:ConcurrencyAdmission 2026-08-05-10:00: the planner must
|
||||||
// own the coordinator's real host reservation before it starts;
|
// own the coordinator's real host reservation before it starts;
|
||||||
// deferring to semaphore.run would reintroduce priority overtaking.
|
// 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 () => {
|
start: async () => {
|
||||||
admittedThisPoll.add(task.id);
|
admittedThisPoll.add(task.id);
|
||||||
this.coordinatorAdmittedTaskIds.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
|
updatePlanningStateIfStillCurrent); this one must be visible — the FN-7977 steps>0
|
||||||
wedge stalled the whole planner for hours precisely because it logged nothing.
|
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(
|
planLog.warn(
|
||||||
`${task.id}: planning claim skipped — live row is no longer in the planning stage; `
|
`${task.id}: planning claim skipped — live row is no longer in the planning stage; `
|
||||||
+ "it will be re-claimed on the next poll",
|
+ "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
|
// Coordinator already owns this top-level slot; run directly so it
|
||||||
// cannot join the priority queue after age-based admission.
|
// cannot join the priority queue after age-based admission.
|
||||||
try {
|
try {
|
||||||
|
|||||||
Reference in New Issue
Block a user