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:
7
.changeset/steady-active-capacity.md
Normal file
7
.changeset/steady-active-capacity.md
Normal 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.
|
||||
@@ -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.
|
||||
|
||||
@@ -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>
|
||||
)}
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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[] = [];
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user