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;
|
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:
|
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.
|
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
|
&& !visualStatus
|
||||||
&& !planReviewRunning
|
&& !planReviewRunning
|
||||||
&& !isAgentActive;
|
&& !isAgentActive;
|
||||||
const showReadyBadge = showIdleTodoBadge && !awaitingPlanning;
|
const showReadyBadge = showIdleTodoBadge && !queued && !awaitingPlanning;
|
||||||
/*
|
/*
|
||||||
FNXC:CodingIdeasWorkflow 2026-07-25-12:05:
|
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
|
"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
|
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.
|
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.
|
// 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
|
// On touch-primary devices the `draggable` attribute still arms the browser's
|
||||||
// touch-drag heuristic, which intermittently hijacks horizontal swipes meant
|
// touch-drag heuristic, which intermittently hijacks horizontal swipes meant
|
||||||
@@ -2130,8 +2130,6 @@ function TaskCardComponent({
|
|||||||
const showAddressPrFeedbackAction = canStartPrFeedbackAddressing(task, taskColumnFlags);
|
const showAddressPrFeedbackAction = canStartPrFeedbackAddressing(task, taskColumnFlags);
|
||||||
const metaRowVisible =
|
const metaRowVisible =
|
||||||
(task.dependencies?.length ?? 0) > 0
|
(task.dependencies?.length ?? 0) > 0
|
||||||
|| queued
|
|
||||||
|| task.status === "queued"
|
|
||||||
|| Boolean(task.blockedBy)
|
|| Boolean(task.blockedBy)
|
||||||
|| Boolean(task.overlapBlockedBy)
|
|| Boolean(task.overlapBlockedBy)
|
||||||
|| Boolean(fanout && fanout.totalCount > 0);
|
|| Boolean(fanout && fanout.totalCount > 0);
|
||||||
@@ -3370,6 +3368,16 @@ function TaskCardComponent({
|
|||||||
&& !visualStatus
|
&& !visualStatus
|
||||||
&& Boolean(task.recentAgentActivityAt)
|
&& Boolean(task.recentAgentActivityAt)
|
||||||
&& isAgentActive;
|
&& 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
|
const showStatusBadge = !isPaused
|
||||||
&& (hasTaskStatusBadge(visualStatus) || isTransientPlannerActive)
|
&& (hasTaskStatusBadge(visualStatus) || isTransientPlannerActive)
|
||||||
&& visualStatus !== "queued";
|
&& visualStatus !== "queued";
|
||||||
@@ -3400,6 +3408,8 @@ function TaskCardComponent({
|
|||||||
*/
|
*/
|
||||||
: showQueuedToPlanBadge
|
: showQueuedToPlanBadge
|
||||||
? t("tasks.queuedToPlan", "Queued to plan")
|
? t("tasks.queuedToPlan", "Queued to plan")
|
||||||
|
: showQueuedBadge
|
||||||
|
? t("tasks.statusQueued", "Queued")
|
||||||
: getTaskStatusLabel(visualStatus ?? "", t, showOptionalGateBadge ? undefined : getRunningWorkflowStepLabel(task), { idle: !isAgentActive, overlapBlockedBy: task.overlapBlockedBy ?? null });
|
: getTaskStatusLabel(visualStatus ?? "", t, showOptionalGateBadge ? undefined : getRunningWorkflowStepLabel(task), { idle: !isAgentActive, overlapBlockedBy: task.overlapBlockedBy ?? null });
|
||||||
const hasCardMetaBadges = showPriorityBadge
|
const hasCardMetaBadges = showPriorityBadge
|
||||||
|| task.executionMode === "fast"
|
|| task.executionMode === "fast"
|
||||||
@@ -3411,6 +3421,7 @@ function TaskCardComponent({
|
|||||||
|| showStatusBadge
|
|| showStatusBadge
|
||||||
|| showOptionalGateBadge
|
|| showOptionalGateBadge
|
||||||
|| showReadyBadge
|
|| showReadyBadge
|
||||||
|
|| showQueuedBadge
|
||||||
// FNXC:CodingIdeasWorkflow 2026-07-25-12:05: the header wrapper only renders when it has a
|
// 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
|
// 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).
|
// 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")}
|
: pausedByAgent ? t("tasks.pausedByAgent", "paused by agent") : t("tasks.paused", "paused")}
|
||||||
</span>
|
</span>
|
||||||
)}
|
)}
|
||||||
{(showStatusBadge || showQueuedToPlanBadge) && (
|
{(showStatusBadge || showQueuedToPlanBadge || showQueuedBadge) && (
|
||||||
<span
|
<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" : ""}`}
|
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={
|
title={
|
||||||
@@ -3573,6 +3584,10 @@ function TaskCardComponent({
|
|||||||
"tasks.queuedToPlanTitle",
|
"tasks.queuedToPlanTitle",
|
||||||
"Waiting for a planning slot — planning starts when an agent slot frees up",
|
"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
|
: task.status === "needs-replan" && !isAgentActive
|
||||||
? t(
|
? t(
|
||||||
"tasks.needsReplanQueuedTitle",
|
"tasks.needsReplanQueuedTitle",
|
||||||
@@ -3585,6 +3600,12 @@ function TaskCardComponent({
|
|||||||
data-awaiting-approval-reason={isAwaitingApproval ? (task.awaitingApprovalReason ?? "manual") : undefined}
|
data-awaiting-approval-reason={isAwaitingApproval ? (task.awaitingApprovalReason ?? "manual") : undefined}
|
||||||
>
|
>
|
||||||
{statusBadgeLabel}
|
{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>
|
</span>
|
||||||
)}
|
)}
|
||||||
{showOptionalGateBadge && optionalGateBadge && (
|
{showOptionalGateBadge && optionalGateBadge && (
|
||||||
@@ -4183,7 +4204,6 @@ function TaskCardComponent({
|
|||||||
</span>
|
</span>
|
||||||
</span>
|
</span>
|
||||||
)}
|
)}
|
||||||
{(queued || task.status === "queued") && !isWipColumn && <span className="queued-badge"><Clock size={12} style={{ verticalAlign: "middle" }} /> {t("tasks.queued", "Queued")}</span>}
|
|
||||||
{placeFooterRightInMeta && footerRightCluster}
|
{placeFooterRightInMeta && footerRightCluster}
|
||||||
</div>
|
</div>
|
||||||
)}
|
)}
|
||||||
|
|||||||
@@ -26,12 +26,12 @@ import { getPriorityColorVar, getPriorityLabel } from "../../utils/priorityIndic
|
|||||||
|
|
||||||
// Mock lucide-react to avoid SVG rendering issues in test env
|
// Mock lucide-react to avoid SVG rendering issues in test env
|
||||||
vi.mock("lucide-react", () => ({
|
vi.mock("lucide-react", () => ({
|
||||||
Link: () => null,
|
Link: (props: React.SVGProps<SVGSVGElement>) => <svg {...props} />,
|
||||||
GitBranch: () => null,
|
GitBranch: () => null,
|
||||||
Gitlab: () => null,
|
Gitlab: () => null,
|
||||||
Clock: () => null,
|
Clock: () => null,
|
||||||
Pencil: () => null,
|
Pencil: () => null,
|
||||||
Layers: () => null,
|
Layers: (props: React.SVGProps<SVGSVGElement>) => <svg {...props} />,
|
||||||
ChevronDown: () => null,
|
ChevronDown: () => null,
|
||||||
Folder: () => null,
|
Folder: () => null,
|
||||||
GitPullRequest: () => null,
|
GitPullRequest: () => null,
|
||||||
@@ -4847,12 +4847,12 @@ describe("TaskCard", () => {
|
|||||||
expect(screen.getByTestId("provider-icon-github")).toBeDefined();
|
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(
|
const { container } = render(
|
||||||
<TaskCard
|
<TaskCard
|
||||||
task={makeTask({
|
task={makeTask({
|
||||||
column: "todo",
|
column: "todo",
|
||||||
status: "queued",
|
status: null,
|
||||||
sourceType: "dashboard_ui",
|
sourceType: "dashboard_ui",
|
||||||
githubTracking: {
|
githubTracking: {
|
||||||
issue: {
|
issue: {
|
||||||
@@ -4864,20 +4864,41 @@ describe("TaskCard", () => {
|
|||||||
},
|
},
|
||||||
},
|
},
|
||||||
})}
|
})}
|
||||||
|
queued
|
||||||
onOpenDetail={noop}
|
onOpenDetail={noop}
|
||||||
addToast={noop}
|
addToast={noop}
|
||||||
/>,
|
/>,
|
||||||
);
|
);
|
||||||
|
|
||||||
const link = screen.getByRole("link", { name: "Linked GitHub issue #42" });
|
const link = screen.getByRole("link", { name: "Linked GitHub issue #42" });
|
||||||
const metaRow = container.querySelector(".card-meta");
|
const queuedBadge = screen.getByText("Queued");
|
||||||
const queuedBadge = container.querySelector(".queued-badge");
|
expect(link.closest(".card-footer-row")).not.toBeNull();
|
||||||
expect(container.querySelector(".card-footer-row")).toBeNull();
|
expect(queuedBadge).toHaveClass("card-status-badge", "card-status-badge--todo");
|
||||||
expect(link.closest(".card-meta")).toBe(metaRow);
|
expect(queuedBadge.closest(".card-header-badges")).not.toBeNull();
|
||||||
expect(link.closest(".card-footer-row-right")?.closest(".card-meta")).toBe(metaRow);
|
expect(container.querySelector(".queued-badge")).toBeNull();
|
||||||
expect(container.querySelector(".card-bottom-right-row")).toBeNull();
|
expect(queuedBadge.querySelector("svg")).toBeNull();
|
||||||
expect(queuedBadge).not.toBeNull();
|
expect(link.closest(".card-meta")).toBeNull();
|
||||||
expect(queuedBadge?.compareDocumentPosition(link) & Node.DOCUMENT_POSITION_FOLLOWING).toBeTruthy();
|
});
|
||||||
|
|
||||||
|
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;
|
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) === */
|
/* === Dependency Dropdown (shared: InlineCreateCard, NewTaskModal, TaskDetailModal, TaskForm) === */
|
||||||
.dep-trigger-wrap {
|
.dep-trigger-wrap {
|
||||||
position: relative;
|
position: relative;
|
||||||
|
|||||||
@@ -1176,6 +1176,75 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
coordinator.releaseReservation("FN-PLANNING");
|
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 () => {
|
it("admits the oldest same-project candidate atomically and partitions projects", async () => {
|
||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const started: string[] = [];
|
const started: string[] = [];
|
||||||
|
|||||||
@@ -2494,6 +2494,74 @@ describe("ProjectEngine paused in-review auto-merge behavior", () => {
|
|||||||
await engine.stop();
|
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 () => {
|
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 });
|
const mockStore = createMockStore({ ...baseSettings, autoMerge: true });
|
||||||
// The dequeued + merged task is a shared branch-group member, so the engine
|
// 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");
|
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 () => {
|
it("leaves a task queued when the authoritative release move rejects after reservation", async () => {
|
||||||
const ready = task({ id: "FN-002", status: "queued" });
|
const ready = task({ id: "FN-002", status: "queued" });
|
||||||
const store = storeWith([ready], { maxConcurrent: 4, maxWorktrees: 4 });
|
const store = storeWith([ready], { maxConcurrent: 4, maxWorktrees: 4 });
|
||||||
|
|||||||
@@ -72,7 +72,9 @@ function createStore(
|
|||||||
const found = tasks.find((candidate) => candidate.id === id);
|
const found = tasks.find((candidate) => candidate.id === id);
|
||||||
return found ? { ...found, prompt: "", attachments: [], comments: [] } : null;
|
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({
|
getSettings: vi.fn().mockResolvedValue({
|
||||||
maxConcurrent: 20,
|
maxConcurrent: 20,
|
||||||
maxWorktrees,
|
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(throttle.length, "a live task must still consume its worktree-capacity slot").toBeGreaterThan(0);
|
||||||
expect(specifyTask).not.toHaveBeenCalled();
|
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;
|
claimed: () => Promise<number> | number;
|
||||||
claimedTaskIds?: () => Promise<Iterable<string>> | Iterable<string>;
|
claimedTaskIds?: () => Promise<Iterable<string>> | Iterable<string>;
|
||||||
}): Promise<number> {
|
}): 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();
|
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());
|
const claimedIds = new Set(await params.claimedTaskIds());
|
||||||
|
for (const taskId of this.reservations.get(params.projectId) ?? []) reservations.add(taskId);
|
||||||
let pendingReservations = 0;
|
let pendingReservations = 0;
|
||||||
for (const taskId of this.reservations.get(params.projectId) ?? []) {
|
for (const taskId of reservations) {
|
||||||
if (!claimedIds.has(taskId)) pendingReservations += 1;
|
if (!claimedIds.has(taskId)) pendingReservations += 1;
|
||||||
}
|
}
|
||||||
return claimed + pendingReservations;
|
return claimed + pendingReservations;
|
||||||
|
|||||||
@@ -457,6 +457,14 @@ export class ProjectEngine {
|
|||||||
// ── Auto-merge state ──
|
// ── Auto-merge state ──
|
||||||
private mergeQueue: string[] = [];
|
private mergeQueue: string[] = [];
|
||||||
private mergeActive = new Set<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. */
|
/** Merge ids selected by the shared coordinator but not yet handed to rawMerge. */
|
||||||
private readonly coordinatorAdmittedMergeTaskIds = new Set<string>();
|
private readonly coordinatorAdmittedMergeTaskIds = new Set<string>();
|
||||||
private unregisterMergeAdmissionProvider?: () => void;
|
private unregisterMergeAdmissionProvider?: () => void;
|
||||||
@@ -805,7 +813,9 @@ export class ProjectEngine {
|
|||||||
window, checking it in addition to `mergeQueue` closes that TOCTOU gap.
|
window, checking it in addition to `mergeQueue` closes that TOCTOU gap.
|
||||||
*/
|
*/
|
||||||
isMergePending(taskId: string): boolean {
|
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);
|
clearInterval(this.mergeActiveReconcileTimer);
|
||||||
this.mergeActiveReconcileTimer = null;
|
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();
|
this.stopPlannerOverseerPoll();
|
||||||
|
|
||||||
/*
|
/*
|
||||||
@@ -2311,6 +2329,15 @@ export class ProjectEngine {
|
|||||||
};
|
};
|
||||||
abort = () => {
|
abort = () => {
|
||||||
this.removeMergeResolver(taskId, resolver);
|
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) {
|
if (this.activeMergeTaskId === taskId) {
|
||||||
this.mergeAbortController?.abort();
|
this.mergeAbortController?.abort();
|
||||||
this.mergeAbortController = null;
|
this.mergeAbortController = null;
|
||||||
@@ -2328,7 +2355,7 @@ export class ProjectEngine {
|
|||||||
|
|
||||||
// If this task is already queued or actively merging, wait for the
|
// If this task is already queued or actively merging, wait for the
|
||||||
// existing merge to finish rather than starting a second one.
|
// 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)) {
|
if (!this.internalEnqueueMerge(taskId)) {
|
||||||
this.removeMergeResolver(taskId, resolver);
|
this.removeMergeResolver(taskId, resolver);
|
||||||
@@ -2813,6 +2840,7 @@ export class ProjectEngine {
|
|||||||
|
|
||||||
private internalEnqueueMerge(taskId: string): boolean {
|
private internalEnqueueMerge(taskId: string): boolean {
|
||||||
if (this.shuttingDown || !this.started) return false;
|
if (this.shuttingDown || !this.started) return false;
|
||||||
|
if (this.capacityDeferredMergeTaskIds.has(taskId)) return false;
|
||||||
if (this.mergeActive.has(taskId)) {
|
if (this.mergeActive.has(taskId)) {
|
||||||
// Distinguish "actually being processed" (queued or active) from a
|
// Distinguish "actually being processed" (queued or active) from a
|
||||||
// leaked entry. Reconcile leaks immediately so recovery paths and fresh
|
// leaked entry. Reconcile leaks immediately so recovery paths and fresh
|
||||||
@@ -3895,7 +3923,9 @@ export class ProjectEngine {
|
|||||||
const admissionSettings = await store.getSettings();
|
const admissionSettings = await store.getSettings();
|
||||||
let mergeClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
let mergeClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
||||||
const getMergeClaimSnapshot = () => mergeClaimSnapshot ??= (async () => {
|
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);
|
const ids = await persistedTopLevelAgentTaskIdsFromStore(store, tasks);
|
||||||
return { count: ids.length, ids };
|
return { count: ids.length, ids };
|
||||||
})();
|
})();
|
||||||
@@ -3943,6 +3973,36 @@ export class ProjectEngine {
|
|||||||
projectAdmissionCoordinator.releaseReservation(taskId);
|
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) {
|
if (mergeStrategy === "pull-request" && this.options.processPullRequestMerge && !routeWorkspaceDirect) {
|
||||||
/*
|
/*
|
||||||
@@ -3966,10 +4026,11 @@ export class ProjectEngine {
|
|||||||
});
|
});
|
||||||
if (result === undefined) {
|
if (result === undefined) {
|
||||||
// Another older lane won the shared capacity pass. Re-queue rather
|
// Another older lane won the shared capacity pass. Re-queue rather
|
||||||
// than treating this deferral as a pull-request merge failure.
|
// than treating this deferral as a pull-request merge failure. End
|
||||||
this.mergeActive.delete(taskId);
|
// this drain: continuing would dequeue the same item immediately
|
||||||
this.internalEnqueueMerge(taskId);
|
// and spin at full speed while capacity remains unavailable.
|
||||||
continue;
|
deferMergeForCapacity();
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
if (result === "merged") {
|
if (result === "merged") {
|
||||||
runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR merged: ${taskId}`);
|
runtimeLog.log(`${hasManualResolver ? "Manual" : "Auto"}-merge PR merged: ${taskId}`);
|
||||||
@@ -4127,9 +4188,10 @@ export class ProjectEngine {
|
|||||||
if (!result) {
|
if (!result) {
|
||||||
// An older lane won this admission pass. Keep this merge queued;
|
// An older lane won this admission pass. Keep this merge queued;
|
||||||
// treating the deferral as a merge failure would consume retries.
|
// treating the deferral as a merge failure would consume retries.
|
||||||
this.mergeActive.delete(taskId);
|
// Exit this drain so the queued item waits for a normal future wake
|
||||||
this.internalEnqueueMerge(taskId);
|
// instead of retrying the same denied admission in a tight loop.
|
||||||
continue;
|
deferMergeForCapacity();
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
this.activeMergeSession = null;
|
this.activeMergeSession = null;
|
||||||
|
|||||||
@@ -2836,22 +2836,11 @@ export class Scheduler {
|
|||||||
topLevelClaimedSlots: reservedWorktreeSlots,
|
topLevelClaimedSlots: reservedWorktreeSlots,
|
||||||
});
|
});
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowScheduling 2026-06-23-20:58:
|
The sweep diagnostic is intentionally descriptive only. Worktree preparation can outlive
|
||||||
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.
|
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
|
||||||
FNXC:GlobalConcurrencyControls 2026-07-14-18:30:
|
diagnostic remains available to explain a rejection from that authoritative check.
|
||||||
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.
|
|
||||||
*/
|
*/
|
||||||
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:
|
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
|
claim the final slot. The reservation remains until the executor observes the persisted
|
||||||
WIP row and takes the handoff.
|
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({
|
const projectSlotReserved = await projectAdmissionCoordinator.reserveIfAvailable({
|
||||||
projectId: this.store.getRootDir(),
|
projectId: this.store.getRootDir(),
|
||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
maxConcurrent: activeTaskLimit,
|
maxConcurrent: activeTaskLimit,
|
||||||
claimed: () => activeWorktreeTaskIds.length,
|
claimed: async () => (await getFinalClaimSnapshot()).count,
|
||||||
claimedTaskIds: () => activeWorktreeTaskIds,
|
claimedTaskIds: async () => (await getFinalClaimSnapshot()).ids,
|
||||||
});
|
});
|
||||||
if (!projectSlotReserved) {
|
if (!projectSlotReserved) {
|
||||||
if (reservedScope) {
|
if (reservedScope) {
|
||||||
|
|||||||
@@ -2273,7 +2273,9 @@ export class TriageProcessor {
|
|||||||
worktreeBudget -= 1;
|
worktreeBudget -= 1;
|
||||||
let freshClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
let freshClaimSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
||||||
const getFreshClaimSnapshot = () => freshClaimSnapshot ??= (async () => {
|
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;
|
let pending = 0;
|
||||||
for (const id of this.processing) {
|
for (const id of this.processing) {
|
||||||
const row = fresh.find((task) => task.id === id);
|
const row = fresh.find((task) => task.id === id);
|
||||||
|
|||||||
Reference in New Issue
Block a user