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.
This commit is contained in:
gsxdsm
2026-08-01 00:32:25 -07:00
parent cced31208e
commit 9b82ff29e1
13 changed files with 425 additions and 59 deletions

View File

@@ -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.

View File

@@ -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.

View File

@@ -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")}
</span>
)}
{(showStatusBadge || showQueuedToPlanBadge) && (
{(showStatusBadge || showQueuedToPlanBadge || showQueuedBadge) && (
<span
className={`card-status-badge card-status-badge--${task.column}${showQueuedToPlanBadge ? " queued-to-plan" : ""}${isAwaitingApproval ? " awaiting-approval" : ""}${isPlanReviewReplanCapApproval ? " awaiting-approval--plan-review-replan-cap" : ""}${isAwaitingInput ? " awaiting-input" : ""}${isAgentActive ? " pulsing" : ""}${isFailed ? " failed" : ""}${isStuck ? " stuck" : ""}`}
title={
@@ -3573,6 +3584,10 @@ function TaskCardComponent({
"tasks.queuedToPlanTitle",
"Waiting for a planning slot — planning starts when an agent slot frees up",
)
: showQueuedBadge && task.overlapBlockedBy
? t("tasks.queuedFileOverlapTitle", "Queued due to file overlap with {{taskId}}", { taskId: task.overlapBlockedBy })
: showQueuedBadge && task.blockedBy
? t("tasks.queuedDependencyTitle", "Queued on dependency {{taskId}}", { taskId: task.blockedBy })
: task.status === "needs-replan" && !isAgentActive
? t(
"tasks.needsReplanQueuedTitle",
@@ -3585,6 +3600,12 @@ function TaskCardComponent({
data-awaiting-approval-reason={isAwaitingApproval ? (task.awaitingApprovalReason ?? "manual") : undefined}
>
{statusBadgeLabel}
{showQueuedBadge && task.overlapBlockedBy && (
<Layers className="card-queued-reason-icon" size={8} aria-hidden="true" data-testid={`card-queued-overlap-icon-${task.id}`} />
)}
{showQueuedBadge && !task.overlapBlockedBy && task.blockedBy && (
<Link className="card-queued-reason-icon" size={8} aria-hidden="true" data-testid={`card-queued-dependency-icon-${task.id}`} />
)}
</span>
)}
{showOptionalGateBadge && optionalGateBadge && (
@@ -4183,7 +4204,6 @@ function TaskCardComponent({
</span>
</span>
)}
{(queued || task.status === "queued") && !isWipColumn && <span className="queued-badge"><Clock size={12} style={{ verticalAlign: "middle" }} /> {t("tasks.queued", "Queued")}</span>}
{placeFooterRightInMeta && footerRightCluster}
</div>
)}

View File

@@ -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<SVGSVGElement>) => <svg {...props} />,
GitBranch: () => null,
Gitlab: () => null,
Clock: () => null,
Pencil: () => null,
Layers: () => null,
Layers: (props: React.SVGProps<SVGSVGElement>) => <svg {...props} />,
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(
<TaskCard
task={makeTask({
column: "todo",
status: "queued",
status: null,
sourceType: "dashboard_ui",
githubTracking: {
issue: {
@@ -4864,20 +4864,41 @@ describe("TaskCard", () => {
},
},
})}
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(
<TaskCard task={queuedTask} onOpenDetail={noop} addToast={noop} />,
);
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();
});

View File

@@ -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;

View File

@@ -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<void>((resolve) => { finishSnapshot = resolve; });
let snapshotStarted!: () => void;
const snapshotDidStart = new Promise<void>((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<void>((resolve) => { releaseDrain = resolve; });
let drainStarted!: () => void;
const drainDidStart = new Promise<void>((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[] = [];

View File

@@ -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<string>;
};
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

View File

@@ -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 });

View File

@@ -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();
});
});

View File

@@ -139,11 +139,22 @@ export class ProjectAdmissionCoordinator {
claimed: () => Promise<number> | number;
claimedTaskIds?: () => Promise<Iterable<string>> | Iterable<string>;
}): Promise<number> {
/*
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;

View File

@@ -457,6 +457,14 @@ export class ProjectEngine {
// ── Auto-merge state ──
private mergeQueue: string[] = [];
private mergeActive = new Set<string>();
/** Capacity-deferred ids stay out of the runnable queue until their retry timer fires. */
private readonly capacityDeferredMergeTaskIds = new Set<string>();
private readonly capacityDeferredMerges = new Map<string, {
timer: ReturnType<typeof setTimeout>;
resolvers: MergeResolver[];
generation: number;
manual: boolean;
}>();
/** Merge ids selected by the shared coordinator but not yet handed to rawMerge. */
private readonly coordinatorAdmittedMergeTaskIds = new Set<string>();
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;

View File

@@ -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) {

View File

@@ -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);