feat(KB-138): add priority-based scheduling to AgentSemaphore
- Add priority levels (PRIORITY_MERGE=2, PRIORITY_EXECUTE=1, PRIORITY_SPECIFY=0) to AgentSemaphore - Update acquire() and run() to accept a numeric priority parameter with FIFO ordering within same level - Wire priority constants into executor (PRIORITY_EXECUTE) and triage (PRIORITY_SPECIFY) callers - Add comprehensive tests for priority ordering, FIFO within same priority, and dynamic limit interaction - Include changeset for the priority-based agent scheduling feature
This commit is contained in:
5
.changeset/priority-agent-scheduling.md
Normal file
5
.changeset/priority-agent-scheduling.md
Normal file
@@ -0,0 +1,5 @@
|
|||||||
|
---
|
||||||
|
"@dustinbyrne/kb": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
Add priority-based agent scheduling: merge agents are served before execution agents, which are served before specification agents, when competing for concurrency slots.
|
||||||
@@ -2,7 +2,7 @@ import { exec } from "node:child_process";
|
|||||||
import type { AddressInfo } from "node:net";
|
import type { AddressInfo } from "node:net";
|
||||||
import { TaskStore } from "@kb/core";
|
import { TaskStore } from "@kb/core";
|
||||||
import { createServer } from "@kb/dashboard";
|
import { createServer } from "@kb/dashboard";
|
||||||
import { TriageProcessor, TaskExecutor, Scheduler, AgentSemaphore, WorktreePool, aiMergeTask } from "@kb/engine";
|
import { TriageProcessor, TaskExecutor, Scheduler, AgentSemaphore, WorktreePool, aiMergeTask, PRIORITY_MERGE } from "@kb/engine";
|
||||||
import { AuthStorage, ModelRegistry } from "@mariozechner/pi-coding-agent";
|
import { AuthStorage, ModelRegistry } from "@mariozechner/pi-coding-agent";
|
||||||
|
|
||||||
function openBrowser(url: string): void {
|
function openBrowser(url: string): void {
|
||||||
@@ -57,7 +57,7 @@ export async function runDashboard(port: number, opts: { open?: boolean } = {})
|
|||||||
onAgentTool: (name) => console.log(`[merger] tool: ${name}`),
|
onAgentTool: (name) => console.log(`[merger] tool: ${name}`),
|
||||||
});
|
});
|
||||||
|
|
||||||
const onMerge = (taskId: string) => semaphore.run(() => rawMerge(taskId));
|
const onMerge = (taskId: string) => semaphore.run(() => rawMerge(taskId), PRIORITY_MERGE);
|
||||||
|
|
||||||
// ── Serialized auto-merge queue ─────────────────────────────────────
|
// ── Serialized auto-merge queue ─────────────────────────────────────
|
||||||
//
|
//
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
import { describe, it, expect, vi } from "vitest";
|
import { describe, it, expect, vi } from "vitest";
|
||||||
import { AgentSemaphore } from "./concurrency.js";
|
import { AgentSemaphore, PRIORITY_MERGE, PRIORITY_EXECUTE, PRIORITY_SPECIFY } from "./concurrency.js";
|
||||||
|
|
||||||
describe("AgentSemaphore", () => {
|
describe("AgentSemaphore", () => {
|
||||||
it("allows immediate acquire when under limit", async () => {
|
it("allows immediate acquire when under limit", async () => {
|
||||||
@@ -243,4 +243,141 @@ describe("AgentSemaphore", () => {
|
|||||||
|
|
||||||
expect(ran).toBe(true);
|
expect(ran).toBe(true);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// ── Priority scheduling tests ──────────────────────────────────────
|
||||||
|
|
||||||
|
it("priority: highest-priority waiter is served first when slot is released", async () => {
|
||||||
|
const sem = new AgentSemaphore(1);
|
||||||
|
await sem.acquire(); // fill the single slot
|
||||||
|
|
||||||
|
const order: string[] = [];
|
||||||
|
|
||||||
|
// Queue three waiters in non-priority order: specify, merge, execute
|
||||||
|
const pSpecify = sem.acquire(PRIORITY_SPECIFY).then(() => order.push("specify"));
|
||||||
|
const pMerge = sem.acquire(PRIORITY_MERGE).then(() => order.push("merge"));
|
||||||
|
const pExecute = sem.acquire(PRIORITY_EXECUTE).then(() => order.push("execute"));
|
||||||
|
|
||||||
|
// Release slots one at a time and observe drain order
|
||||||
|
sem.release();
|
||||||
|
await pMerge;
|
||||||
|
expect(order).toEqual(["merge"]);
|
||||||
|
|
||||||
|
sem.release();
|
||||||
|
await pExecute;
|
||||||
|
expect(order).toEqual(["merge", "execute"]);
|
||||||
|
|
||||||
|
sem.release();
|
||||||
|
await pSpecify;
|
||||||
|
expect(order).toEqual(["merge", "execute", "specify"]);
|
||||||
|
|
||||||
|
sem.release(); // cleanup
|
||||||
|
expect(sem.activeCount).toBe(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("priority: FIFO order is preserved among equal-priority waiters", async () => {
|
||||||
|
const sem = new AgentSemaphore(1);
|
||||||
|
await sem.acquire(); // fill the slot
|
||||||
|
|
||||||
|
const order: number[] = [];
|
||||||
|
|
||||||
|
const p1 = sem.acquire(PRIORITY_EXECUTE).then(() => order.push(1));
|
||||||
|
const p2 = sem.acquire(PRIORITY_EXECUTE).then(() => order.push(2));
|
||||||
|
const p3 = sem.acquire(PRIORITY_EXECUTE).then(() => order.push(3));
|
||||||
|
|
||||||
|
sem.release();
|
||||||
|
await p1;
|
||||||
|
sem.release();
|
||||||
|
await p2;
|
||||||
|
sem.release();
|
||||||
|
await p3;
|
||||||
|
|
||||||
|
expect(order).toEqual([1, 2, 3]);
|
||||||
|
sem.release();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("priority: run() forwards priority to acquire()", async () => {
|
||||||
|
const sem = new AgentSemaphore(1);
|
||||||
|
await sem.acquire(); // fill the slot
|
||||||
|
|
||||||
|
const order: string[] = [];
|
||||||
|
|
||||||
|
const pLow = sem.run(async () => { order.push("low"); }, PRIORITY_SPECIFY);
|
||||||
|
const pHigh = sem.run(async () => { order.push("high"); }, PRIORITY_MERGE);
|
||||||
|
|
||||||
|
// Release — high priority should go first, then its run() releases the
|
||||||
|
// slot automatically, allowing the low-priority waiter to proceed.
|
||||||
|
sem.release();
|
||||||
|
await pHigh;
|
||||||
|
await pLow;
|
||||||
|
expect(order).toEqual(["high", "low"]);
|
||||||
|
expect(sem.activeCount).toBe(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("priority: mixed-priority integration — 1 slot, arbitrary enqueue order, correct drain", async () => {
|
||||||
|
const sem = new AgentSemaphore(1);
|
||||||
|
await sem.acquire(); // hold the single slot
|
||||||
|
|
||||||
|
const order: string[] = [];
|
||||||
|
|
||||||
|
// Enqueue in a scrambled order: execute, specify, merge, specify, execute, merge
|
||||||
|
const promises = [
|
||||||
|
sem.acquire(PRIORITY_EXECUTE).then(() => { order.push("execute-1"); }),
|
||||||
|
sem.acquire(PRIORITY_SPECIFY).then(() => { order.push("specify-1"); }),
|
||||||
|
sem.acquire(PRIORITY_MERGE).then(() => { order.push("merge-1"); }),
|
||||||
|
sem.acquire(PRIORITY_SPECIFY).then(() => { order.push("specify-2"); }),
|
||||||
|
sem.acquire(PRIORITY_EXECUTE).then(() => { order.push("execute-2"); }),
|
||||||
|
sem.acquire(PRIORITY_MERGE).then(() => { order.push("merge-2"); }),
|
||||||
|
];
|
||||||
|
|
||||||
|
// Expected drain order:
|
||||||
|
// merge-1, merge-2 (highest, FIFO within),
|
||||||
|
// execute-1, execute-2 (middle, FIFO within),
|
||||||
|
// specify-1, specify-2 (lowest, FIFO within)
|
||||||
|
for (let i = 0; i < 6; i++) {
|
||||||
|
sem.release();
|
||||||
|
// Wait for the next promise to settle
|
||||||
|
await Promise.resolve();
|
||||||
|
await Promise.resolve();
|
||||||
|
}
|
||||||
|
|
||||||
|
await Promise.all(promises);
|
||||||
|
|
||||||
|
expect(order).toEqual([
|
||||||
|
"merge-1", "merge-2",
|
||||||
|
"execute-1", "execute-2",
|
||||||
|
"specify-1", "specify-2",
|
||||||
|
]);
|
||||||
|
|
||||||
|
sem.release(); // cleanup
|
||||||
|
expect(sem.activeCount).toBe(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("priority constants have correct values", () => {
|
||||||
|
expect(PRIORITY_MERGE).toBe(2);
|
||||||
|
expect(PRIORITY_EXECUTE).toBe(1);
|
||||||
|
expect(PRIORITY_SPECIFY).toBe(0);
|
||||||
|
expect(PRIORITY_MERGE).toBeGreaterThan(PRIORITY_EXECUTE);
|
||||||
|
expect(PRIORITY_EXECUTE).toBeGreaterThan(PRIORITY_SPECIFY);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("priority: default priority (no argument) behaves as PRIORITY_SPECIFY (0)", async () => {
|
||||||
|
const sem = new AgentSemaphore(1);
|
||||||
|
await sem.acquire(); // fill the slot
|
||||||
|
|
||||||
|
const order: string[] = [];
|
||||||
|
|
||||||
|
// acquire() with no priority arg — should be treated as 0
|
||||||
|
const pDefault = sem.acquire().then(() => order.push("default"));
|
||||||
|
const pMerge = sem.acquire(PRIORITY_MERGE).then(() => order.push("merge"));
|
||||||
|
|
||||||
|
sem.release();
|
||||||
|
await pMerge;
|
||||||
|
expect(order).toEqual(["merge"]);
|
||||||
|
|
||||||
|
sem.release();
|
||||||
|
await pDefault;
|
||||||
|
expect(order).toEqual(["merge", "default"]);
|
||||||
|
|
||||||
|
sem.release();
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,3 +1,16 @@
|
|||||||
|
/** Priority level for merge agents — served first. */
|
||||||
|
export const PRIORITY_MERGE = 2;
|
||||||
|
/** Priority level for execution agents — served after merge, before specify. */
|
||||||
|
export const PRIORITY_EXECUTE = 1;
|
||||||
|
/** Priority level for specification/triage agents — served last (default). */
|
||||||
|
export const PRIORITY_SPECIFY = 0;
|
||||||
|
|
||||||
|
/** A waiter entry that tracks both the priority and the resolve callback. */
|
||||||
|
interface PriorityWaiter {
|
||||||
|
priority: number;
|
||||||
|
resolve: () => void;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* A concurrency semaphore that gates all agentic activities (triage specification,
|
* A concurrency semaphore that gates all agentic activities (triage specification,
|
||||||
* task execution, and merge operations) behind a shared slot limit.
|
* task execution, and merge operations) behind a shared slot limit.
|
||||||
@@ -5,6 +18,15 @@
|
|||||||
* The semaphore ensures that the total number of concurrently running AI agents
|
* The semaphore ensures that the total number of concurrently running AI agents
|
||||||
* never exceeds `maxConcurrent`, regardless of which subsystem spawned them.
|
* never exceeds `maxConcurrent`, regardless of which subsystem spawned them.
|
||||||
*
|
*
|
||||||
|
* **Priority-based draining:** When a slot becomes available and multiple agents
|
||||||
|
* are waiting, the waiter with the highest `priority` value is served first.
|
||||||
|
* Among waiters with the same priority, FIFO order is preserved. The built-in
|
||||||
|
* priority constants are:
|
||||||
|
*
|
||||||
|
* - {@link PRIORITY_MERGE} (`2`) — merge agents (highest)
|
||||||
|
* - {@link PRIORITY_EXECUTE} (`1`) — execution agents
|
||||||
|
* - {@link PRIORITY_SPECIFY} (`0`) — specification/triage agents (lowest, default)
|
||||||
|
*
|
||||||
* The limit is read dynamically at `acquire()` time via a getter callback, so
|
* The limit is read dynamically at `acquire()` time via a getter callback, so
|
||||||
* live changes to `settings.maxConcurrent` take effect on the next acquire
|
* live changes to `settings.maxConcurrent` take effect on the next acquire
|
||||||
* without restarting the engine. Reducing the limit below the current
|
* without restarting the engine. Reducing the limit below the current
|
||||||
@@ -16,12 +38,12 @@
|
|||||||
* const sem = new AgentSemaphore(() => store.getSettings().then(s => s.maxConcurrent));
|
* const sem = new AgentSemaphore(() => store.getSettings().then(s => s.maxConcurrent));
|
||||||
* await sem.run(async () => {
|
* await sem.run(async () => {
|
||||||
* // at most maxConcurrent agents run this block concurrently
|
* // at most maxConcurrent agents run this block concurrently
|
||||||
* });
|
* }, PRIORITY_EXECUTE);
|
||||||
* ```
|
* ```
|
||||||
*/
|
*/
|
||||||
export class AgentSemaphore {
|
export class AgentSemaphore {
|
||||||
private _active = 0;
|
private _active = 0;
|
||||||
private _waiters: Array<() => void> = [];
|
private _waiters: PriorityWaiter[] = [];
|
||||||
private _getLimit: () => number;
|
private _getLimit: () => number;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -51,17 +73,27 @@ export class AgentSemaphore {
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* Acquire a slot. Resolves immediately if a slot is available, otherwise
|
* Acquire a slot. Resolves immediately if a slot is available, otherwise
|
||||||
* queues the caller and resolves in FIFO order when a slot is released.
|
* queues the caller and resolves when a slot is released.
|
||||||
|
*
|
||||||
|
* When multiple callers are waiting, the highest-priority waiter is served
|
||||||
|
* first. Among waiters with equal priority, FIFO order is preserved.
|
||||||
|
*
|
||||||
|
* @param priority - Numeric priority (higher = served first). Defaults to `0`
|
||||||
|
* ({@link PRIORITY_SPECIFY}). Use {@link PRIORITY_MERGE} (`2`) for merge
|
||||||
|
* agents and {@link PRIORITY_EXECUTE} (`1`) for execution agents.
|
||||||
*/
|
*/
|
||||||
acquire(): Promise<void> {
|
acquire(priority: number = 0): Promise<void> {
|
||||||
if (this._active < this._getLimit()) {
|
if (this._active < this._getLimit()) {
|
||||||
this._active++;
|
this._active++;
|
||||||
return Promise.resolve();
|
return Promise.resolve();
|
||||||
}
|
}
|
||||||
return new Promise<void>((resolve) => {
|
return new Promise<void>((resolve) => {
|
||||||
this._waiters.push(() => {
|
this._waiters.push({
|
||||||
this._active++;
|
priority,
|
||||||
resolve();
|
resolve: () => {
|
||||||
|
this._active++;
|
||||||
|
resolve();
|
||||||
|
},
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
@@ -78,9 +110,13 @@ export class AgentSemaphore {
|
|||||||
/**
|
/**
|
||||||
* Convenience wrapper: acquires a slot, runs `fn`, and releases the slot
|
* Convenience wrapper: acquires a slot, runs `fn`, and releases the slot
|
||||||
* when `fn` settles (whether it resolves or rejects).
|
* when `fn` settles (whether it resolves or rejects).
|
||||||
|
*
|
||||||
|
* @param fn - The async function to run while holding the slot.
|
||||||
|
* @param priority - Numeric priority forwarded to {@link acquire}. Defaults
|
||||||
|
* to `0` ({@link PRIORITY_SPECIFY}).
|
||||||
*/
|
*/
|
||||||
async run<T>(fn: () => Promise<T>): Promise<T> {
|
async run<T>(fn: () => Promise<T>, priority: number = 0): Promise<T> {
|
||||||
await this.acquire();
|
await this.acquire(priority);
|
||||||
try {
|
try {
|
||||||
return await fn();
|
return await fn();
|
||||||
} finally {
|
} finally {
|
||||||
@@ -88,11 +124,34 @@ export class AgentSemaphore {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Unblock waiters while slots are available. */
|
/**
|
||||||
|
* Unblock waiters while slots are available.
|
||||||
|
*
|
||||||
|
* Picks the highest-priority waiter first. Among waiters with the same
|
||||||
|
* priority, the one that was enqueued first (FIFO) is chosen.
|
||||||
|
*/
|
||||||
private _drain(): void {
|
private _drain(): void {
|
||||||
while (this._waiters.length > 0 && this._active < this._getLimit()) {
|
while (this._waiters.length > 0 && this._active < this._getLimit()) {
|
||||||
const next = this._waiters.shift()!;
|
const idx = this._highestPriorityIndex();
|
||||||
next();
|
const [waiter] = this._waiters.splice(idx, 1);
|
||||||
|
waiter.resolve();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Find the index of the highest-priority waiter. When multiple waiters
|
||||||
|
* share the highest priority, the first one (lowest index = earliest
|
||||||
|
* enqueued) is returned, preserving FIFO within the same priority level.
|
||||||
|
*/
|
||||||
|
private _highestPriorityIndex(): number {
|
||||||
|
let bestIdx = 0;
|
||||||
|
let bestPriority = this._waiters[0].priority;
|
||||||
|
for (let i = 1; i < this._waiters.length; i++) {
|
||||||
|
if (this._waiters[i].priority > bestPriority) {
|
||||||
|
bestPriority = this._waiters[i].priority;
|
||||||
|
bestIdx = i;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return bestIdx;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,7 +8,7 @@ import { Type, type Static } from "@mariozechner/pi-ai";
|
|||||||
import { createKbAgent } from "./pi.js";
|
import { createKbAgent } from "./pi.js";
|
||||||
import { reviewStep, type ReviewVerdict } from "./reviewer.js";
|
import { reviewStep, type ReviewVerdict } from "./reviewer.js";
|
||||||
import type { ToolDefinition, AgentSession, SessionManager } from "@mariozechner/pi-coding-agent";
|
import type { ToolDefinition, AgentSession, SessionManager } from "@mariozechner/pi-coding-agent";
|
||||||
import type { AgentSemaphore } from "./concurrency.js";
|
import { PRIORITY_EXECUTE, type AgentSemaphore } from "./concurrency.js";
|
||||||
import type { WorktreePool } from "./worktree-pool.js";
|
import type { WorktreePool } from "./worktree-pool.js";
|
||||||
import { AgentLogger } from "./agent-logger.js";
|
import { AgentLogger } from "./agent-logger.js";
|
||||||
import { executorLog, reviewerLog } from "./logger.js";
|
import { executorLog, reviewerLog } from "./logger.js";
|
||||||
@@ -413,7 +413,7 @@ export class TaskExecutor {
|
|||||||
};
|
};
|
||||||
|
|
||||||
if (this.options.semaphore) {
|
if (this.options.semaphore) {
|
||||||
await this.options.semaphore.run(agentWork);
|
await this.options.semaphore.run(agentWork, PRIORITY_EXECUTE);
|
||||||
} else {
|
} else {
|
||||||
await agentWork();
|
await agentWork();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
export { AgentLogger, type AgentLoggerOptions, summarizeToolArgs } from "./agent-logger.js";
|
export { AgentLogger, type AgentLoggerOptions, summarizeToolArgs } from "./agent-logger.js";
|
||||||
export { AgentSemaphore } from "./concurrency.js";
|
export { AgentSemaphore, PRIORITY_MERGE, PRIORITY_EXECUTE, PRIORITY_SPECIFY } from "./concurrency.js";
|
||||||
export { TriageProcessor, type TriageProcessorOptions } from "./triage.js";
|
export { TriageProcessor, type TriageProcessorOptions } from "./triage.js";
|
||||||
export { TaskExecutor, type TaskExecutorOptions } from "./executor.js";
|
export { TaskExecutor, type TaskExecutorOptions } from "./executor.js";
|
||||||
export { Scheduler, type SchedulerOptions } from "./scheduler.js";
|
export { Scheduler, type SchedulerOptions } from "./scheduler.js";
|
||||||
|
|||||||
@@ -3,7 +3,7 @@ import type { ImageContent } from "@mariozechner/pi-ai";
|
|||||||
import { Type, type Static } from "@mariozechner/pi-ai";
|
import { Type, type Static } from "@mariozechner/pi-ai";
|
||||||
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
|
import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
|
||||||
import { createKbAgent } from "./pi.js";
|
import { createKbAgent } from "./pi.js";
|
||||||
import type { AgentSemaphore } from "./concurrency.js";
|
import { PRIORITY_SPECIFY, type AgentSemaphore } from "./concurrency.js";
|
||||||
import { AgentLogger } from "./agent-logger.js";
|
import { AgentLogger } from "./agent-logger.js";
|
||||||
import { triageLog } from "./logger.js";
|
import { triageLog } from "./logger.js";
|
||||||
|
|
||||||
@@ -314,7 +314,7 @@ export class TriageProcessor {
|
|||||||
};
|
};
|
||||||
|
|
||||||
if (this.options.semaphore) {
|
if (this.options.semaphore) {
|
||||||
await this.options.semaphore.run(agentWork);
|
await this.options.semaphore.run(agentWork, PRIORITY_SPECIFY);
|
||||||
} else {
|
} else {
|
||||||
await agentWork();
|
await agentWork();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user