From 9b82ff29e15aa5b352a5986037a52ea7ef205a06 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sat, 1 Aug 2026 00:32:25 -0700 Subject: [PATCH] fix: enforce active capacity and clarify queued cards Refresh canonical live-task claims at serialized scheduler, planning, and merge admission boundaries, including workflow-step leases and reservation handoffs. Back off capacity-denied merges safely across abort and restart lifecycles. Render queued planning cards in the header badge family with compact reason-specific icons. --- .changeset/steady-active-capacity.md | 7 ++ .../dashboard/app/components/TaskCard.css | 5 ++ .../dashboard/app/components/TaskCard.tsx | 32 ++++++-- .../components/__tests__/TaskCard.test.tsx | 45 +++++++--- packages/dashboard/app/styles.css | 10 --- .../engine/src/__tests__/concurrency.test.ts | 69 ++++++++++++++++ .../src/__tests__/project-engine.test.ts | 68 +++++++++++++++ .../scheduler-workflow-cutover.test.ts | 82 +++++++++++++++++++ ...sion-worktree-ledger-renamed-lanes.test.ts | 26 +++++- packages/engine/src/concurrency.ts | 15 +++- packages/engine/src/project-engine.ts | 82 ++++++++++++++++--- packages/engine/src/scheduler.ts | 39 +++++---- packages/engine/src/triage.ts | 4 +- 13 files changed, 425 insertions(+), 59 deletions(-) create mode 100644 .changeset/steady-active-capacity.md 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 && ( + )} {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);