summaryrefslogtreecommitdiffhomepage
path: root/packages/provider-concurrency/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-27 03:03:53 +0900
committerAdam Malczewski <[email protected]>2026-06-27 03:03:53 +0900
commita6b95188a110464b6ffa0334c8af58463f2a36f2 (patch)
treeeb6ef57909e164be4ae721ea1fb25585354d351e /packages/provider-concurrency/src
parentad9d135e583c99a0d93327115defa43187cde1c3 (diff)
downloaddispatch-a6b95188a110464b6ffa0334c8af58463f2a36f2.tar.gz
dispatch-a6b95188a110464b6ffa0334c8af58463f2a36f2.zip
feat(provider-concurrency): implement per-provider in-memory concurrency limits with oldest-agent-first scheduling
Diffstat (limited to 'packages/provider-concurrency/src')
-rw-r--r--packages/provider-concurrency/src/concurrency-manager.test.ts386
-rw-r--r--packages/provider-concurrency/src/concurrency-manager.ts327
-rw-r--r--packages/provider-concurrency/src/extension.ts61
-rw-r--r--packages/provider-concurrency/src/index.ts10
-rw-r--r--packages/provider-concurrency/src/provider-wrapper.test.ts173
-rw-r--r--packages/provider-concurrency/src/provider-wrapper.ts59
-rw-r--r--packages/provider-concurrency/src/service.ts11
7 files changed, 1027 insertions, 0 deletions
diff --git a/packages/provider-concurrency/src/concurrency-manager.test.ts b/packages/provider-concurrency/src/concurrency-manager.test.ts
new file mode 100644
index 0000000..bfefb4a
--- /dev/null
+++ b/packages/provider-concurrency/src/concurrency-manager.test.ts
@@ -0,0 +1,386 @@
+import { describe, expect, it } from "vitest";
+import { type ConcurrencyService, createConcurrencyManager } from "./concurrency-manager.js";
+
+// ─── Fake timers ──────────────────────────────────────────────────────────────
+
+interface FakeTimer {
+ fire: () => void;
+ cleared: boolean;
+}
+
+function createFakeTimers() {
+ let currentTime = 0;
+ const intervals: FakeTimer[] = [];
+ const timeouts: { time: number; fire: () => void; cleared: boolean }[] = [];
+
+ const setInterval = ((_fn: () => void, _ms: number) => {
+ const timer: FakeTimer = { fire: () => _fn(), cleared: false };
+ intervals.push(timer);
+ return timer as unknown as ReturnType<typeof setInterval>;
+ }) as typeof setInterval;
+
+ const clearInterval = ((timer: ReturnType<typeof setInterval>) => {
+ const t = timer as unknown as FakeTimer;
+ t.cleared = true;
+ }) as typeof clearInterval;
+
+ const setTimeout = ((_fn: () => void, ms: number) => {
+ const entry = { time: currentTime + ms, fire: () => _fn(), cleared: false };
+ timeouts.push(entry);
+ return entry as unknown as ReturnType<typeof setTimeout>;
+ }) as typeof setTimeout;
+
+ const clearTimeout = ((timer: ReturnType<typeof setTimeout>) => {
+ const t = timer as unknown as { cleared: boolean };
+ t.cleared = true;
+ }) as typeof clearTimeout;
+
+ return {
+ now: () => currentTime,
+ advance(ms: number) {
+ currentTime += ms;
+ // Fire any due timeouts.
+ for (const entry of timeouts) {
+ if (!entry.cleared && entry.time <= currentTime) {
+ entry.cleared = true;
+ entry.fire();
+ }
+ }
+ },
+ fireIntervals() {
+ for (const timer of intervals) {
+ if (!timer.cleared) timer.fire();
+ }
+ },
+ setInterval,
+ clearInterval,
+ setTimeout,
+ clearTimeout,
+ };
+}
+
+function createManager(): {
+ manager: ConcurrencyService;
+ timers: ReturnType<typeof createFakeTimers>;
+} {
+ const timers = createFakeTimers();
+ const manager = createConcurrencyManager({
+ now: timers.now,
+ slotTimeoutMs: 5000,
+ watchdogIntervalMs: 1000,
+ defaultPauseMs: 30000,
+ setTimeout: timers.setTimeout,
+ clearTimeout: timers.clearTimeout,
+ setInterval: timers.setInterval,
+ clearInterval: timers.clearInterval,
+ });
+ return { manager, timers };
+}
+
+describe("createConcurrencyManager", () => {
+ it("returns no-op release for providers with no configured limit", async () => {
+ const { manager } = createManager();
+ const release = await manager.acquire("unknown", "conv1", 0);
+ expect(typeof release).toBe("function");
+ // No state → release is a no-op, no error.
+ release();
+ expect(manager.getStatus("unknown")).toBeUndefined();
+ });
+
+ it("grants immediately when under the limit", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 4);
+
+ const release1 = await manager.acquire("umans", "conv1", 0);
+ const status = manager.getStatus("umans");
+ expect(status).toEqual({
+ providerId: "umans",
+ limit: 4,
+ inFlight: 1,
+ queued: 0,
+ paused: false,
+ });
+ release1();
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+ });
+
+ it("queues when at the limit and grants on release (FIFO when same priority)", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 1);
+
+ const release1 = await manager.acquire("umans", "conv1", 100);
+
+ // Second request should block (at limit).
+ let resolved = false;
+ const promise2 = manager.acquire("umans", "conv2", 200).then((r) => {
+ resolved = true;
+ return r;
+ });
+
+ // Let microtasks settle.
+ await Promise.resolve();
+ await Promise.resolve();
+ expect(resolved).toBe(false);
+ expect(manager.getStatus("umans")?.queued).toBe(1);
+
+ // Release the first slot.
+ release1();
+
+ const release2 = await promise2;
+ expect(resolved).toBe(true);
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+ expect(manager.getStatus("umans")?.queued).toBe(0);
+ release2();
+ });
+
+ it("grants to the oldest agent first (priority queue by promptStartedAt)", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 1);
+
+ // Hold the single slot.
+ const release0 = await manager.acquire("umans", "holder", 0);
+
+ // Three agents queue with different prompt start times.
+ // Agent C started latest (t=300), Agent A started earliest (t=100).
+ const results: string[] = [];
+ const acquireAndRecord = (conv: string, promptAt: number) =>
+ manager.acquire("umans", conv, promptAt).then((r) => {
+ results.push(conv);
+ return r;
+ });
+
+ // Queue in non-sorted order: B (t=200), A (t=100), C (t=300).
+ const pB = acquireAndRecord("convB", 200);
+ const pA = acquireAndRecord("convA", 100);
+ const pC = acquireAndRecord("convC", 300);
+
+ await Promise.resolve();
+ await Promise.resolve();
+ expect(results).toEqual([]); // none resolved yet.
+
+ // Release the holder. The oldest agent (A, t=100) should get the slot first.
+ release0();
+
+ const rA = await pA;
+ expect(results).toEqual(["convA"]);
+
+ rA.release ? rA.release() : rA();
+
+ // Now B (t=200) should be next.
+ const rB = await pB;
+ expect(results).toEqual(["convA", "convB"]);
+ rB.release ? rB.release() : rB();
+
+ // Then C (t=300).
+ const rC = await pC;
+ expect(results).toEqual(["convA", "convB", "convC"]);
+ rC.release ? rC.release() : rC();
+ });
+
+ it("does not grant slots while paused (429 backoff)", async () => {
+ const { manager, timers } = createManager();
+ manager.setLimit("umans", 1);
+
+ const release1 = await manager.acquire("umans", "conv1", 0);
+ release1();
+
+ // Simulate a 429 → queue pauses.
+ manager.reportRateLimit("umans");
+ const status = manager.getStatus("umans");
+ expect(status?.paused).toBe(true);
+ expect(status?.pausedUntil).toBe(30000);
+
+ // A new acquire should block (paused, even though under limit).
+ let resolved = false;
+ const promise = manager.acquire("umans", "conv2", 0).then((r) => {
+ resolved = true;
+ return r;
+ });
+ await Promise.resolve();
+ await Promise.resolve();
+ expect(resolved).toBe(false);
+
+ // Advance past the pause duration.
+ timers.advance(30000);
+
+ const release2 = await promise;
+ expect(resolved).toBe(true);
+ expect(manager.getStatus("umans")?.paused).toBe(false);
+ release2();
+ });
+
+ it("respects retryAfterMs for 429 backoff", () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 2);
+
+ manager.reportRateLimit("umans", 5000);
+ expect(manager.getStatus("umans")?.pausedUntil).toBe(5000);
+ });
+
+ it("watchdog reclaims slots held beyond the timeout", async () => {
+ const { manager, timers } = createManager();
+ manager.setLimit("umans", 1);
+
+ const release = await manager.acquire("umans", "conv1", 0);
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+
+ // Advance past the slot timeout (5000ms) and fire the watchdog.
+ timers.advance(5001);
+ timers.fireIntervals();
+
+ // The watchdog should have force-released the slot.
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+
+ // Calling release again (from the holder) should be a no-op (idempotent).
+ release();
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+ });
+
+ it("watchdog grants the next waiter after reclaiming a stale slot", async () => {
+ const { manager, timers } = createManager();
+ manager.setLimit("umans", 1);
+
+ // Hold the slot.
+ await manager.acquire("umans", "holder", 0);
+
+ // Queue a waiter.
+ let resolved = false;
+ const promise = manager.acquire("umans", "waiter", 10).then((r) => {
+ resolved = true;
+ return r;
+ });
+ await Promise.resolve();
+ await Promise.resolve();
+ expect(resolved).toBe(false);
+
+ // Watchdog reclaims the held slot.
+ timers.advance(5001);
+ timers.fireIntervals();
+
+ // The waiter should now be granted.
+ const release = await promise;
+ expect(resolved).toBe(true);
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+ release();
+ });
+
+ it("setLimit grants queued requests when the limit increases", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 1);
+
+ const release1 = await manager.acquire("umans", "conv1", 0);
+
+ // Queue a waiter.
+ let resolved = false;
+ const promise = manager.acquire("umans", "conv2", 100).then((r) => {
+ resolved = true;
+ return r;
+ });
+ await Promise.resolve();
+ await Promise.resolve();
+ expect(resolved).toBe(false);
+
+ // Increase the limit → the queued request should be granted.
+ manager.setLimit("umans", 2);
+
+ const release2 = await promise;
+ expect(resolved).toBe(true);
+ expect(manager.getStatus("umans")?.inFlight).toBe(2);
+
+ release2();
+ release1();
+ });
+
+ it("removeLimit grants all queued requests and removes the state", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 1);
+
+ const release1 = await manager.acquire("umans", "conv1", 0);
+
+ // Queue two waiters.
+ const p2 = manager.acquire("umans", "conv2", 100);
+ const p3 = manager.acquire("umans", "conv3", 200);
+ await Promise.resolve();
+ await Promise.resolve();
+
+ // Remove the limit → all queued requests should be granted.
+ manager.removeLimit("umans");
+
+ const r2 = await p2;
+ const r3 = await p3;
+ expect(manager.getStatus("umans")).toBeUndefined();
+
+ // Releases work (no error after state removal).
+ r2();
+ r3();
+ release1();
+ });
+
+ it("getLimits returns all configured limits", () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 4);
+ manager.setLimit("openai-compat", 5);
+
+ const limits = manager.getLimits();
+ expect(limits).toHaveLength(2);
+ expect(limits).toContainEqual({ providerId: "umans", limit: 4 });
+ expect(limits).toContainEqual({ providerId: "openai-compat", limit: 5 });
+ });
+
+ it("getStatusAll returns status for all configured providers", () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 4);
+ manager.setLimit("anthropic", 3);
+
+ const statuses = manager.getStatusAll();
+ expect(statuses).toHaveLength(2);
+ const umans = statuses.find((s) => s.providerId === "umans");
+ expect(umans).toEqual({
+ providerId: "umans",
+ limit: 4,
+ inFlight: 0,
+ queued: 0,
+ paused: false,
+ });
+ });
+
+ it("destroy clears timers without error", () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 4);
+ manager.reportRateLimit("umans", 5000);
+ expect(() => manager.destroy()).not.toThrow();
+ });
+
+ it("release is idempotent (double-release does not overshoot)", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 2);
+
+ const release = await manager.acquire("umans", "conv1", 0);
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+
+ release();
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+
+ // Double-release should not decrement below 0.
+ release();
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+ });
+
+ it("multiple concurrent acquires up to the limit all resolve immediately", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 3);
+
+ const releases = await Promise.all([
+ manager.acquire("umans", "conv1", 0),
+ manager.acquire("umans", "conv2", 0),
+ manager.acquire("umans", "conv3", 0),
+ ]);
+
+ expect(manager.getStatus("umans")?.inFlight).toBe(3);
+
+ for (const release of releases) {
+ release();
+ }
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+ });
+});
diff --git a/packages/provider-concurrency/src/concurrency-manager.ts b/packages/provider-concurrency/src/concurrency-manager.ts
new file mode 100644
index 0000000..e7535cf
--- /dev/null
+++ b/packages/provider-concurrency/src/concurrency-manager.ts
@@ -0,0 +1,327 @@
+/**
+ * In-memory per-provider concurrency limiter.
+ *
+ * Tracks and limits how many concurrent API requests (token-generating
+ * requests) are in flight per provider. When the limit is reached, additional
+ * requests queue and are granted slots based on oldest-agent-first priority
+ * (the agent whose current prompt started the longest ago wins the next slot).
+ *
+ * A watchdog reclaims slots held beyond a timeout (deadlock / stuck-agent
+ * recovery). 429 backoff pauses a provider's queue for a configurable duration.
+ *
+ * This module is the PURE decision logic. It takes an injected clock (`now`)
+ * and injected timers (`setTimeout`/`clearTimeout`/`setInterval`/`clearInterval`)
+ * so it is fully testable with deterministic fake time. The extension layer
+ * wires real timers.
+ */
+
+// ─── Types ───────────────────────────────────────────────────────────────────
+
+/** Status snapshot for a single provider's concurrency state. */
+export interface ProviderConcurrencyStatus {
+ readonly providerId: string;
+ /** Configured concurrency limit. Always present (status is only returned for providers with a limit). */
+ readonly limit: number;
+ /** Currently in-flight (held) slots. */
+ readonly inFlight: number;
+ /** Agents waiting in the queue for a slot. */
+ readonly queued: number;
+ /** Whether the queue is paused (429 backoff). */
+ readonly paused: boolean;
+ /** When the pause expires (epoch-ms). Present only when paused. */
+ readonly pausedUntil?: number;
+}
+
+/**
+ * The limiter surface a consumer (session-orchestrator) needs: acquire a
+ * slot before a provider stream starts, release it when the stream completes,
+ * and report rate-limit (429) events so the manager can back off.
+ */
+export interface ConcurrencyLimiter {
+ /**
+ * Acquire a concurrency slot for `providerId`. Resolves immediately when a
+ * slot is available; otherwise blocks (queued by oldest-agent-first) until
+ * one frees up. The returned function MUST be called when the response
+ * stream completes (in a `finally` block). For providers with no configured
+ * limit, resolves instantly with a no-op release.
+ *
+ * @param providerId The provider to limit (e.g. "umans", "openai-compat").
+ * @param conversationId The agent requesting the slot.
+ * @param promptStartedAt When the agent's current prompt (turn) started
+ * (epoch-ms). Used for oldest-agent-first scheduling.
+ */
+ acquire(providerId: string, conversationId: string, promptStartedAt: number): Promise<() => void>;
+
+ /**
+ * Report a 429 from a provider. Pauses the queue for that provider for
+ * `retryAfterMs` (or a default duration when omitted). Queued and in-flight
+ * requests are unaffected; new `acquire` calls block until the pause expires.
+ */
+ reportRateLimit(providerId: string, retryAfterMs?: number): void;
+}
+
+/**
+ * The full service surface (limiter + config + status) for HTTP routes.
+ */
+export interface ConcurrencyService extends ConcurrencyLimiter {
+ /** Set the concurrency limit for a provider. Creates the state if new. */
+ setLimit(providerId: string, limit: number): void;
+ /** Get the configured limit, or `undefined` when none. */
+ getLimit(providerId: string): number | undefined;
+ /** Remove the limit for a provider (makes it unlimited). */
+ removeLimit(providerId: string): void;
+ /** All configured limits as `{ providerId, limit }` entries. */
+ getLimits(): readonly { providerId: string; limit: number }[];
+ /** Status for one provider, or `undefined` when no limit is configured. */
+ getStatus(providerId: string): ProviderConcurrencyStatus | undefined;
+ /** Status for every provider with a configured limit. */
+ getStatusAll(): readonly ProviderConcurrencyStatus[];
+ /** Stop the watchdog + clear all timers. */
+ destroy(): void;
+}
+
+// ─── Internal state ───────────────────────────────────────────────────────────
+
+interface Slot {
+ readonly conversationId: string;
+ readonly acquiredAt: number;
+ /** Idempotent release — safe to call from the holder or the watchdog. */
+ readonly releaseFn: () => void;
+}
+
+interface QueuedWaiter {
+ readonly conversationId: string;
+ readonly promptStartedAt: number;
+ readonly resolve: (release: () => void) => void;
+}
+
+interface ProviderState {
+ limit: number;
+ inFlight: number;
+ slots: Map<number, Slot>;
+ queue: QueuedWaiter[];
+ paused: boolean;
+ pausedUntil: number | undefined;
+ pauseTimer: ReturnType<typeof setTimeout> | undefined;
+}
+
+export interface ConcurrencyManagerOpts {
+ /** Monotonic-ish clock (epoch-ms). */
+ readonly now: () => number;
+ /** Max time a slot may be held before the watchdog reclaims it (ms). */
+ readonly slotTimeoutMs: number;
+ /** How often the watchdog sweeps (ms). */
+ readonly watchdogIntervalMs: number;
+ /** Default pause duration when a 429 arrives without Retry-After (ms). */
+ readonly defaultPauseMs: number;
+ /** Injected timers (default: global). Override in tests for deterministic time. */
+ readonly setTimeout?: typeof setTimeout;
+ readonly clearTimeout?: typeof clearTimeout;
+ readonly setInterval?: typeof setInterval;
+ readonly clearInterval?: typeof clearInterval;
+ /** Optional logger for watchdog + pause events. */
+ readonly onWatchdogReclaim?: (providerId: string, conversationId: string, heldMs: number) => void;
+ readonly onPause?: (providerId: string, durationMs: number) => void;
+}
+
+function noopRelease(): void {
+ // No limit configured → nothing to release.
+}
+
+export function createConcurrencyManager(opts: ConcurrencyManagerOpts): ConcurrencyService {
+ const now = opts.now;
+ const slotTimeoutMs = opts.slotTimeoutMs;
+ const defaultPauseMs = opts.defaultPauseMs;
+ const setTimeout = opts.setTimeout ?? globalThis.setTimeout.bind(globalThis);
+ const clearTimeout = opts.clearTimeout ?? globalThis.clearTimeout.bind(globalThis);
+ const setInterval = opts.setInterval ?? globalThis.setInterval.bind(globalThis);
+ const clearInterval = opts.clearInterval ?? globalThis.clearInterval.bind(globalThis);
+
+ const states = new Map<string, ProviderState>();
+ let slotIdCounter = 0;
+
+ // ── Slot granting ──────────────────────────────────────────────────────────
+
+ function grantSlot(state: ProviderState, providerId: string, conversationId: string): () => void {
+ const id = slotIdCounter++;
+ let released = false;
+ const releaseFn = () => {
+ if (released) return;
+ released = true;
+ state.slots.delete(id);
+ state.inFlight--;
+ tryGrantNext(providerId);
+ };
+ state.slots.set(id, {
+ conversationId,
+ acquiredAt: now(),
+ releaseFn,
+ });
+ state.inFlight++;
+ return releaseFn;
+ }
+
+ function tryGrantNext(providerId: string): void {
+ const state = states.get(providerId);
+ if (state === undefined) return;
+ if (state.paused) return;
+ while (state.queue.length > 0 && state.inFlight < state.limit) {
+ const waiter = state.queue[0];
+ if (waiter === undefined) break;
+ state.queue.shift();
+ const releaseFn = grantSlot(state, providerId, waiter.conversationId);
+ waiter.resolve(releaseFn);
+ }
+ }
+
+ // ── Watchdog ──────────────────────────────────────────────────────────────────
+
+ function sweep(): void {
+ const currentNow = now();
+ for (const [providerId, state] of states) {
+ for (const [, slot] of state.slots) {
+ const heldMs = currentNow - slot.acquiredAt;
+ if (heldMs > slotTimeoutMs) {
+ opts.onWatchdogReclaim?.(providerId, slot.conversationId, heldMs);
+ slot.releaseFn();
+ }
+ }
+ }
+ }
+
+ const watchdogTimer = setInterval(sweep, opts.watchdogIntervalMs);
+
+ // ── Public API ─────────────────────────────────────────────────────────────
+
+ const manager: ConcurrencyService = {
+ acquire(providerId, conversationId, promptStartedAt) {
+ const state = states.get(providerId);
+ if (state === undefined) {
+ // No limit configured → unlimited.
+ return Promise.resolve(noopRelease);
+ }
+
+ if (!state.paused && state.inFlight < state.limit) {
+ return Promise.resolve(grantSlot(state, providerId, conversationId));
+ }
+
+ // Queue (oldest-agent-first by promptStartedAt).
+ return new Promise<() => void>((resolve) => {
+ state.queue.push({ conversationId, promptStartedAt, resolve });
+ // Keep sorted ascending by promptStartedAt (oldest first).
+ // Insertion sort would be O(n), but the queue is typically tiny (<20),
+ // so a simple sort is fine and keeps the code simple.
+ state.queue.sort((a, b) => a.promptStartedAt - b.promptStartedAt);
+ });
+ },
+
+ reportRateLimit(providerId, retryAfterMs) {
+ const state = states.get(providerId);
+ if (state === undefined) return;
+
+ const pauseDuration = retryAfterMs ?? defaultPauseMs;
+ state.paused = true;
+ state.pausedUntil = now() + pauseDuration;
+
+ if (state.pauseTimer !== undefined) {
+ clearTimeout(state.pauseTimer);
+ }
+ opts.onPause?.(providerId, pauseDuration);
+ state.pauseTimer = setTimeout(() => {
+ state.paused = false;
+ state.pausedUntil = undefined;
+ state.pauseTimer = undefined;
+ tryGrantNext(providerId);
+ }, pauseDuration);
+ },
+
+ setLimit(providerId, limit) {
+ let state = states.get(providerId);
+ if (state === undefined) {
+ state = {
+ limit,
+ inFlight: 0,
+ slots: new Map(),
+ queue: [],
+ paused: false,
+ pausedUntil: undefined,
+ pauseTimer: undefined,
+ };
+ states.set(providerId, state);
+ } else {
+ state.limit = limit;
+ }
+ // A higher limit may let queued requests through.
+ tryGrantNext(providerId);
+ },
+
+ getLimit(providerId) {
+ return states.get(providerId)?.limit;
+ },
+
+ removeLimit(providerId) {
+ const state = states.get(providerId);
+ if (state === undefined) return;
+
+ // Clear pause.
+ state.paused = false;
+ state.pausedUntil = undefined;
+ if (state.pauseTimer !== undefined) {
+ clearTimeout(state.pauseTimer);
+ state.pauseTimer = undefined;
+ }
+
+ // Grant all queued requests (they become unlimited now).
+ while (state.queue.length > 0) {
+ const waiter = state.queue[0];
+ if (waiter === undefined) break;
+ state.queue.shift();
+ const releaseFn = grantSlot(state, providerId, waiter.conversationId);
+ waiter.resolve(releaseFn);
+ }
+
+ // Remove the state. In-flight slots' release functions still work —
+ // they close over `state` and call `tryGrantNext` which finds no state
+ // and returns early. The watchdog won't sweep removed states.
+ states.delete(providerId);
+ },
+
+ getLimits() {
+ return [...states.entries()].map(([providerId, s]) => ({
+ providerId,
+ limit: s.limit,
+ }));
+ },
+
+ getStatus(providerId) {
+ const state = states.get(providerId);
+ if (state === undefined) return undefined;
+ return {
+ providerId,
+ limit: state.limit,
+ inFlight: state.inFlight,
+ queued: state.queue.length,
+ paused: state.paused,
+ ...(state.pausedUntil !== undefined ? { pausedUntil: state.pausedUntil } : {}),
+ };
+ },
+
+ getStatusAll() {
+ return [...states.keys()]
+ .map((providerId) => manager.getStatus(providerId))
+ .filter((s): s is ProviderConcurrencyStatus => s !== undefined);
+ },
+
+ destroy() {
+ clearInterval(watchdogTimer);
+ for (const state of states.values()) {
+ if (state.pauseTimer !== undefined) {
+ clearTimeout(state.pauseTimer);
+ }
+ }
+ states.clear();
+ },
+ };
+
+ return manager;
+}
diff --git a/packages/provider-concurrency/src/extension.ts b/packages/provider-concurrency/src/extension.ts
new file mode 100644
index 0000000..c741173
--- /dev/null
+++ b/packages/provider-concurrency/src/extension.ts
@@ -0,0 +1,61 @@
+import type { Extension, HostAPI, Manifest } from "@dispatch/kernel";
+import { type ConcurrencyService, createConcurrencyManager } from "./concurrency-manager.js";
+import { concurrencyServiceHandle } from "./service.js";
+
+export const manifest: Manifest = {
+ id: "provider-concurrency",
+ name: "Provider Concurrency Limits",
+ version: "0.0.0",
+ apiVersion: "^0.1.0",
+ trust: "bundled",
+ activation: "eager",
+ contributes: { services: ["provider-concurrency/service"] },
+};
+
+/**
+ * Default tuning constants.
+ *
+ * - `SLOT_TIMEOUT_MS` (5 min): a slot held longer than this is force-reclaimed
+ * by the watchdog (deadlock / stuck-agent recovery). Generation streams
+ * rarely exceed 2–3 minutes; 5 min is a generous safety margin.
+ * - `WATCHDOG_INTERVAL_MS` (30s): how often the watchdog sweeps for stale slots.
+ * - `DEFAULT_PAUSE_MS` (30s): default 429 backoff when no Retry-After is given.
+ * Umans docs note each concurrency 429 deprioritizes the account for ~30 min,
+ * but a 30s queue pause prevents immediate re-overshoot while still allowing
+ * recovery.
+ */
+const SLOT_TIMEOUT_MS = 5 * 60 * 1000;
+const WATCHDOG_INTERVAL_MS = 30 * 1000;
+const DEFAULT_PAUSE_MS = 30 * 1000;
+
+export function activate(host: HostAPI): void {
+ const logger = host.logger;
+
+ const manager: ConcurrencyService = createConcurrencyManager({
+ now: () => Date.now(),
+ slotTimeoutMs: SLOT_TIMEOUT_MS,
+ watchdogIntervalMs: WATCHDOG_INTERVAL_MS,
+ defaultPauseMs: DEFAULT_PAUSE_MS,
+ onWatchdogReclaim: (providerId, conversationId, heldMs) => {
+ logger.warn("provider-concurrency: watchdog reclaimed stale slot", {
+ providerId,
+ conversationId,
+ heldMs,
+ });
+ },
+ onPause: (providerId, durationMs) => {
+ logger.warn("provider-concurrency: 429 backoff — pausing queue", {
+ providerId,
+ durationMs,
+ });
+ },
+ });
+
+ host.provideService(concurrencyServiceHandle, manager);
+ logger.info("provider-concurrency: registered");
+}
+
+export const extension: Extension = {
+ manifest,
+ activate,
+};
diff --git a/packages/provider-concurrency/src/index.ts b/packages/provider-concurrency/src/index.ts
new file mode 100644
index 0000000..f35c070
--- /dev/null
+++ b/packages/provider-concurrency/src/index.ts
@@ -0,0 +1,10 @@
+export {
+ type ConcurrencyLimiter,
+ type ConcurrencyManagerOpts,
+ type ConcurrencyService,
+ createConcurrencyManager,
+ type ProviderConcurrencyStatus,
+} from "./concurrency-manager.js";
+export { extension, manifest } from "./extension.js";
+export { wrapProviderWithConcurrency } from "./provider-wrapper.js";
+export { concurrencyServiceHandle } from "./service.js";
diff --git a/packages/provider-concurrency/src/provider-wrapper.test.ts b/packages/provider-concurrency/src/provider-wrapper.test.ts
new file mode 100644
index 0000000..e024d59
--- /dev/null
+++ b/packages/provider-concurrency/src/provider-wrapper.test.ts
@@ -0,0 +1,173 @@
+import type { ProviderContract, ProviderEvent } from "@dispatch/kernel";
+import { describe, expect, it } from "vitest";
+import type { ConcurrencyLimiter } from "./concurrency-manager.js";
+import { wrapProviderWithConcurrency } from "./provider-wrapper.js";
+
+/** Build a fake provider that yields a sequence of events. */
+function fakeProvider(events: ProviderEvent[]): ProviderContract {
+ return {
+ id: "test-provider",
+ stream: async function* (): AsyncIterable<ProviderEvent> {
+ for (const e of events) {
+ yield e;
+ }
+ },
+ };
+}
+
+/** A fake limiter that records acquire/release calls. */
+function recordingLimiter(): ConcurrencyLimiter & {
+ acquireCalls: { providerId: string; conversationId: string; promptStartedAt: number }[];
+ releaseCalls: number;
+ rateLimitReports: string[];
+} {
+ const acquireCalls: { providerId: string; conversationId: string; promptStartedAt: number }[] =
+ [];
+ const releaseCalls: { count: number } = { count: 0 };
+ const rateLimitReports: string[] = [];
+
+ return {
+ acquireCalls,
+ get releaseCalls() {
+ return releaseCalls.count;
+ },
+ rateLimitReports,
+ acquire(providerId, conversationId, promptStartedAt) {
+ acquireCalls.push({ providerId, conversationId, promptStartedAt });
+ return Promise.resolve(() => {
+ releaseCalls.count++;
+ });
+ },
+ reportRateLimit(providerId) {
+ rateLimitReports.push(providerId);
+ },
+ };
+}
+
+describe("wrapProviderWithConcurrency", () => {
+ it("acquires a slot before streaming and releases after the stream completes", async () => {
+ const provider = fakeProvider([
+ { type: "text-delta", delta: "hello" },
+ { type: "finish", reason: "stop" },
+ ]);
+ const limiter = recordingLimiter();
+
+ const wrapped = wrapProviderWithConcurrency(provider, limiter, "conv1", 12345);
+
+ const events: ProviderEvent[] = [];
+ for await (const e of wrapped.stream([], [])) {
+ events.push(e);
+ }
+
+ // Slot acquired before stream, released after.
+ expect(limiter.acquireCalls).toEqual([
+ { providerId: "test-provider", conversationId: "conv1", promptStartedAt: 12345 },
+ ]);
+ expect(limiter.releaseCalls).toBe(1);
+ expect(events).toEqual([
+ { type: "text-delta", delta: "hello" },
+ { type: "finish", reason: "stop" },
+ ]);
+ });
+
+ it("releases the slot even when the stream throws", async () => {
+ const provider: ProviderContract = {
+ id: "err-provider",
+ stream: async function* (): AsyncIterable<ProviderEvent> {
+ yield { type: "text-delta", delta: "partial" };
+ throw new Error("stream exploded");
+ },
+ };
+ const limiter = recordingLimiter();
+ const wrapped = wrapProviderWithConcurrency(provider, limiter, "conv1", 0);
+
+ await expect(async () => {
+ for await (const _e of wrapped.stream([], [])) {
+ // consume
+ }
+ }).rejects.toThrow("stream exploded");
+
+ expect(limiter.releaseCalls).toBe(1);
+ });
+
+ it("reports 429 errors to the limiter", async () => {
+ const provider = fakeProvider([
+ { type: "error", message: "Too many requests", code: "429", retryable: true },
+ ]);
+ const limiter = recordingLimiter();
+ const wrapped = wrapProviderWithConcurrency(provider, limiter, "conv1", 0);
+
+ const events: ProviderEvent[] = [];
+ for await (const e of wrapped.stream([], [])) {
+ events.push(e);
+ }
+
+ expect(limiter.rateLimitReports).toEqual(["test-provider"]);
+ // The 429 error event is still yielded to the consumer (kernel handles retry).
+ expect(events).toHaveLength(1);
+ expect(events[0]?.type).toBe("error");
+ });
+
+ it("does not report non-429 errors", async () => {
+ const provider = fakeProvider([
+ { type: "error", message: "Internal error", code: "500", retryable: true },
+ ]);
+ const limiter = recordingLimiter();
+ const wrapped = wrapProviderWithConcurrency(provider, limiter, "conv1", 0);
+
+ for await (const _e of wrapped.stream([], [])) {
+ // consume
+ }
+
+ expect(limiter.rateLimitReports).toEqual([]);
+ });
+
+ it("preserves the provider id and listModels", async () => {
+ const provider: ProviderContract = {
+ id: "my-provider",
+ stream: async function* (): AsyncIterable<ProviderEvent> {
+ yield { type: "finish", reason: "stop" };
+ },
+ listModels: async () => [{ id: "model-1" }],
+ };
+ const limiter = recordingLimiter();
+ const wrapped = wrapProviderWithConcurrency(provider, limiter, "conv1", 0);
+
+ expect(wrapped.id).toBe("my-provider");
+ expect(wrapped.listModels).toBeDefined();
+ const models = await wrapped.listModels?.();
+ expect(models).toEqual([{ id: "model-1" }]);
+ });
+
+ it("passes through messages, tools, and opts to the inner stream", async () => {
+ let receivedArgs:
+ | {
+ messages: unknown;
+ tools: unknown;
+ opts: unknown;
+ }
+ | undefined;
+
+ const provider: ProviderContract = {
+ id: "passthrough",
+ stream: async function* (messages, tools, opts): AsyncIterable<ProviderEvent> {
+ receivedArgs = { messages, tools, opts };
+ yield { type: "finish", reason: "stop" };
+ },
+ };
+ const limiter = recordingLimiter();
+ const wrapped = wrapProviderWithConcurrency(provider, limiter, "conv1", 0);
+
+ const messages = [{ role: "user" as const, chunks: [{ type: "text" as const, text: "hi" }] }];
+ const tools = [{ name: "test_tool", description: "test", parameters: {} }];
+ const opts = { model: "gpt-4" };
+
+ for await (const _e of wrapped.stream(messages, tools, opts)) {
+ // consume
+ }
+
+ expect(receivedArgs?.messages).toBe(messages);
+ expect(receivedArgs?.tools).toBe(tools);
+ expect(receivedArgs?.opts).toBe(opts);
+ });
+});
diff --git a/packages/provider-concurrency/src/provider-wrapper.ts b/packages/provider-concurrency/src/provider-wrapper.ts
new file mode 100644
index 0000000..ee3ca85
--- /dev/null
+++ b/packages/provider-concurrency/src/provider-wrapper.ts
@@ -0,0 +1,59 @@
+import type {
+ ChatMessage,
+ ProviderContract,
+ ProviderEvent,
+ ProviderStreamOptions,
+ ToolContract,
+} from "@dispatch/kernel";
+import type { ConcurrencyLimiter } from "./concurrency-manager.js";
+
+/**
+ * Wrap a provider's `stream` method with concurrency limiting.
+ *
+ * A slot is acquired BEFORE the first event is yielded (before the HTTP
+ * request is sent — the `await limiter.acquire()` runs before the generator
+ * body starts iterating the inner stream). The slot is released in a `finally`
+ * block AFTER the inner stream completes (the full response stream, not just
+ * HTTP headers — matching the Umans concurrency model where a slot is held
+ * only while tokens are actually generating).
+ *
+ * 429 detection: if the provider yields an `error` event with `code: "429"`,
+ * the limiter is notified so it can pause the queue for that provider.
+ *
+ * @param provider The underlying provider to wrap.
+ * @param limiter The concurrency limiter (acquire/release/reportRateLimit).
+ * @param conversationId The agent requesting the stream (for slot attribution).
+ * @param promptStartedAt When the agent's current prompt (turn) started
+ * (epoch-ms, for oldest-agent-first scheduling).
+ */
+export function wrapProviderWithConcurrency(
+ provider: ProviderContract,
+ limiter: ConcurrencyLimiter,
+ conversationId: string,
+ promptStartedAt: number,
+): ProviderContract {
+ const innerStream = provider.stream;
+ const providerId = provider.id;
+
+ return {
+ id: provider.id,
+ stream: async function* (
+ messages: readonly ChatMessage[],
+ tools: readonly ToolContract[],
+ opts?: ProviderStreamOptions,
+ ): AsyncIterable<ProviderEvent> {
+ const release = await limiter.acquire(providerId, conversationId, promptStartedAt);
+ try {
+ for await (const event of innerStream(messages, tools, opts)) {
+ if (event.type === "error" && event.code === "429") {
+ limiter.reportRateLimit(providerId);
+ }
+ yield event;
+ }
+ } finally {
+ release();
+ }
+ },
+ ...(provider.listModels !== undefined ? { listModels: provider.listModels } : {}),
+ };
+}
diff --git a/packages/provider-concurrency/src/service.ts b/packages/provider-concurrency/src/service.ts
new file mode 100644
index 0000000..aa578e8
--- /dev/null
+++ b/packages/provider-concurrency/src/service.ts
@@ -0,0 +1,11 @@
+import { defineService } from "@dispatch/kernel";
+import type { ConcurrencyService } from "./concurrency-manager.js";
+
+/**
+ * Typed service handle for the provider-concurrency service. The
+ * `provider-concurrency` extension provides the implementation; the
+ * session-orchestrator + transport-http consume it.
+ */
+export const concurrencyServiceHandle = defineService<ConcurrencyService>(
+ "provider-concurrency/service",
+);