FN-8705: prioritize review and execution slot admission
Prioritize lifecycle-critical work whenever project capacity becomes available. - Rank admission candidates as review/merge, execution, then planning. - Coordinate scheduler handoffs with project and host capacity reservations. - Cover lane priority and document the updated operator behavior. Files changed: .changeset/fn-8705-slot-priority.md | 7 ++ docs/architecture.md | 2 +- docs/dashboard-guide.md | 2 +- packages/engine/src/__tests__/concurrency.test.ts | 95 +++++++++++----------- .../engine/src/__tests__/project-engine.test.ts | 4 +- .../starved-refinement-x-triage-poll.test.ts | 11 ++- .../__tests__/triage-refinement-routing.test.ts | 16 ++-- .../workflow-continuation-capacity.test.ts | 3 +- packages/engine/src/concurrency.ts | 37 +++++++-- packages/engine/src/project-engine.ts | 10 ++- packages/engine/src/runtimes/in-process-runtime.ts | 3 +- packages/engine/src/scheduler.ts | 40 ++++----- packages/engine/src/triage.ts | 7 +- 13 files changed, 133 insertions(+), 104 deletions(-) Fusion-Task-Id: FN-8705 Fusion-Task-Lineage: ef66360b-e504-4e3f-b25e-b032a719d8c0 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8705-slot-priority.md
Normal file
7
.changeset/fn-8705-slot-priority.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": minor
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Prioritize review and execution work before planning when a project slot opens.
|
||||||
|
category: feature
|
||||||
|
dev: Project admission ranks review, execute, then planning; age and task ID break ties within a lane.
|
||||||
@@ -1141,7 +1141,7 @@ The run-audit system records every mutation performed by the engine across four
|
|||||||
|
|
||||||
Events are tied to specific run IDs for end-to-end traceability.
|
Events are tied to specific run IDs for end-to-end traceability.
|
||||||
|
|
||||||
For scheduler concurrency diagnostics, the queued reason names the active limiter(s) and usage (for example `gate=maxConcurrent ...`). The reason includes the `bindingGates` (`maxConcurrent`/`maxWorktrees`/`semaphore`), per-gate `{ used, limit, slack }`, `holders`, and computed `available`. `maxConcurrent` is a per-project cap on enriched live top-level planning, execution, and review/merge agents; the host semaphore is the separate process-global pool. Free project capacity is admitted oldest-first across lanes, rather than by lane priority. `maxWorktrees` is also enforced inside `TaskStore.moveTaskInternal` when committing an allocated move into `in-progress`, making it a hard active execution worktree cap even when workflow WIP/`maxConcurrent` would allow more tasks. These queued-reason logs are transition-only: a newly emitted line indicates the limiter signature changed or the condition cleared and later reappeared, not that a poll loop simply observed the same blocked state again.
|
For scheduler concurrency diagnostics, the queued reason names the active limiter(s) and usage (for example `gate=maxConcurrent ...`). The reason includes the `bindingGates` (`maxConcurrent`/`maxWorktrees`/`semaphore`), per-gate `{ used, limit, slack }`, `holders`, and computed `available`. `maxConcurrent` is a per-project cap on enriched live top-level planning, execution, and review/merge agents; the host semaphore is the separate process-global pool. Each newly free project slot admits review/merge first, ready execution second, and planning last; oldest valid `createdAt` and task ID break ties only within that lifecycle lane. `maxWorktrees` is also enforced inside `TaskStore.moveTaskInternal` when committing an allocated move into `in-progress`, making it a hard active execution worktree cap even when workflow WIP/`maxConcurrent` would allow more tasks. These queued-reason logs are transition-only: a newly emitted line indicates the limiter signature changed or the condition cleared and later reappeared, not that a poll loop simply observed the same blocked state again.
|
||||||
|
|
||||||
**Run audit endpoints:**
|
**Run audit endpoints:**
|
||||||
- `GET /api/agents/:id/runs/:runId/audit` — Returns audit trail for a specific agent run
|
- `GET /api/agents/:id/runs/:runId/audit` — Returns audit trail for a specific agent run
|
||||||
|
|||||||
@@ -1332,7 +1332,7 @@ Features:
|
|||||||
<!-- FNXC:CommandCenter 2026-06-27-10:03: Tokens detail charts must show every model bucket returned by analytics for accurate spend attribution; Overview remains a compact top-model summary because its copy explicitly frames those cards as top consumers/share. -->
|
<!-- FNXC:CommandCenter 2026-06-27-10:03: Tokens detail charts must show every model bucket returned by analytics for accurate spend attribution; Overview remains a compact top-model summary because its copy explicitly frames those cards as top consumers/share. -->
|
||||||
<!-- FNXC:CommandCenterActivity 2026-06-30-00:00: Activity active-agent counts include both durable-agent usage events and ephemeral task-worker execution runs from agentRuns, because task execution can be visible without a matching usage_events row. -->
|
<!-- FNXC:CommandCenterActivity 2026-06-30-00:00: Activity active-agent counts include both durable-agent usage events and ephemeral task-worker execution runs from agentRuns, because task execution can be visible without a matching usage_events row. -->
|
||||||
<!-- FNXC:CommandCenterActivity 2026-07-01-00:00: Graph-owned workflow step sessions publish active-to-terminal agentRuns lifecycle rows with task lineage and step metadata, so daily activity and Activity throughput charts include new workflow execution without dashboard-side recounting. -->
|
<!-- FNXC:CommandCenterActivity 2026-07-01-00:00: Graph-owned workflow step sessions publish active-to-terminal agentRuns lifecycle rows with task lineage and step metadata, so daily activity and Activity throughput charts include new workflow execution without dashboard-side recounting. -->
|
||||||
- **Overview controls dashboard** includes AI engine stop/start backed by `globalPause` and current-project **Max concurrent tasks** plus **Max worktrees** controls. Both capacity sliders remain visible while settings load or fail, but are disabled until settings are editable; a failed load shows its error and an intentionally disabled worktree limit explains how to enable it in Settings. Max concurrency caps top-level working agents across planning, execution, and review/merge; free capacity is admitted oldest-first within the project. The footer reports **Waiting**, **Running (N/max)**, and **Blocked**; column headers report executing/total (live agents in the lane over card count). Nested helper agents remain parent-internal and may temporarily exceed the displayed top-level count.
|
- **Overview controls dashboard** includes AI engine stop/start backed by `globalPause` and current-project **Max concurrent tasks** plus **Max worktrees** controls. Both capacity sliders remain visible while settings load or fail, but are disabled until settings are editable; a failed load shows its error and an intentionally disabled worktree limit explains how to enable it in Settings. Max concurrency caps top-level working agents across planning, execution, and review/merge; a free project slot serves review/merge first, ready execution second, then planning, with age and task ID deciding order only within each lane. The footer reports **Waiting**, **Running (N/max)**, and **Blocked**; column headers report executing/total (live agents in the lane over card count). Nested helper agents remain parent-internal and may temporarily exceed the displayed top-level count.
|
||||||
<!-- FNXC:TeamArea 2026-07-18-12:30: FN-8351 moves organization export and import to the Team tab so team-level portability controls are not presented as Overview dashboard controls. -->
|
<!-- FNXC:TeamArea 2026-07-18-12:30: FN-8351 moves organization export and import to the Team tab so team-level portability controls are not presented as Overview dashboard controls. -->
|
||||||
- **Team tab — Org export / import** lets an operator download a portable organization JSON bundle or paste one for a dry-run preview before confirming the apply step. Exports are secret-scrubbed by default: credentials and tokens are never included, while safe secret references can remain for setup in the destination project.
|
- **Team tab — Org export / import** lets an operator download a portable organization JSON bundle or paste one for a dry-run preview before confirming the apply step. Exports are secret-scrubbed by default: credentials and tokens are never included, while safe secret references can remain for setup in the destination project.
|
||||||
- **Configuration versions** lives in **Settings → Project → Configuration Versions**. It lists recorded project-setting revisions newest first; select **Roll back** on any revision and confirm once to restore it. The restore is recorded as a new forward revision, so it can itself be undone without manually reconstructing settings.
|
- **Configuration versions** lives in **Settings → Project → Configuration Versions**. It lists recorded project-setting revisions newest first; select **Roll back** on any revision and confirm once to restore it. The restore is recorded as a new forward revision, so it can itself be undone without manually reconstructing settings.
|
||||||
|
|||||||
@@ -1136,31 +1136,32 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
for (const [lane, taskId, createdAt] of [
|
for (const [lane, taskId, createdAt] of [
|
||||||
["planning", "FN-PLANNING", "2026-01-01T00:00:00.000Z"],
|
["planning", "FN-PLANNING", "2026-01-01T00:00:00.000Z"],
|
||||||
["execute", "FN-EXECUTE", "2026-01-02T00:00:00.000Z"],
|
["execute", "FN-EXECUTE", "2026-01-02T00:00:00.000Z"],
|
||||||
["merge", "FN-MERGE", "2026-01-03T00:00:00.000Z"],
|
["review", "FN-MERGE", "2026-01-03T00:00:00.000Z"],
|
||||||
] as const) {
|
] as const) {
|
||||||
coordinator.registerProvider(lane, {
|
coordinator.registerProvider(lane, {
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
refresh: async () => [{
|
refresh: async () => [{
|
||||||
taskId,
|
taskId,
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
|
lane,
|
||||||
createdAt,
|
createdAt,
|
||||||
start: async () => { started.push(taskId); },
|
start: async () => { started.push(taskId); },
|
||||||
}],
|
}],
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
expect(await coordinator.admitOldest({
|
expect(await coordinator.admitNext({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
maxConcurrent: activeTaskLimit,
|
maxConcurrent: activeTaskLimit,
|
||||||
claimed: () => 8,
|
claimed: () => 8,
|
||||||
})).toBe("FN-PLANNING");
|
})).toBe("FN-MERGE");
|
||||||
expect(await coordinator.reserveIfAvailable({
|
expect(await coordinator.reserveIfAvailable({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
taskId: "FN-DIRECT-SCHEDULER",
|
taskId: "FN-DIRECT-SCHEDULER",
|
||||||
maxConcurrent: activeTaskLimit,
|
maxConcurrent: activeTaskLimit,
|
||||||
claimed: () => 8,
|
claimed: () => 8,
|
||||||
})).toBe(false);
|
})).toBe(false);
|
||||||
expect(started).toEqual(["FN-PLANNING"]);
|
expect(started).toEqual(["FN-MERGE"]);
|
||||||
|
|
||||||
// Once the selected task is durably live, its matching reservation is the
|
// Once the selected task is durably live, its matching reservation is the
|
||||||
// same slot—not a second occupant—so the next real slot remains usable.
|
// same slot—not a second occupant—so the next real slot remains usable.
|
||||||
@@ -1169,11 +1170,11 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
taskId: "FN-DIRECT-SCHEDULER",
|
taskId: "FN-DIRECT-SCHEDULER",
|
||||||
maxConcurrent: 10,
|
maxConcurrent: 10,
|
||||||
claimed: () => 9,
|
claimed: () => 9,
|
||||||
claimedTaskIds: () => ["FN-PLANNING"],
|
claimedTaskIds: () => ["FN-MERGE"],
|
||||||
})).toBe(true);
|
})).toBe(true);
|
||||||
|
|
||||||
coordinator.releaseReservation("FN-DIRECT-SCHEDULER");
|
coordinator.releaseReservation("FN-DIRECT-SCHEDULER");
|
||||||
coordinator.releaseReservation("FN-PLANNING");
|
coordinator.releaseReservation("FN-MERGE");
|
||||||
});
|
});
|
||||||
|
|
||||||
it("does not lose a holder that transfers from reservation to durable state during a claim read", async () => {
|
it("does not lose a holder that transfers from reservation to durable state during a claim read", async () => {
|
||||||
@@ -1249,25 +1250,25 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const started: string[] = [];
|
const started: string[] = [];
|
||||||
const candidates = [
|
const candidates = [
|
||||||
{ taskId: "FN-20", projectId: "a", createdAt: "2026-01-02T00:00:00.000Z", start: async () => { started.push("new"); } },
|
{ taskId: "FN-20", projectId: "a", lane: "execute" as const, createdAt: "2026-01-02T00:00:00.000Z", start: async () => { started.push("new"); } },
|
||||||
{ taskId: "FN-10", projectId: "a", createdAt: "2026-01-01T00:00:00.000Z", start: async () => { started.push("old"); } },
|
{ taskId: "FN-10", projectId: "a", lane: "execute" as const, createdAt: "2026-01-01T00:00:00.000Z", start: async () => { started.push("old"); } },
|
||||||
{ taskId: "FN-1", projectId: "b", createdAt: "2026-01-03T00:00:00.000Z", start: async () => { started.push("other-project"); } },
|
{ taskId: "FN-1", projectId: "b", lane: "execute" as const, createdAt: "2026-01-03T00:00:00.000Z", start: async () => { started.push("other-project"); } },
|
||||||
];
|
];
|
||||||
const sem = new AgentSemaphore(2);
|
const sem = new AgentSemaphore(2);
|
||||||
await Promise.all([
|
await Promise.all([
|
||||||
coordinator.admitOldest({ projectId: "a", maxConcurrent: 1, claimed: () => 0, refresh: async () => candidates, semaphore: sem }),
|
coordinator.admitNext({ projectId: "a", maxConcurrent: 1, claimed: () => 0, refresh: async () => candidates, semaphore: sem }),
|
||||||
coordinator.admitOldest({ projectId: "a", maxConcurrent: 1, claimed: () => started.length, refresh: async () => candidates, semaphore: sem }),
|
coordinator.admitNext({ projectId: "a", maxConcurrent: 1, claimed: () => started.length, refresh: async () => candidates, semaphore: sem }),
|
||||||
]);
|
]);
|
||||||
expect(started).toEqual(["old"]);
|
expect(started).toEqual(["old"]);
|
||||||
sem.release();
|
sem.release();
|
||||||
await coordinator.admitOldest({ projectId: "b", maxConcurrent: 1, claimed: () => 0, refresh: async () => candidates, semaphore: sem });
|
await coordinator.admitNext({ projectId: "b", maxConcurrent: 1, claimed: () => 0, refresh: async () => candidates, semaphore: sem });
|
||||||
expect(started).toEqual(["old", "other-project"]);
|
expect(started).toEqual(["old", "other-project"]);
|
||||||
});
|
});
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:ConcurrencyAdmission 2026-07-26-09:45:
|
FNXC:ConcurrencyAdmission 2026-07-26-09:45:
|
||||||
Regression for the planning-starvation half of the FN-8600 incident: a card sat "Queued to plan"
|
Regression for the planning-starvation half of the FN-8600 incident: a card sat "Queued to plan"
|
||||||
while capacity was free, because admitOldest only ever evaluated candidates[0]. When the oldest
|
while capacity was free, because admitNext only ever evaluated candidates[0]. When the oldest
|
||||||
candidate's lane declines the handoff, younger work in another lane must still be admitted.
|
candidate's lane declines the handoff, younger work in another lane must still be admitted.
|
||||||
|
|
||||||
Invariant under test (not just the reported repro): a declining candidate is SKIPPED, not
|
Invariant under test (not just the reported repro): a declining candidate is SKIPPED, not
|
||||||
@@ -1280,7 +1281,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
const semaphore = new AgentSemaphore(4);
|
const semaphore = new AgentSemaphore(4);
|
||||||
const started: string[] = [];
|
const started: string[] = [];
|
||||||
|
|
||||||
const admitted = await coordinator.admitOldest({
|
const admitted = await coordinator.admitNext({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
maxConcurrent: 4,
|
maxConcurrent: 4,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1288,22 +1289,22 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
refresh: async () => [
|
refresh: async () => [
|
||||||
// Oldest, but its lane cannot start it (e.g. a merge id no longer queued).
|
// Oldest, but its lane cannot start it (e.g. a merge id no longer queued).
|
||||||
{
|
{
|
||||||
taskId: "FN-OLDEST", projectId: "project-a", createdAt: "2026-01-01T00:00:00.000Z",
|
taskId: "FN-OLDEST", projectId: "project-a", lane: "review", createdAt: "2026-01-01T00:00:00.000Z",
|
||||||
start: async () => { started.push("FN-OLDEST"); return false; },
|
start: async () => { started.push("FN-OLDEST"); return false; },
|
||||||
},
|
},
|
||||||
// Also declines — proves the walk continues past more than one.
|
// Also declines — proves the walk continues past more than one.
|
||||||
{
|
{
|
||||||
taskId: "FN-MIDDLE", projectId: "project-a", createdAt: "2026-01-02T00:00:00.000Z",
|
taskId: "FN-MIDDLE", projectId: "project-a", lane: "review", createdAt: "2026-01-02T00:00:00.000Z",
|
||||||
start: async () => { started.push("FN-MIDDLE"); return false; },
|
start: async () => { started.push("FN-MIDDLE"); return false; },
|
||||||
},
|
},
|
||||||
// The planning candidate that was starving behind them.
|
// The planning candidate that was starving behind them.
|
||||||
{
|
{
|
||||||
taskId: "FN-PLANNING", projectId: "project-a", createdAt: "2026-01-03T00:00:00.000Z",
|
taskId: "FN-PLANNING", projectId: "project-a", lane: "planning", createdAt: "2026-01-03T00:00:00.000Z",
|
||||||
start: async () => { started.push("FN-PLANNING"); },
|
start: async () => { started.push("FN-PLANNING"); },
|
||||||
},
|
},
|
||||||
// Younger still: must NOT be admitted, so skipping never becomes overtaking.
|
// Younger still: must NOT be admitted, so skipping never becomes overtaking.
|
||||||
{
|
{
|
||||||
taskId: "FN-YOUNGEST", projectId: "project-a", createdAt: "2026-01-04T00:00:00.000Z",
|
taskId: "FN-YOUNGEST", projectId: "project-a", lane: "planning", createdAt: "2026-01-04T00:00:00.000Z",
|
||||||
start: async () => { started.push("FN-YOUNGEST"); },
|
start: async () => { started.push("FN-YOUNGEST"); },
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
@@ -1332,11 +1333,11 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
// A pre-tryAcquire shim: release only, no tryAcquire.
|
// A pre-tryAcquire shim: release only, no tryAcquire.
|
||||||
const shim = { release: () => { releases.push(1); } };
|
const shim = { release: () => { releases.push(1); } };
|
||||||
|
|
||||||
const admitted = await coordinator.admitOldest({
|
const admitted = await coordinator.admitNext({
|
||||||
projectId: "project-shim",
|
projectId: "project-shim",
|
||||||
maxConcurrent: 4,
|
maxConcurrent: 4,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
semaphore: shim as unknown as Parameters<ProjectAdmissionCoordinator["admitOldest"]>[0]["semaphore"],
|
semaphore: shim as unknown as Parameters<ProjectAdmissionCoordinator["admitNext"]>[0]["semaphore"],
|
||||||
refresh: async () => [
|
refresh: async () => [
|
||||||
{ taskId: "FN-A", projectId: "project-shim", createdAt: "2026-01-01T00:00:00.000Z", start: async () => false },
|
{ taskId: "FN-A", projectId: "project-shim", createdAt: "2026-01-01T00:00:00.000Z", start: async () => false },
|
||||||
{ taskId: "FN-B", projectId: "project-shim", createdAt: "2026-01-02T00:00:00.000Z", start: async () => false },
|
{ taskId: "FN-B", projectId: "project-shim", createdAt: "2026-01-02T00:00:00.000Z", start: async () => false },
|
||||||
@@ -1362,7 +1363,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const semaphore = new AgentSemaphore(4);
|
const semaphore = new AgentSemaphore(4);
|
||||||
|
|
||||||
const admitted = await coordinator.admitOldest({
|
const admitted = await coordinator.admitNext({
|
||||||
projectId: "project-prehold",
|
projectId: "project-prehold",
|
||||||
maxConcurrent: 4,
|
maxConcurrent: 4,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1399,7 +1400,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const semaphore = new AgentSemaphore(2);
|
const semaphore = new AgentSemaphore(2);
|
||||||
|
|
||||||
await expect(coordinator.admitOldest({
|
await expect(coordinator.admitNext({
|
||||||
projectId: "project-throw",
|
projectId: "project-throw",
|
||||||
maxConcurrent: 4,
|
maxConcurrent: 4,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1421,7 +1422,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
expect(semaphore.tryAcquire()).toBe(true);
|
expect(semaphore.tryAcquire()).toBe(true);
|
||||||
const started: string[] = [];
|
const started: string[] = [];
|
||||||
|
|
||||||
const admitted = await coordinator.admitOldest({
|
const admitted = await coordinator.admitNext({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
maxConcurrent: 4,
|
maxConcurrent: 4,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1441,7 +1442,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
it("releases a rejected handoff and retains an accepted reservation until lane transfer", async () => {
|
it("releases a rejected handoff and retains an accepted reservation until lane transfer", async () => {
|
||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const semaphore = new AgentSemaphore(1);
|
const semaphore = new AgentSemaphore(1);
|
||||||
const rejected = await coordinator.admitOldest({
|
const rejected = await coordinator.admitNext({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
maxConcurrent: 1,
|
maxConcurrent: 1,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1456,7 +1457,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
|
|
||||||
let releaseStart!: () => void;
|
let releaseStart!: () => void;
|
||||||
const startBlocked = new Promise<void>((resolve) => { releaseStart = resolve; });
|
const startBlocked = new Promise<void>((resolve) => { releaseStart = resolve; });
|
||||||
const first = coordinator.admitOldest({
|
const first = coordinator.admitNext({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
maxConcurrent: 1,
|
maxConcurrent: 1,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1467,7 +1468,7 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
}],
|
}],
|
||||||
});
|
});
|
||||||
await Promise.resolve();
|
await Promise.resolve();
|
||||||
const second = coordinator.admitOldest({
|
const second = coordinator.admitNext({
|
||||||
projectId: "project-a",
|
projectId: "project-a",
|
||||||
maxConcurrent: 1,
|
maxConcurrent: 1,
|
||||||
claimed: () => 0,
|
claimed: () => 0,
|
||||||
@@ -1484,35 +1485,31 @@ describe("ProjectAdmissionCoordinator", () => {
|
|||||||
semaphore.release();
|
semaphore.release();
|
||||||
});
|
});
|
||||||
|
|
||||||
it("refreshes every registered lane before selecting the cross-lane oldest task", async () => {
|
it("refreshes every lane and admits review before older execution and planning", async () => {
|
||||||
const coordinator = new ProjectAdmissionCoordinator();
|
const coordinator = new ProjectAdmissionCoordinator();
|
||||||
const started: string[] = [];
|
const started: string[] = [];
|
||||||
coordinator.registerProvider("planning", {
|
const register = (lane: "review" | "execute" | "planning", taskId: string, createdAt: string, name: string) => {
|
||||||
projectId: "project-a",
|
coordinator.registerProvider(name, {
|
||||||
refresh: async () => [{
|
projectId: "project-a",
|
||||||
taskId: "FN-20", projectId: "project-a", createdAt: "2026-01-02T00:00:00.000Z",
|
refresh: async () => [{ taskId, projectId: "project-a", lane, createdAt, start: async () => { started.push(name); } }],
|
||||||
start: async () => { started.push("planner"); },
|
});
|
||||||
}],
|
};
|
||||||
});
|
register("planning", "FN-1", "2026-01-01T00:00:00.000Z", "planner");
|
||||||
coordinator.registerProvider("execute", {
|
register("execute", "FN-2", "2026-01-02T00:00:00.000Z", "executor");
|
||||||
projectId: "project-a",
|
register("review", "FN-3", "2026-01-03T00:00:00.000Z", "merge");
|
||||||
refresh: async () => [{
|
|
||||||
taskId: "FN-10", projectId: "project-a", createdAt: "2026-01-01T00:00:00.000Z",
|
|
||||||
start: async () => { started.push("executor"); },
|
|
||||||
}],
|
|
||||||
});
|
|
||||||
|
|
||||||
await coordinator.admitOldest({ projectId: "project-a", maxConcurrent: 1, claimed: () => 0 });
|
await coordinator.admitNext({ projectId: "project-a", maxConcurrent: 1, claimed: () => 0 });
|
||||||
expect(started).toEqual(["executor"]);
|
expect(started).toEqual(["merge"]);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("uses a stable total order for invalid timestamps and malformed ids", () => {
|
it("uses oldest valid age then task ID only within one lifecycle lane", () => {
|
||||||
const ordered = [
|
const ordered = [
|
||||||
{ taskId: "bad", createdAt: "not-a-date" },
|
{ taskId: "bad", lane: "execute" as const, createdAt: "not-a-date" },
|
||||||
{ taskId: "FN-12", createdAt: "2026-01-01T00:00:00.000Z" },
|
{ taskId: "FN-12", lane: "execute" as const, createdAt: "2026-01-01T00:00:00.000Z" },
|
||||||
{ taskId: "FN-2", createdAt: "2026-01-01T00:00:00.000Z" },
|
{ taskId: "FN-2", lane: "execute" as const, createdAt: "2026-01-01T00:00:00.000Z" },
|
||||||
{ taskId: "also-bad" },
|
{ taskId: "also-bad", lane: "execute" as const },
|
||||||
|
{ taskId: "FN-older-planning", lane: "planning" as const, createdAt: "2020-01-01T00:00:00.000Z" },
|
||||||
].sort(compareAdmissionCandidates);
|
].sort(compareAdmissionCandidates);
|
||||||
expect(ordered.map((item) => item.taskId)).toEqual(["FN-2", "FN-12", "also-bad", "bad"]);
|
expect(ordered.map((item) => item.taskId)).toEqual(["FN-2", "FN-12", "also-bad", "bad", "FN-older-planning"]);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -3864,8 +3864,8 @@ describe("U9 merge safeguards without prior coverage", () => {
|
|||||||
mergeRetries: 0,
|
mergeRetries: 0,
|
||||||
createdAt: new Date(0).toISOString(),
|
createdAt: new Date(0).toISOString(),
|
||||||
});
|
});
|
||||||
const admitted = (await mergeProvider.refresh()) as Array<{ taskId: string }>;
|
const admitted = (await mergeProvider.refresh()) as Array<{ taskId: string; lane: string }>;
|
||||||
expect(admitted.map((c) => c.taskId)).toEqual(["FN-paused"]);
|
expect(admitted).toMatchObject([{ taskId: "FN-paused", lane: "review" }]);
|
||||||
|
|
||||||
// Engine-level `paused` is the sibling half of the same filter.
|
// Engine-level `paused` is the sibling half of the same filter.
|
||||||
mockStore.store.getTask.mockResolvedValue({
|
mockStore.store.getTask.mockResolvedValue({
|
||||||
|
|||||||
@@ -83,12 +83,11 @@ describe("reliability interaction: starved refinement x triage poll", () => {
|
|||||||
|
|
||||||
(triage as any).running = true;
|
(triage as any).running = true;
|
||||||
/*
|
/*
|
||||||
FNXC:EngineTests 2026-07-23-21:30:
|
FNXC:ConcurrencyAdmission 2026-08-01-15:42:
|
||||||
FN-8453 (commit eef5eb751) replaced priority-based triage ordering with the unified
|
FN-8705 makes review and execution higher priority than planning, while
|
||||||
oldest-createdAt-first admission coordinator, so the self-healing priority bump no
|
preserving FIFO fairness among planning candidates. With higher-priority
|
||||||
longer reorders admission. The surviving reliability invariant is FIFO fairness:
|
lanes empty, maxConcurrent=1, and 6 older planning tasks, the starved
|
||||||
with maxConcurrent=1 and 6 older backlog tasks, the starved refinement must be
|
refinement must be admitted within 7 bounded polls (one admission per poll).
|
||||||
admitted within 7 bounded polls (one admission per poll).
|
|
||||||
*/
|
*/
|
||||||
for (let i = 0; i < 7; i++) {
|
for (let i = 0; i < 7; i++) {
|
||||||
await (triage as any).poll();
|
await (triage as any).poll();
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import { tmpdir } from "node:os";
|
|||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import type { Task } from "@fusion/core";
|
import type { Task } from "@fusion/core";
|
||||||
import { TriageProcessor } from "../triage.js";
|
import { TriageProcessor } from "../triage.js";
|
||||||
|
import { projectAdmissionCoordinator } from "../concurrency.js";
|
||||||
|
|
||||||
function withStoreEvents<T extends Record<string, unknown>>(store: T): T & { on: () => void; off: () => void } {
|
function withStoreEvents<T extends Record<string, unknown>>(store: T): T & { on: () => void; off: () => void } {
|
||||||
return {
|
return {
|
||||||
@@ -89,16 +90,19 @@ describe("refinement routing from triage", () => {
|
|||||||
const specifySpy = vi.spyOn(processor, "specifyTask").mockImplementation(async (task) => {
|
const specifySpy = vi.spyOn(processor, "specifyTask").mockImplementation(async (task) => {
|
||||||
const idx = tasks.findIndex((t) => t.id === task.id);
|
const idx = tasks.findIndex((t) => t.id === task.id);
|
||||||
if (idx >= 0) tasks[idx] = { ...tasks[idx], column: "todo" };
|
if (idx >= 0) tasks[idx] = { ...tasks[idx], column: "todo" };
|
||||||
|
// The production planner transfers this bridge when it becomes durable;
|
||||||
|
// this lightweight mock must model that handoff between bounded polls.
|
||||||
|
projectAdmissionCoordinator.releaseReservation(task.id);
|
||||||
});
|
});
|
||||||
|
|
||||||
(processor as any).running = true;
|
(processor as any).running = true;
|
||||||
/*
|
/*
|
||||||
FNXC:EngineTests 2026-07-23-21:30:
|
FNXC:ConcurrencyAdmission 2026-08-01-15:42:
|
||||||
FN-8453 (commit eef5eb751) replaced priority-then-refinement triage ordering with the
|
FN-8705 keeps FIFO fairness within the planning lane even though cross-lane
|
||||||
unified oldest-createdAt-first admission coordinator. Refinements no longer jump the
|
admission now prioritizes review and execution. Refinements do not jump the
|
||||||
same-priority backlog; the no-starvation invariant is now FIFO fairness — the newest
|
planning backlog; the newest refinement behind an 8-task backlog at
|
||||||
refinement behind an 8-task backlog at maxConcurrent=2 must be admitted within
|
maxConcurrent=2 must be admitted within ceil(9/2)=5 bounded polls once
|
||||||
ceil(9/2)=5 bounded polls.
|
higher-priority lanes have no accepting candidate.
|
||||||
*/
|
*/
|
||||||
for (let i = 0; i < 5; i++) {
|
for (let i = 0; i < 5; i++) {
|
||||||
await (processor as any).poll();
|
await (processor as any).poll();
|
||||||
|
|||||||
@@ -121,7 +121,7 @@ describe("workflow continuation active-slot admission", () => {
|
|||||||
const dispatch = vi.fn(async () => {});
|
const dispatch = vi.fn(async () => {});
|
||||||
let releaseBlocker!: () => void;
|
let releaseBlocker!: () => void;
|
||||||
const blockerStarted = new Promise<void>((resolveStarted) => {
|
const blockerStarted = new Promise<void>((resolveStarted) => {
|
||||||
void projectAdmissionCoordinator.admitOldest({
|
void projectAdmissionCoordinator.admitNext({
|
||||||
projectId: PROJECT_ID,
|
projectId: PROJECT_ID,
|
||||||
maxConcurrent: 9,
|
maxConcurrent: 9,
|
||||||
claimed: () => 8,
|
claimed: () => 8,
|
||||||
@@ -129,6 +129,7 @@ describe("workflow continuation active-slot admission", () => {
|
|||||||
refresh: async () => [{
|
refresh: async () => [{
|
||||||
taskId: "FN-BLOCKER",
|
taskId: "FN-BLOCKER",
|
||||||
projectId: PROJECT_ID,
|
projectId: PROJECT_ID,
|
||||||
|
lane: "execute",
|
||||||
createdAt: "2026-07-31T23:59:59.000Z",
|
createdAt: "2026-07-31T23:59:59.000Z",
|
||||||
start: async () => {
|
start: async () => {
|
||||||
resolveStarted();
|
resolveStarted();
|
||||||
|
|||||||
@@ -36,10 +36,21 @@ export function resolveActiveTaskCapacityLimit(params: {
|
|||||||
: Math.min(params.maxConcurrent, maxWorktrees);
|
: Math.min(params.maxConcurrent, maxWorktrees);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Lifecycle lanes ordered by the project admission coordinator. */
|
||||||
|
export type AdmissionLane = "review" | "execute" | "planning";
|
||||||
|
|
||||||
|
const admissionLanePriority: Record<AdmissionLane, number> = {
|
||||||
|
review: 0,
|
||||||
|
execute: 1,
|
||||||
|
planning: 2,
|
||||||
|
};
|
||||||
|
|
||||||
/** A task waiting to enter one of the top-level agent lanes. */
|
/** A task waiting to enter one of the top-level agent lanes. */
|
||||||
export interface AdmissionCandidate {
|
export interface AdmissionCandidate {
|
||||||
taskId: string;
|
taskId: string;
|
||||||
projectId: string;
|
projectId: string;
|
||||||
|
/** Explicit lifecycle ownership; priority never depends on provider or column names. */
|
||||||
|
lane: AdmissionLane;
|
||||||
createdAt?: string;
|
createdAt?: string;
|
||||||
/** Records ownership of the host reservation before the lane starts. */
|
/** Records ownership of the host reservation before the lane starts. */
|
||||||
reserve?: () => void;
|
reserve?: () => void;
|
||||||
@@ -57,13 +68,21 @@ export interface AdmissionProvider {
|
|||||||
refresh: () => Promise<AdmissionCandidate[]>;
|
refresh: () => Promise<AdmissionCandidate[]>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:ConcurrencyAdmission 2026-08-01-15:42:
|
||||||
|
FN-8705 requires every newly available project slot to finish review/merge work
|
||||||
|
before ready execution and planning. The lane is explicit on each candidate so
|
||||||
|
custom workflow column names and provider IDs cannot change lifecycle priority;
|
||||||
|
age and task ID only preserve fairness within the same lane.
|
||||||
|
*/
|
||||||
/**
|
/**
|
||||||
* Deterministic oldest-first ordering used for all task-lane admission.
|
* Deterministic lifecycle-lane ordering for project admission. Invalid/missing
|
||||||
* Invalid/missing timestamps deliberately sort after valid timestamps; numeric
|
* timestamps sort after valid timestamps only within one lane; numeric task ids
|
||||||
* task ids break normal ties before lexical ids so a malformed fixture cannot
|
* then lexical ids make malformed data deterministic.
|
||||||
* make Array.sort's NaN handling decide capacity admission.
|
|
||||||
*/
|
*/
|
||||||
export function compareAdmissionCandidates(a: Pick<AdmissionCandidate, "taskId" | "createdAt">, b: Pick<AdmissionCandidate, "taskId" | "createdAt">): number {
|
export function compareAdmissionCandidates(a: Pick<AdmissionCandidate, "taskId" | "createdAt" | "lane">, b: Pick<AdmissionCandidate, "taskId" | "createdAt" | "lane">): number {
|
||||||
|
const laneOrder = admissionLanePriority[a.lane] - admissionLanePriority[b.lane];
|
||||||
|
if (laneOrder !== 0) return laneOrder;
|
||||||
const aTime = a.createdAt ? Date.parse(a.createdAt) : Number.NaN;
|
const aTime = a.createdAt ? Date.parse(a.createdAt) : Number.NaN;
|
||||||
const bTime = b.createdAt ? Date.parse(b.createdAt) : Number.NaN;
|
const bTime = b.createdAt ? Date.parse(b.createdAt) : Number.NaN;
|
||||||
const aValid = Number.isFinite(aTime);
|
const aValid = Number.isFinite(aTime);
|
||||||
@@ -160,7 +179,11 @@ export class ProjectAdmissionCoordinator {
|
|||||||
return claimed + pendingReservations;
|
return claimed + pendingReservations;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Atomically reserve one live-task slot for a lane that performs its own dispatch sweep. */
|
/**
|
||||||
|
* Reserve only for compatibility callers that cannot supply a lifecycle candidate.
|
||||||
|
* Top-level production lanes must use admitNext so refreshed higher-priority work
|
||||||
|
* is considered before this project capacity is claimed.
|
||||||
|
*/
|
||||||
async reserveIfAvailable(params: {
|
async reserveIfAvailable(params: {
|
||||||
projectId: string;
|
projectId: string;
|
||||||
taskId: string;
|
taskId: string;
|
||||||
@@ -202,7 +225,7 @@ export class ProjectAdmissionCoordinator {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
async admitOldest(params: {
|
async admitNext(params: {
|
||||||
projectId: string;
|
projectId: string;
|
||||||
maxConcurrent: number;
|
maxConcurrent: number;
|
||||||
claimed: () => Promise<number> | number;
|
claimed: () => Promise<number> | number;
|
||||||
|
|||||||
@@ -687,6 +687,7 @@ export class ProjectEngine {
|
|||||||
return [{
|
return [{
|
||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
projectId,
|
projectId,
|
||||||
|
lane: "review",
|
||||||
createdAt: task.createdAt,
|
createdAt: task.createdAt,
|
||||||
start: async () => {
|
start: async () => {
|
||||||
// Do not run merge work in the coordinator; hand the exact queued
|
// Do not run merge work in the coordinator; hand the exact queued
|
||||||
@@ -3906,7 +3907,7 @@ export class ProjectEngine {
|
|||||||
a merge IS an agent, so it still consumes one of the project's slots; it
|
a merge IS an agent, so it still consumes one of the project's slots; it
|
||||||
just no longer consumes a machine-wide slot too.
|
just no longer consumes a machine-wide slot too.
|
||||||
|
|
||||||
`admitOldest` already takes `semaphore` as optional and enforces
|
`admitNext` already takes `semaphore` as optional and enforces
|
||||||
`maxConcurrent` independently of it (see its `claimed() + reservations >=
|
`maxConcurrent` independently of it (see its `claimed() + reservations >=
|
||||||
maxConcurrent` check), so dropping the argument keeps per-project
|
maxConcurrent` check), so dropping the argument keeps per-project
|
||||||
admission and oldest-first fairness exactly as they were.
|
admission and oldest-first fairness exactly as they were.
|
||||||
@@ -3933,7 +3934,7 @@ export class ProjectEngine {
|
|||||||
FNXC:ConcurrencyAdmission 2026-08-01-01:50 (ROOT CAUSE — triage admission died during every merge):
|
FNXC:ConcurrencyAdmission 2026-08-01-01:50 (ROOT CAUSE — triage admission died during every merge):
|
||||||
This lane previously ran `value = await start()` INSIDE its admission `start()` callback —
|
This lane previously ran `value = await start()` INSIDE its admission `start()` callback —
|
||||||
i.e. the ENTIRE merge (git rebase, verification, landing: minutes, or forever when the
|
i.e. the ENTIRE merge (git rebase, verification, landing: minutes, or forever when the
|
||||||
merge wedges) executed inside `admitOldest`'s single-flight drain. The coordinator is a
|
merge wedges) executed inside `admitNext`'s single-flight drain. The coordinator is a
|
||||||
project-wide singleton and every caller awaits the previous drain, so triage's poll parked
|
project-wide singleton and every caller awaits the previous drain, so triage's poll parked
|
||||||
at `await existing` for the whole merge window, its `polling` re-entrance guard stayed
|
at `await existing` for the whole merge window, its `polling` re-entrance guard stayed
|
||||||
closed, and every 15s tick + task:created wake dropped silently. Observed twice on the
|
closed, and every 15s tick + task:created wake dropped silently. Observed twice on the
|
||||||
@@ -3942,12 +3943,12 @@ export class ProjectEngine {
|
|||||||
the merge finished. With merge pinned at 1, every merge was a planning outage.
|
the merge finished. With merge pinned at 1, every merge was a planning outage.
|
||||||
|
|
||||||
The lane start now only CLAIMS the admission and returns; the merge body runs after
|
The lane start now only CLAIMS the admission and returns; the merge body runs after
|
||||||
`admitOldest` settles, outside the drain. Capacity stays honest: the merge row's own
|
`admitNext` settles, outside the drain. Capacity stays honest: the merge row's own
|
||||||
merging/landing status is what `claimed()` counts, and at-most-once merging is enforced
|
merging/landing status is what `claimed()` counts, and at-most-once merging is enforced
|
||||||
by the merge lease, not by this drain. The transient admit→status-write gap is the same
|
by the merge lease, not by this drain. The transient admit→status-write gap is the same
|
||||||
one every other lane (triage `void specifyTask`, scheduler `void schedule`) already has.
|
one every other lane (triage `void specifyTask`, scheduler `void schedule`) already has.
|
||||||
*/
|
*/
|
||||||
await projectAdmissionCoordinator.admitOldest({
|
await projectAdmissionCoordinator.admitNext({
|
||||||
projectId: cwd,
|
projectId: cwd,
|
||||||
maxConcurrent: resolveActiveTaskCapacityLimit({
|
maxConcurrent: resolveActiveTaskCapacityLimit({
|
||||||
maxConcurrent: admissionSettings.maxConcurrent ?? 2,
|
maxConcurrent: admissionSettings.maxConcurrent ?? 2,
|
||||||
@@ -3959,6 +3960,7 @@ export class ProjectEngine {
|
|||||||
refresh: async () => [{
|
refresh: async () => [{
|
||||||
taskId,
|
taskId,
|
||||||
projectId: cwd,
|
projectId: cwd,
|
||||||
|
lane: "review",
|
||||||
createdAt: mergeCandidate?.createdAt,
|
createdAt: mergeCandidate?.createdAt,
|
||||||
start: async () => {
|
start: async () => {
|
||||||
selected = true;
|
selected = true;
|
||||||
|
|||||||
@@ -478,7 +478,7 @@ export async function admitPlanningContinuation(input: {
|
|||||||
// turn; a pre-drain project snapshot can admit into its newly occupied slot.
|
// turn; a pre-drain project snapshot can admit into its newly occupied slot.
|
||||||
let admissionSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
let admissionSnapshot: Promise<{ count: number; ids: string[] }> | undefined;
|
||||||
const getAdmissionSnapshot = () => admissionSnapshot ??= loadClaimSnapshot();
|
const getAdmissionSnapshot = () => admissionSnapshot ??= loadClaimSnapshot();
|
||||||
await projectAdmissionCoordinator.admitOldest({
|
await projectAdmissionCoordinator.admitNext({
|
||||||
projectId: input.projectId,
|
projectId: input.projectId,
|
||||||
maxConcurrent: resolveActiveTaskCapacityLimit({
|
maxConcurrent: resolveActiveTaskCapacityLimit({
|
||||||
maxConcurrent: settings.maxConcurrent ?? 2,
|
maxConcurrent: settings.maxConcurrent ?? 2,
|
||||||
@@ -490,6 +490,7 @@ export async function admitPlanningContinuation(input: {
|
|||||||
refresh: async () => [{
|
refresh: async () => [{
|
||||||
taskId: input.task.id,
|
taskId: input.task.id,
|
||||||
projectId: input.projectId,
|
projectId: input.projectId,
|
||||||
|
lane: "execute",
|
||||||
createdAt: input.item.createdAt ?? input.task.createdAt,
|
createdAt: input.item.createdAt ?? input.task.createdAt,
|
||||||
start: async () => {
|
start: async () => {
|
||||||
// The preflight above is only a fast path. This serialized check is the
|
// The preflight above is only a fast path. This serialized check is the
|
||||||
|
|||||||
@@ -17,7 +17,6 @@ import { readFile } from "node:fs/promises";
|
|||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import {
|
import {
|
||||||
dropPreHeldExecutorSlot,
|
dropPreHeldExecutorSlot,
|
||||||
hasPreHeldExecutorSlot,
|
|
||||||
projectAdmissionCoordinator,
|
projectAdmissionCoordinator,
|
||||||
persistedTopLevelAgentTaskIdsFromStore,
|
persistedTopLevelAgentTaskIdsFromStore,
|
||||||
recoverIdleSemaphoreLeakCandidate,
|
recoverIdleSemaphoreLeakCandidate,
|
||||||
@@ -974,6 +973,7 @@ export class Scheduler {
|
|||||||
.map((task) => ({
|
.map((task) => ({
|
||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
projectId,
|
projectId,
|
||||||
|
lane: "execute",
|
||||||
createdAt: task.createdAt,
|
createdAt: task.createdAt,
|
||||||
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
|
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
|
||||||
start: async () => {
|
start: async () => {
|
||||||
@@ -2825,12 +2825,24 @@ export class Scheduler {
|
|||||||
const ids = await persistedTopLevelAgentTaskIdsFromStore(this.store, liveTasks);
|
const ids = await persistedTopLevelAgentTaskIdsFromStore(this.store, liveTasks);
|
||||||
return { count: ids.length, ids };
|
return { count: ids.length, ids };
|
||||||
})();
|
})();
|
||||||
const projectSlotReserved = await projectAdmissionCoordinator.reserveIfAvailable({
|
let projectSlotReserved = false;
|
||||||
|
await projectAdmissionCoordinator.admitNext({
|
||||||
projectId: this.store.getRootDir(),
|
projectId: this.store.getRootDir(),
|
||||||
taskId: task.id,
|
|
||||||
maxConcurrent: activeTaskLimit,
|
maxConcurrent: activeTaskLimit,
|
||||||
claimed: async () => (await getFinalClaimSnapshot()).count,
|
claimed: async () => (await getFinalClaimSnapshot()).count,
|
||||||
claimedTaskIds: async () => (await getFinalClaimSnapshot()).ids,
|
claimedTaskIds: async () => (await getFinalClaimSnapshot()).ids,
|
||||||
|
semaphore: this.options.semaphore,
|
||||||
|
refresh: async () => [{
|
||||||
|
taskId: task.id,
|
||||||
|
projectId: this.store.getRootDir(),
|
||||||
|
lane: "execute",
|
||||||
|
createdAt: task.createdAt,
|
||||||
|
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
|
||||||
|
start: async () => {
|
||||||
|
projectSlotReserved = true;
|
||||||
|
return true;
|
||||||
|
},
|
||||||
|
}],
|
||||||
});
|
});
|
||||||
if (!projectSlotReserved) {
|
if (!projectSlotReserved) {
|
||||||
if (reservedScope) {
|
if (reservedScope) {
|
||||||
@@ -2850,27 +2862,9 @@ export class Scheduler {
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// admitNext acquired and registered the host slot atomically with the
|
||||||
|
// project reservation, so this handoff cannot bypass review candidates.
|
||||||
const sem = this.options.semaphore;
|
const sem = this.options.semaphore;
|
||||||
const hostSlotReserved = hasPreHeldExecutorSlot(task.id);
|
|
||||||
registerPreHeldExecutorSlot(task.id, hostSlotReserved);
|
|
||||||
if (sem && !hostSlotReserved && !sem.tryAcquire()) {
|
|
||||||
dropPreHeldExecutorSlot(task.id);
|
|
||||||
if (reservedScope) {
|
|
||||||
activeScopes.delete(task.id);
|
|
||||||
activeScopeColumns.delete(task.id);
|
|
||||||
}
|
|
||||||
const reason = formatConcurrencyLimitReason({
|
|
||||||
...concurrencyDiagnostic,
|
|
||||||
available: 0,
|
|
||||||
bindingGates: [...new Set([...concurrencyDiagnostic.bindingGates, "semaphore" as const])],
|
|
||||||
});
|
|
||||||
await this.store.updateTask(task.id, { status: "queued" });
|
|
||||||
await this.logDispatchQueuedReason(task.id, reason, formatConcurrencyLimitMemoKey(concurrencyDiagnostic));
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
if (sem && !hostSlotReserved) {
|
|
||||||
registerPreHeldExecutorSlot(task.id);
|
|
||||||
}
|
|
||||||
|
|
||||||
let acquiredSymbols: string[] | undefined;
|
let acquiredSymbols: string[] | undefined;
|
||||||
try {
|
try {
|
||||||
|
|||||||
@@ -595,7 +595,7 @@ export class TriageProcessor {
|
|||||||
now,
|
now,
|
||||||
);
|
);
|
||||||
return tasks.filter((task) => !this.coordinatorAdmittedTaskIds.has(task.id)).map((task) => ({
|
return tasks.filter((task) => !this.coordinatorAdmittedTaskIds.has(task.id)).map((task) => ({
|
||||||
taskId: task.id, projectId: this.rootDir, createdAt: task.createdAt,
|
taskId: task.id, projectId: this.rootDir, lane: "planning", createdAt: task.createdAt,
|
||||||
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
|
reserve: () => registerPreHeldExecutorSlot(task.id, this.options.semaphore !== undefined),
|
||||||
start: async () => {
|
start: async () => {
|
||||||
this.coordinatorAdmittedTaskIds.add(task.id);
|
this.coordinatorAdmittedTaskIds.add(task.id);
|
||||||
@@ -1231,7 +1231,7 @@ export class TriageProcessor {
|
|||||||
duplicate-claim guard), so a promise that never settles — exactly the case this eviction
|
duplicate-claim guard), so a promise that never settles — exactly the case this eviction
|
||||||
exists for — left the id in the set permanently. Planning discovery does not consult that
|
exists for — left the id in the set permanently. Planning discovery does not consult that
|
||||||
set, so the card stayed in `triageTasks` and `maxToStart` stayed positive, which means the
|
set, so the card stayed in `triageTasks` and `maxToStart` stayed positive, which means the
|
||||||
throttle branch (the only thing that logs or emits) never fired; but `admitOldest`'s
|
throttle branch (the only thing that logs or emits) never fired; but `admitNext`'s
|
||||||
`refresh()` filters on the set, so the coordinator saw no candidate. Silent stall until
|
`refresh()` filters on the set, so the coordinator saw no candidate. Silent stall until
|
||||||
engine restart, and the badge (a pure client-side "unplanned + idle in Todo" inference) kept
|
engine restart, and the badge (a pure client-side "unplanned + idle in Todo" inference) kept
|
||||||
claiming the card was queued.
|
claiming the card was queued.
|
||||||
@@ -2237,7 +2237,7 @@ export class TriageProcessor {
|
|||||||
const ids = await persistedTopLevelAgentTaskIdsFromStore(this.store, fresh);
|
const ids = await persistedTopLevelAgentTaskIdsFromStore(this.store, fresh);
|
||||||
return { count: ids.length + pending, ids: [...new Set([...ids, ...this.processing])] };
|
return { count: ids.length + pending, ids: [...new Set([...ids, ...this.processing])] };
|
||||||
})();
|
})();
|
||||||
await projectAdmissionCoordinator.admitOldest({
|
await projectAdmissionCoordinator.admitNext({
|
||||||
// rootDir is the stable per-project identity held by this processor.
|
// rootDir is the stable per-project identity held by this processor.
|
||||||
projectId: this.rootDir,
|
projectId: this.rootDir,
|
||||||
maxConcurrent: activeTaskLimit,
|
maxConcurrent: activeTaskLimit,
|
||||||
@@ -2249,6 +2249,7 @@ export class TriageProcessor {
|
|||||||
.map((task) => ({
|
.map((task) => ({
|
||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
projectId: this.rootDir,
|
projectId: this.rootDir,
|
||||||
|
lane: "planning",
|
||||||
createdAt: task.createdAt,
|
createdAt: task.createdAt,
|
||||||
// FNXC:ConcurrencyAdmission 2026-08-05-10:00: the planner must
|
// FNXC:ConcurrencyAdmission 2026-08-05-10:00: the planner must
|
||||||
// own the coordinator's real host reservation before it starts;
|
// own the coordinator's real host reservation before it starts;
|
||||||
|
|||||||
Reference in New Issue
Block a user