diff --git a/.changeset/steady-active-capacity.md b/.changeset/steady-active-capacity.md
new file mode 100644
index 0000000000..e668e10c72
--- /dev/null
+++ b/.changeset/steady-active-capacity.md
@@ -0,0 +1,7 @@
+---
+"@runfusion/fusion": patch
+---
+
+summary: Keep active task admission within worktree capacity and show queued cards as status badges.
+category: fix
+dev: Refreshes full live-task claims at serialized admission and defers capacity-blocked merge retries.
diff --git a/packages/dashboard/app/components/TaskCard.css b/packages/dashboard/app/components/TaskCard.css
index fc7be8ecfc..6fc6a01834 100644
--- a/packages/dashboard/app/components/TaskCard.css
+++ b/packages/dashboard/app/components/TaskCard.css
@@ -259,6 +259,11 @@ Lock the id to the same chip-height row as .card-header-actions so the task id,
letter-spacing: 0.5px;
}
+.card-queued-reason-icon {
+ flex: 0 0 auto;
+ margin-inline-start: calc(var(--space-xs) / 2);
+}
+
/*
FNXC:WorkflowBoard 2026-06-29-00:00:
Aggregate Board cards need a compact workflow-name badge so operators can identify each card's source workflow while shared columns combine multiple workflows. The badge is opt-in metadata from Board and uses neutral card tokens so non-board TaskCard callers do not render empty workflow shells.
diff --git a/packages/dashboard/app/components/TaskCard.tsx b/packages/dashboard/app/components/TaskCard.tsx
index 4cfccd58de..53cf572435 100644
--- a/packages/dashboard/app/components/TaskCard.tsx
+++ b/packages/dashboard/app/components/TaskCard.tsx
@@ -1581,7 +1581,7 @@ function TaskCardComponent({
&& !visualStatus
&& !planReviewRunning
&& !isAgentActive;
- const showReadyBadge = showIdleTodoBadge && !awaitingPlanning;
+ const showReadyBadge = showIdleTodoBadge && !queued && !awaitingPlanning;
/*
FNXC:CodingIdeasWorkflow 2026-07-25-12:05:
"Queued to plan" is the exact complement of Ready: same idle-in-Todo conditions, but the card has
@@ -1598,7 +1598,7 @@ function TaskCardComponent({
Both badges now derive from the single `awaitingPlanning` value above, so "exact complement" is
structural rather than a property of the step count that two independent conditions had to agree on.
*/
- const showQueuedToPlanBadge = showIdleTodoBadge && awaitingPlanning;
+ const showQueuedToPlanBadge = showIdleTodoBadge && !queued && awaitingPlanning;
// Native HTML5 drag is desktop-mouse only — it doesn't move cards via touch.
// On touch-primary devices the `draggable` attribute still arms the browser's
// touch-drag heuristic, which intermittently hijacks horizontal swipes meant
@@ -2130,8 +2130,6 @@ function TaskCardComponent({
const showAddressPrFeedbackAction = canStartPrFeedbackAddressing(task, taskColumnFlags);
const metaRowVisible =
(task.dependencies?.length ?? 0) > 0
- || queued
- || task.status === "queued"
|| Boolean(task.blockedBy)
|| Boolean(task.overlapBlockedBy)
|| Boolean(fanout && fanout.totalCount > 0);
@@ -3370,6 +3368,16 @@ function TaskCardComponent({
&& !visualStatus
&& Boolean(task.recentAgentActivityAt)
&& isAgentActive;
+ /*
+ FNXC:TaskStatusBadge 2026-08-01-07:20 (operator: queued belongs with Planning and Ready):
+ Queued used to render as a clock-and-text footer tag, separating the waiting state from the
+ Planning and Ready badges operators compare it with. Treat every non-WIP queued card as a normal
+ header status badge instead. The shared badge geometry and column color carry desktop/mobile
+ behavior; no standalone queued visual remains at the bottom of the card.
+ */
+ const showQueuedBadge = !isPaused
+ && !isWipColumn
+ && (queued || visualStatus === "queued");
const showStatusBadge = !isPaused
&& (hasTaskStatusBadge(visualStatus) || isTransientPlannerActive)
&& visualStatus !== "queued";
@@ -3400,6 +3408,8 @@ function TaskCardComponent({
*/
: showQueuedToPlanBadge
? t("tasks.queuedToPlan", "Queued to plan")
+ : showQueuedBadge
+ ? t("tasks.statusQueued", "Queued")
: getTaskStatusLabel(visualStatus ?? "", t, showOptionalGateBadge ? undefined : getRunningWorkflowStepLabel(task), { idle: !isAgentActive, overlapBlockedBy: task.overlapBlockedBy ?? null });
const hasCardMetaBadges = showPriorityBadge
|| task.executionMode === "fast"
@@ -3411,6 +3421,7 @@ function TaskCardComponent({
|| showStatusBadge
|| showOptionalGateBadge
|| showReadyBadge
+ || showQueuedBadge
// FNXC:CodingIdeasWorkflow 2026-07-25-12:05: the header wrapper only renders when it has a
// real child, so a new badge must be declared here or it never mounts (Queued to plan is the
// only badge on an unplanned idle Todo card — without this the whole cluster stays absent).
@@ -3544,7 +3555,7 @@ function TaskCardComponent({
: pausedByAgent ? t("tasks.pausedByAgent", "paused by agent") : t("tasks.paused", "paused")}
)}
- {(showStatusBadge || showQueuedToPlanBadge) && (
+ {(showStatusBadge || showQueuedToPlanBadge || showQueuedBadge) && (
{statusBadgeLabel}
+ {showQueuedBadge && task.overlapBlockedBy && (
+
+ )}
+ {showQueuedBadge && !task.overlapBlockedBy && task.blockedBy && (
+
+ )}
)}
{showOptionalGateBadge && optionalGateBadge && (
@@ -4183,7 +4204,6 @@ function TaskCardComponent({
)}
- {(queued || task.status === "queued") && !isWipColumn && {t("tasks.queued", "Queued")}}
{placeFooterRightInMeta && footerRightCluster}
)}
diff --git a/packages/dashboard/app/components/__tests__/TaskCard.test.tsx b/packages/dashboard/app/components/__tests__/TaskCard.test.tsx
index 3859d21347..fd499e80c6 100644
--- a/packages/dashboard/app/components/__tests__/TaskCard.test.tsx
+++ b/packages/dashboard/app/components/__tests__/TaskCard.test.tsx
@@ -26,12 +26,12 @@ import { getPriorityColorVar, getPriorityLabel } from "../../utils/priorityIndic
// Mock lucide-react to avoid SVG rendering issues in test env
vi.mock("lucide-react", () => ({
- Link: () => null,
+ Link: (props: React.SVGProps) => ,
GitBranch: () => null,
Gitlab: () => null,
Clock: () => null,
Pencil: () => null,
- Layers: () => null,
+ Layers: (props: React.SVGProps) => ,
ChevronDown: () => null,
Folder: () => null,
GitPullRequest: () => null,
@@ -4847,12 +4847,12 @@ describe("TaskCard", () => {
expect(screen.getByTestId("provider-icon-github")).toBeDefined();
});
- it("renders the GitHub tracking link inline with queued metadata when the footer has no leading content", () => {
+ it("renders queued as a header status badge and leaves no clock tag at the bottom", () => {
const { container } = render(
{
},
},
})}
+ queued
onOpenDetail={noop}
addToast={noop}
/>,
);
const link = screen.getByRole("link", { name: "Linked GitHub issue #42" });
- const metaRow = container.querySelector(".card-meta");
- const queuedBadge = container.querySelector(".queued-badge");
- expect(container.querySelector(".card-footer-row")).toBeNull();
- expect(link.closest(".card-meta")).toBe(metaRow);
- expect(link.closest(".card-footer-row-right")?.closest(".card-meta")).toBe(metaRow);
- expect(container.querySelector(".card-bottom-right-row")).toBeNull();
- expect(queuedBadge).not.toBeNull();
- expect(queuedBadge?.compareDocumentPosition(link) & Node.DOCUMENT_POSITION_FOLLOWING).toBeTruthy();
+ const queuedBadge = screen.getByText("Queued");
+ expect(link.closest(".card-footer-row")).not.toBeNull();
+ expect(queuedBadge).toHaveClass("card-status-badge", "card-status-badge--todo");
+ expect(queuedBadge.closest(".card-header-badges")).not.toBeNull();
+ expect(container.querySelector(".queued-badge")).toBeNull();
+ expect(queuedBadge.querySelector("svg")).toBeNull();
+ expect(link.closest(".card-meta")).toBeNull();
+ });
+
+ it.each([
+ ["file overlap", { overlapBlockedBy: "FN-OVERLAP" }, "card-queued-overlap-icon", "card-queued-dependency-icon", "Queued due to file overlap with FN-OVERLAP"],
+ ["dependency", { blockedBy: "FN-DEPENDENCY" }, "card-queued-dependency-icon", "card-queued-overlap-icon", "Queued on dependency FN-DEPENDENCY"],
+ ["file overlap when both blockers are present", { overlapBlockedBy: "FN-OVERLAP", blockedBy: "FN-DEPENDENCY" }, "card-queued-overlap-icon", "card-queued-dependency-icon", "Queued due to file overlap with FN-OVERLAP"],
+ ] as const)("shows the %s icon after Queued without putting the blocker id in the badge", (_case, blocker, expectedIcon, absentIcon, title) => {
+ const queuedTask = makeTask({ column: "todo", status: "queued", ...blocker });
+ const { container } = render(
+ ,
+ );
+
+ const badge = screen.getByText("Queued").closest(".card-status-badge") as HTMLElement;
+ expect(badge).toHaveTextContent(/^Queued$/);
+ const icon = badge.querySelector(`[data-testid="${expectedIcon}-${queuedTask.id}"]`);
+ expect(icon).not.toBeNull();
+ expect(icon).toHaveClass("card-queued-reason-icon");
+ expect(icon).toHaveAttribute("size", "8");
+ expect(badge.querySelector(`[data-testid="${absentIcon}-${queuedTask.id}"]`)).toBeNull();
+ expect(badge).toHaveAttribute("title", title);
+ expect(container.querySelector(".queued-badge")).toBeNull();
});
diff --git a/packages/dashboard/app/styles.css b/packages/dashboard/app/styles.css
index e34010fe2a..d584e44210 100644
--- a/packages/dashboard/app/styles.css
+++ b/packages/dashboard/app/styles.css
@@ -3020,16 +3020,6 @@ Toast text must contrast its status background across every dashboard theme and
cursor: default;
}
-.queued-badge {
- display: inline-flex;
- align-items: center;
- gap: 3px;
- font-size: 11px;
- color: var(--todo);
-}
-
-
-
/* === Dependency Dropdown (shared: InlineCreateCard, NewTaskModal, TaskDetailModal, TaskForm) === */
.dep-trigger-wrap {
position: relative;
diff --git a/packages/engine/src/__tests__/concurrency.test.ts b/packages/engine/src/__tests__/concurrency.test.ts
index 061fe8d437..e6effa4325 100644
--- a/packages/engine/src/__tests__/concurrency.test.ts
+++ b/packages/engine/src/__tests__/concurrency.test.ts
@@ -1176,6 +1176,75 @@ describe("ProjectAdmissionCoordinator", () => {
coordinator.releaseReservation("FN-PLANNING");
});
+ it("does not lose a holder that transfers from reservation to durable state during a claim read", async () => {
+ const coordinator = new ProjectAdmissionCoordinator();
+ expect(await coordinator.reserveIfAvailable({
+ projectId: "project-transfer",
+ taskId: "FN-HANDOFF",
+ maxConcurrent: 1,
+ claimed: () => 0,
+ })).toBe(true);
+
+ let finishSnapshot!: () => void;
+ const snapshotBlocked = new Promise((resolve) => { finishSnapshot = resolve; });
+ let snapshotStarted!: () => void;
+ const snapshotDidStart = new Promise((resolve) => { snapshotStarted = resolve; });
+ const candidate = coordinator.reserveIfAvailable({
+ projectId: "project-transfer",
+ taskId: "FN-CANDIDATE",
+ maxConcurrent: 1,
+ claimed: async () => {
+ snapshotStarted();
+ await snapshotBlocked;
+ // This is the pre-persistence snapshot: the handoff is not durable in it yet.
+ return 0;
+ },
+ claimedTaskIds: () => [],
+ });
+
+ await snapshotDidStart;
+ // The handoff becomes durable and releases its transient reservation while the stale read is open.
+ coordinator.releaseReservation("FN-HANDOFF");
+ finishSnapshot();
+
+ expect(await candidate).toBe(false);
+ coordinator.releaseReservation("FN-CANDIDATE");
+ });
+
+ it("evaluates a waiting admission claim only after the prior project drain finishes", async () => {
+ const coordinator = new ProjectAdmissionCoordinator();
+ let releaseDrain!: () => void;
+ const drainBlocked = new Promise((resolve) => { releaseDrain = resolve; });
+ let drainStarted!: () => void;
+ const drainDidStart = new Promise((resolve) => { drainStarted = resolve; });
+ const first = coordinator.reserveIfAvailable({
+ projectId: "project-serialized-snapshot",
+ taskId: "FN-BLOCKER",
+ maxConcurrent: 0,
+ claimed: async () => {
+ drainStarted();
+ await drainBlocked;
+ return 0;
+ },
+ });
+ await drainDidStart;
+
+ const freshClaim = vi.fn(() => 1);
+ const second = coordinator.reserveIfAvailable({
+ projectId: "project-serialized-snapshot",
+ taskId: "FN-WAITING",
+ maxConcurrent: 1,
+ claimed: freshClaim,
+ });
+ await Promise.resolve();
+ expect(freshClaim).not.toHaveBeenCalled();
+
+ releaseDrain();
+ expect(await first).toBe(false);
+ expect(await second).toBe(false);
+ expect(freshClaim).toHaveBeenCalledOnce();
+ });
+
it("admits the oldest same-project candidate atomically and partitions projects", async () => {
const coordinator = new ProjectAdmissionCoordinator();
const started: string[] = [];
diff --git a/packages/engine/src/__tests__/project-engine.test.ts b/packages/engine/src/__tests__/project-engine.test.ts
index 05b113f4be..27884ee9ca 100644
--- a/packages/engine/src/__tests__/project-engine.test.ts
+++ b/packages/engine/src/__tests__/project-engine.test.ts
@@ -2494,6 +2494,74 @@ describe("ProjectEngine paused in-review auto-merge behavior", () => {
await engine.stop();
});
+ it("does not admit a merge over a pending optional workflow-step lease", async () => {
+ const mockStore = createMockStore({ ...baseSettings, autoMerge: true, maxConcurrent: 1, maxWorktrees: 1 });
+ mockStore.store.getTask.mockResolvedValue({
+ id: "FN-MERGE-WAITING",
+ column: "in-review",
+ paused: false,
+ mergeRetries: 0,
+ status: null,
+ branch: "fusion/fn-merge-waiting",
+ createdAt: "2026-01-01T00:00:00.000Z",
+ });
+ mocks.currentStore = mockStore.store;
+
+ const engine = createEngine();
+ await engine.start();
+ mockStore.store.listTasks.mockImplementation(async (options?: { slim?: boolean }) => options?.slim === false
+ ? [{
+ id: "FN-LIVE-REVIEW",
+ column: "todo",
+ paused: false,
+ status: null,
+ workflowStepResults: [{
+ workflowStepId: "code-review",
+ workflowStepName: "Code Review",
+ phase: "pre-merge",
+ source: "optional-group",
+ status: "pending",
+ startedAt: "2026-08-01T00:00:00.000Z",
+ }],
+ }]
+ : []);
+
+ engine.enqueueMerge("FN-MERGE-WAITING");
+
+ await vi.waitFor(() => {
+ expect(mockStore.store.listTasks).toHaveBeenCalledWith({ slim: false, includeArchived: false });
+ });
+ expect(mocks.runAiMerge).not.toHaveBeenCalled();
+ expect(mockStore.store.listTasks.mock.calls.filter(([options]) => options?.slim === false)).toHaveLength(1);
+ const privateEngine = engine as unknown as {
+ mergeQueue: string[];
+ capacityDeferredMergeTaskIds: Set;
+ };
+ expect(privateEngine.mergeQueue).not.toContain("FN-MERGE-WAITING");
+ expect(privateEngine.capacityDeferredMergeTaskIds.has("FN-MERGE-WAITING")).toBe(true);
+ expect(engine.isMergePending("FN-MERGE-WAITING")).toBe(true);
+
+ // An unrelated queue wake must not make the deferred task runnable before its timer.
+ mockStore.store.getTask.mockResolvedValue({
+ id: "FN-OTHER-MERGE",
+ column: "in-review",
+ paused: false,
+ mergeRetries: 0,
+ status: null,
+ branch: "fusion/fn-other-merge",
+ createdAt: "2026-01-02T00:00:00.000Z",
+ });
+ engine.enqueueMerge("FN-OTHER-MERGE");
+ await vi.waitFor(() => {
+ expect(privateEngine.capacityDeferredMergeTaskIds.has("FN-OTHER-MERGE")).toBe(true);
+ });
+ expect(privateEngine.mergeQueue).not.toContain("FN-MERGE-WAITING");
+ expect(mocks.runAiMerge).not.toHaveBeenCalled();
+ await engine.stop();
+ expect(privateEngine.capacityDeferredMergeTaskIds.size).toBe(0);
+ expect(engine.isMergePending("FN-MERGE-WAITING")).toBe(false);
+ });
+
it("records an audit event (not silent) when auto-promotion of a branch-group member fails (Fix #4)", async () => {
const mockStore = createMockStore({ ...baseSettings, autoMerge: true });
// The dequeued + merged task is a shared branch-group member, so the engine
diff --git a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts
index dddce4448e..b8cf695269 100644
--- a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts
+++ b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts
@@ -625,6 +625,88 @@ describe("Scheduler workflow cutover", () => {
expect(second.status).toBe("queued");
});
+ it("rechecks canonical live tasks inside final admission after the sweep snapshot goes stale", async () => {
+ const active = Array.from({ length: 8 }, (_, index) =>
+ task({ id: `FN-ACTIVE-${index}`, column: "in-progress" }),
+ );
+ const ready = task({ id: "FN-READY", status: "queued" });
+ const latePlanner = task({ id: "FN-LATE-PLANNER", status: "planning" });
+ const store = storeWith([...active, ready], { maxConcurrent: 12, maxWorktrees: 9 });
+ vi.mocked(store.listTasks)
+ .mockResolvedValueOnce([...active, ready])
+ .mockResolvedValue([...active, latePlanner, ready]);
+ const onSchedule = vi.fn();
+ const scheduler = new Scheduler(store, { onSchedule });
+ (scheduler as unknown as { running: boolean }).running = true;
+
+ await scheduler.schedule();
+
+ expect(store.listTasks).toHaveBeenCalledWith({ slim: false, includeArchived: false });
+ expect(store.moveTaskIf).not.toHaveBeenCalledWith(
+ ready.id,
+ "in-progress",
+ expect.any(Function),
+ expect.anything(),
+ );
+ expect(onSchedule).not.toHaveBeenCalled();
+ expect(ready.column).toBe("todo");
+ expect(ready.status).toBe("queued");
+ });
+
+ it("counts a pending optional workflow-step lease in final scheduler admission", async () => {
+ const active = Array.from({ length: 8 }, (_, index) =>
+ task({ id: `FN-LEASE-ACTIVE-${index}`, column: "in-progress" }),
+ );
+ const liveLease = task({
+ id: "FN-LIVE-LEASE",
+ status: null,
+ workflowStepResults: [{
+ workflowStepId: "code-review",
+ workflowStepName: "Code Review",
+ phase: "pre-merge",
+ source: "optional-group",
+ status: "pending",
+ startedAt: "2026-08-01T00:00:00.000Z",
+ }],
+ });
+ const ready = task({ id: "FN-READY", status: "queued" });
+ const store = storeWith([...active, liveLease, ready], { 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(store.listTasks).toHaveBeenCalledWith({ slim: false, includeArchived: false });
+ expect(onSchedule).not.toHaveBeenCalled();
+ expect(ready.column).toBe("todo");
+ });
+
+ it("does not let a stale full sweep hide capacity that freed before final admission", async () => {
+ const initiallyActive = Array.from({ length: 9 }, (_, index) =>
+ task({ id: `FN-STALE-ACTIVE-${index}`, column: "in-progress" }),
+ );
+ const ready = task({ id: "FN-READY", status: "queued" });
+ const store = storeWith([...initiallyActive, ready], { maxConcurrent: 12, maxWorktrees: 9 });
+ vi.mocked(store.listTasks)
+ .mockResolvedValueOnce([...initiallyActive, ready])
+ .mockResolvedValue(initiallyActive.slice(0, 8).concat(ready));
+ const onSchedule = vi.fn();
+ const scheduler = new Scheduler(store, { onSchedule });
+ (scheduler as unknown as { running: boolean }).running = true;
+
+ await scheduler.schedule();
+
+ expect(store.listTasks).toHaveBeenCalledWith({ slim: false, includeArchived: false });
+ expect(store.moveTaskIf).toHaveBeenCalledWith(
+ ready.id,
+ "in-progress",
+ expect.any(Function),
+ expect.anything(),
+ );
+ expect(onSchedule).toHaveBeenCalledWith(expect.objectContaining({ id: ready.id }));
+ });
+
it("leaves a task queued when the authoritative release move rejects after reservation", async () => {
const ready = task({ id: "FN-002", status: "queued" });
const store = storeWith([ready], { maxConcurrent: 4, maxWorktrees: 4 });
diff --git a/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts b/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts
index be4a463c92..5a5550d322 100644
--- a/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts
+++ b/packages/engine/src/__tests__/triage-admission-worktree-ledger-renamed-lanes.test.ts
@@ -72,7 +72,9 @@ function createStore(
const found = tasks.find((candidate) => candidate.id === id);
return found ? { ...found, prompt: "", attachments: [], comments: [] } : null;
}),
- listTasks: vi.fn().mockResolvedValue(tasks),
+ listTasks: vi.fn(async (options?: { slim?: boolean }) => options?.slim === true
+ ? tasks.map(({ workflowStepResults: _omitted, ...candidate }) => candidate as Task)
+ : tasks),
getSettings: vi.fn().mockResolvedValue({
maxConcurrent: 20,
maxWorktrees,
@@ -215,4 +217,26 @@ describe("planning admission's worktree ledger on a renamed board", () => {
expect(throttle.length, "a live task must still consume its worktree-capacity slot").toBeGreaterThan(0);
expect(specifyTask).not.toHaveBeenCalled();
});
+
+ it("counts a pending optional workflow-step lease in final planning admission", async () => {
+ const store = createStore([
+ task("FN-WAITING", "drafting"),
+ task("FN-LIVE-REVIEW", "drafting", {
+ steps: [{ name: "Implementation", status: "pending" }] as Task["steps"],
+ workflowStepResults: [{
+ workflowStepId: "code-review",
+ workflowStepName: "Code Review",
+ phase: "pre-merge",
+ source: "optional-group",
+ status: "pending",
+ startedAt: "2026-08-01T00:00:00.000Z",
+ }],
+ }),
+ ], recorded, 1);
+
+ const specifyTask = await pollOnce(store);
+
+ expect(store.listTasks).toHaveBeenCalledWith({ slim: false, includeArchived: false });
+ expect(specifyTask).not.toHaveBeenCalled();
+ });
});
diff --git a/packages/engine/src/concurrency.ts b/packages/engine/src/concurrency.ts
index dd23da2a1a..5615a03d9f 100644
--- a/packages/engine/src/concurrency.ts
+++ b/packages/engine/src/concurrency.ts
@@ -139,11 +139,22 @@ export class ProjectAdmissionCoordinator {
claimed: () => Promise | number;
claimedTaskIds?: () => Promise> | Iterable;
}): Promise {
+ /*
+ FNXC:ConcurrencyAdmission 2026-08-01-07:35:
+ A durable handoff can release its in-memory reservation while an asynchronous task snapshot is
+ being read. The snapshot then predates the durable row while a post-read reservation lookup
+ postdates its release, making one real holder disappear between the two ledgers. Union the
+ reservations observed on both sides of the read so that transfer window is conservatively
+ counted once. The next admission gets a fresh durable snapshot and naturally sheds the old id.
+ */
+ const reservations = new Set(this.reservations.get(params.projectId) ?? []);
const claimed = await params.claimed();
- if (!params.claimedTaskIds) return claimed + this.reservationCount(params.projectId);
+ for (const taskId of this.reservations.get(params.projectId) ?? []) reservations.add(taskId);
+ if (!params.claimedTaskIds) return claimed + reservations.size;
const claimedIds = new Set(await params.claimedTaskIds());
+ for (const taskId of this.reservations.get(params.projectId) ?? []) reservations.add(taskId);
let pendingReservations = 0;
- for (const taskId of this.reservations.get(params.projectId) ?? []) {
+ for (const taskId of reservations) {
if (!claimedIds.has(taskId)) pendingReservations += 1;
}
return claimed + pendingReservations;
diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts
index 82e62d7aa3..53a07f7a45 100644
--- a/packages/engine/src/project-engine.ts
+++ b/packages/engine/src/project-engine.ts
@@ -457,6 +457,14 @@ export class ProjectEngine {
// ── Auto-merge state ──
private mergeQueue: string[] = [];
private mergeActive = new Set();
+ /** Capacity-deferred ids stay out of the runnable queue until their retry timer fires. */
+ private readonly capacityDeferredMergeTaskIds = new Set();
+ private readonly capacityDeferredMerges = new Map;
+ resolvers: MergeResolver[];
+ generation: number;
+ manual: boolean;
+ }>();
/** Merge ids selected by the shared coordinator but not yet handed to rawMerge. */
private readonly coordinatorAdmittedMergeTaskIds = new Set();
private unregisterMergeAdmissionProvider?: () => void;
@@ -805,7 +813,9 @@ export class ProjectEngine {
window, checking it in addition to `mergeQueue` closes that TOCTOU gap.
*/
isMergePending(taskId: string): boolean {
- return this.mergeActive.has(taskId) || this.mergeQueue.includes(taskId);
+ return this.mergeActive.has(taskId)
+ || this.mergeQueue.includes(taskId)
+ || this.capacityDeferredMergeTaskIds.has(taskId);
}
/**
@@ -1316,6 +1326,14 @@ export class ProjectEngine {
clearInterval(this.mergeActiveReconcileTimer);
this.mergeActiveReconcileTimer = null;
}
+ for (const [taskId, deferred] of this.capacityDeferredMerges) {
+ clearTimeout(deferred.timer);
+ for (const resolver of deferred.resolvers) {
+ resolver.reject(new Error(`Engine shutting down — deferred merge for ${taskId} aborted`));
+ }
+ }
+ this.capacityDeferredMerges.clear();
+ this.capacityDeferredMergeTaskIds.clear();
this.stopPlannerOverseerPoll();
/*
@@ -2311,6 +2329,15 @@ export class ProjectEngine {
};
abort = () => {
this.removeMergeResolver(taskId, resolver);
+ const deferred = this.capacityDeferredMerges.get(taskId);
+ if (deferred) {
+ deferred.resolvers = deferred.resolvers.filter((candidate) => candidate !== resolver);
+ if (deferred.manual && deferred.resolvers.length === 0 && !this.hasMergeResolvers(taskId)) {
+ clearTimeout(deferred.timer);
+ this.capacityDeferredMerges.delete(taskId);
+ this.capacityDeferredMergeTaskIds.delete(taskId);
+ }
+ }
if (this.activeMergeTaskId === taskId) {
this.mergeAbortController?.abort();
this.mergeAbortController = null;
@@ -2328,7 +2355,7 @@ export class ProjectEngine {
// If this task is already queued or actively merging, wait for the
// existing merge to finish rather than starting a second one.
- if (this.mergeActive.has(taskId)) return;
+ if (this.mergeActive.has(taskId) || this.capacityDeferredMergeTaskIds.has(taskId)) return;
if (!this.internalEnqueueMerge(taskId)) {
this.removeMergeResolver(taskId, resolver);
@@ -2813,6 +2840,7 @@ export class ProjectEngine {
private internalEnqueueMerge(taskId: string): boolean {
if (this.shuttingDown || !this.started) return false;
+ if (this.capacityDeferredMergeTaskIds.has(taskId)) return false;
if (this.mergeActive.has(taskId)) {
// Distinguish "actually being processed" (queued or active) from a
// leaked entry. Reconcile leaks immediately so recovery paths and fresh
@@ -3895,7 +3923,9 @@ export class ProjectEngine {
const admissionSettings = await store.getSettings();
let mergeClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
const getMergeClaimSnapshot = () => mergeClaimSnapshot ??= (async () => {
- const tasks = await store.listTasks({ slim: true, includeArchived: false });
+ // Full rows preserve pending optional workflow-step leases, which may be the only
+ // live-agent signal for a task while its ordinary status is null.
+ const tasks = await store.listTasks({ slim: false, includeArchived: false });
const ids = await persistedTopLevelAgentTaskIdsFromStore(store, tasks);
return { count: ids.length, ids };
})();
@@ -3943,6 +3973,36 @@ export class ProjectEngine {
projectAdmissionCoordinator.releaseReservation(taskId);
}
};
+ const deferMergeForCapacity = (): void => {
+ const retryMs = settings.pollIntervalMs ?? 15_000;
+ const stashedResolvers = this.takeMergeResolvers(taskId);
+ const generation = this.startupGeneration;
+ this.mergeActive.delete(taskId);
+ this.capacityDeferredMergeTaskIds.add(taskId);
+ const timer = setTimeout(() => {
+ const deferred = this.capacityDeferredMerges.get(taskId);
+ if (!deferred || deferred.timer !== timer) return;
+ this.capacityDeferredMerges.delete(taskId);
+ this.capacityDeferredMergeTaskIds.delete(taskId);
+ if (this.shuttingDown || deferred.generation !== this.startupGeneration) {
+ for (const resolver of deferred.resolvers) resolver.reject(new Error("Engine shutting down"));
+ return;
+ }
+ for (const resolver of deferred.resolvers) this.addMergeResolver(taskId, resolver);
+ if (!this.internalEnqueueMerge(taskId)) {
+ for (const resolver of this.takeMergeResolvers(taskId)) {
+ resolver.reject(new Error(`Deferred merge enqueue rejected for ${taskId}`));
+ }
+ }
+ }, retryMs);
+ timer.unref?.();
+ this.capacityDeferredMerges.set(taskId, {
+ timer,
+ resolvers: stashedResolvers,
+ generation,
+ manual: stashedResolvers.length > 0,
+ });
+ };
if (mergeStrategy === "pull-request" && this.options.processPullRequestMerge && !routeWorkspaceDirect) {
/*
@@ -3966,10 +4026,11 @@ export class ProjectEngine {
});
if (result === undefined) {
// Another older lane won the shared capacity pass. Re-queue rather
- // than treating this deferral as a pull-request merge failure.
- this.mergeActive.delete(taskId);
- this.internalEnqueueMerge(taskId);
- continue;
+ // than treating this deferral as a pull-request merge failure. End
+ // this drain: continuing would dequeue the same item immediately
+ // and spin at full speed while capacity remains unavailable.
+ deferMergeForCapacity();
+ break;
}
if (result === "merged") {
runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR merged: ${taskId}`);
@@ -4127,9 +4188,10 @@ export class ProjectEngine {
if (!result) {
// An older lane won this admission pass. Keep this merge queued;
// treating the deferral as a merge failure would consume retries.
- this.mergeActive.delete(taskId);
- this.internalEnqueueMerge(taskId);
- continue;
+ // Exit this drain so the queued item waits for a normal future wake
+ // instead of retrying the same denied admission in a tight loop.
+ deferMergeForCapacity();
+ break;
}
this.activeMergeSession = null;
diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts
index d1ac94e356..a75f3c3274 100644
--- a/packages/engine/src/scheduler.ts
+++ b/packages/engine/src/scheduler.ts
@@ -2836,22 +2836,11 @@ export class Scheduler {
topLevelClaimedSlots: reservedWorktreeSlots,
});
/*
- FNXC:WorkflowScheduling 2026-06-23-20:58:
- The workflow hold/release sweep is the only todo pickup path, so it must honor the same maxConcurrent, maxWorktrees, and shared semaphore pressure before releasing a task to in-progress.
-
- FNXC:GlobalConcurrencyControls 2026-07-14-18:30:
- Preflight is no longer non-mutating for the shared semaphore: tryAcquire reserves a real slot before the move so triage cannot fill the global cap while this card is already counted as an in-progress runner. On move failure the reservation is released; on success the pre-held slot is transferred to the executor/graph run.
+ The sweep diagnostic is intentionally descriptive only. Worktree preparation can outlive
+ a holder, so using this older snapshot as an admission gate strands queued cards even after
+ capacity frees. The serialized fresh reservation below is the sole capacity authority; the
+ diagnostic remains available to explain a rejection from that authoritative check.
*/
- if (concurrencyDiagnostic.available <= 0) {
- if (reservedScope) {
- activeScopes.delete(task.id);
- activeScopeColumns.delete(task.id);
- }
- const reason = formatConcurrencyLimitReason(concurrencyDiagnostic);
- await this.store.updateTask(task.id, { status: "queued" });
- await this.logDispatchQueuedReason(task.id, reason, formatConcurrencyLimitMemoKey(concurrencyDiagnostic));
- return null;
- }
/*
FNXC:WorktreeCapacity 2026-08-01-04:38:
@@ -2860,12 +2849,28 @@ export class Scheduler {
claim the final slot. The reservation remains until the executor observes the persisted
WIP row and takes the handoff.
*/
+ let finalClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
+ const getFinalClaimSnapshot = () => finalClaimSnapshot ??= (async () => {
+ /*
+ FNXC:WorkflowContinuationCapacity 2026-08-01-07:10:
+ Worktree preparation and startup recovery can make the sweep's original task list stale
+ before this serialized admission point. A planner that became live after that snapshot
+ was absent from `activeWorktreeTaskIds`; once its handoff reservation transferred to the
+ durable planning status, the coordinator could no longer see either claim and admitted a
+ tenth active task against maxWorktrees=9. Re-read full rows lazily inside the coordinator
+ drain so pending workflow-step leases and every newly durable lane holder participate in
+ the final decision. Same-sweep transient starts remain covered by coordinator reservations.
+ */
+ const liveTasks = await this.store.listTasks({ slim: false, includeArchived: false });
+ const ids = await persistedTopLevelAgentTaskIdsFromStore(this.store, liveTasks);
+ return { count: ids.length, ids };
+ })();
const projectSlotReserved = await projectAdmissionCoordinator.reserveIfAvailable({
projectId: this.store.getRootDir(),
taskId: task.id,
maxConcurrent: activeTaskLimit,
- claimed: () => activeWorktreeTaskIds.length,
- claimedTaskIds: () => activeWorktreeTaskIds,
+ claimed: async () => (await getFinalClaimSnapshot()).count,
+ claimedTaskIds: async () => (await getFinalClaimSnapshot()).ids,
});
if (!projectSlotReserved) {
if (reservedScope) {
diff --git a/packages/engine/src/triage.ts b/packages/engine/src/triage.ts
index b67c7f54ab..8d69097c65 100644
--- a/packages/engine/src/triage.ts
+++ b/packages/engine/src/triage.ts
@@ -2273,7 +2273,9 @@ export class TriageProcessor {
worktreeBudget -= 1;
let freshClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
const getFreshClaimSnapshot = () => freshClaimSnapshot ??= (async () => {
- const fresh = await this.store.listTasks({ slim: true, includeArchived: false });
+ // Full rows are required here: a pending optional workflow-step lease can be the task's
+ // only live-agent signal, and slim rows intentionally omit workflowStepResults.
+ const fresh = await this.store.listTasks({ slim: false, includeArchived: false });
let pending = 0;
for (const id of this.processing) {
const row = fresh.find((task) => task.id === id);