summaryrefslogtreecommitdiffhomepage
diff options
context:
space:
mode:
-rw-r--r--bun.lock12
-rw-r--r--notes/concurrency-library-investigation.md228
-rw-r--r--packages/host-bin/package.json1
-rw-r--r--packages/host-bin/src/main.ts2
-rw-r--r--packages/host-bin/tsconfig.json3
-rw-r--r--packages/provider-concurrency/package.json11
-rw-r--r--packages/provider-concurrency/src/concurrency-manager.test.ts487
-rw-r--r--packages/provider-concurrency/src/concurrency-manager.ts376
-rw-r--r--packages/provider-concurrency/src/extension.ts140
-rw-r--r--packages/provider-concurrency/src/index.ts10
-rw-r--r--packages/provider-concurrency/src/provider-wrapper.test.ts246
-rw-r--r--packages/provider-concurrency/src/provider-wrapper.ts68
-rw-r--r--packages/provider-concurrency/src/service.ts11
-rw-r--r--packages/provider-concurrency/tsconfig.json6
-rw-r--r--packages/session-orchestrator/package.json1
-rw-r--r--packages/session-orchestrator/src/extension.ts14
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts69
-rw-r--r--packages/session-orchestrator/tsconfig.json1
-rw-r--r--packages/transport-contract/src/index.ts53
-rw-r--r--packages/transport-http/package.json1
-rw-r--r--packages/transport-http/src/app.ts90
-rw-r--r--packages/transport-http/src/extension.ts15
-rw-r--r--packages/transport-http/src/logic.ts2
-rw-r--r--packages/transport-http/src/seam.ts2
-rw-r--r--packages/transport-http/tsconfig.json2
-rw-r--r--packages/wire/src/index.ts10
-rw-r--r--tsconfig.json3
27 files changed, 1861 insertions, 3 deletions
diff --git a/bun.lock b/bun.lock
index d9762f7..493da15 100644
--- a/bun.lock
+++ b/bun.lock
@@ -84,6 +84,7 @@
"@dispatch/lsp": "workspace:*",
"@dispatch/mcp": "workspace:*",
"@dispatch/message-queue": "workspace:*",
+ "@dispatch/provider-concurrency": "workspace:*",
"@dispatch/provider-openai-compat": "workspace:*",
"@dispatch/provider-umans": "workspace:*",
"@dispatch/session-orchestrator": "workspace:*",
@@ -162,6 +163,13 @@
"@dispatch/wire": "workspace:*",
},
},
+ "packages/provider-concurrency": {
+ "name": "@dispatch/provider-concurrency",
+ "version": "0.0.0",
+ "dependencies": {
+ "@dispatch/kernel": "workspace:*",
+ },
+ },
"packages/provider-openai-compat": {
"name": "@dispatch/provider-openai-compat",
"version": "0.0.0",
@@ -187,6 +195,7 @@
"@dispatch/credential-store": "workspace:*",
"@dispatch/kernel": "workspace:*",
"@dispatch/message-queue": "workspace:*",
+ "@dispatch/provider-concurrency": "workspace:*",
"@dispatch/system-prompt": "workspace:*",
},
},
@@ -339,6 +348,7 @@
"@dispatch/kernel": "workspace:*",
"@dispatch/lsp": "workspace:*",
"@dispatch/mcp": "workspace:*",
+ "@dispatch/provider-concurrency": "workspace:*",
"@dispatch/session-orchestrator": "workspace:*",
"@dispatch/system-prompt": "workspace:*",
"@dispatch/throughput-store": "workspace:*",
@@ -426,6 +436,8 @@
"@dispatch/openai-stream": ["@dispatch/openai-stream@workspace:packages/openai-stream"],
+ "@dispatch/provider-concurrency": ["@dispatch/provider-concurrency@workspace:packages/provider-concurrency"],
+
"@dispatch/provider-openai-compat": ["@dispatch/provider-openai-compat@workspace:packages/provider-openai-compat"],
"@dispatch/provider-umans": ["@dispatch/provider-umans@workspace:packages/provider-umans"],
diff --git a/notes/concurrency-library-investigation.md b/notes/concurrency-library-investigation.md
new file mode 100644
index 0000000..abf7dab
--- /dev/null
+++ b/notes/concurrency-library-investigation.md
@@ -0,0 +1,228 @@
+# Concurrency Library Investigation
+
+> Research-only. No code changes. Investigated 2026-06-26.
+
+## 1. ai-concurrency-shaper (joeycumines)
+
+**Repository:** https://github.com/joeycumines/ai-concurrency-shaper
+
+### What it is
+A **standalone reverse proxy** written in **Go**. It sits between your
+application and an upstream AI/LLM API (e.g. Anthropic, OpenAI), limiting
+concurrent requests to configured HTTP routes. Requests exceeding the limit
+block until a slot opens. Non-matching requests pass through unmodified.
+
+### How it works
+- Uses a **token-bucket channel** (Go's `chan struct{}`) as a semaphore. Each
+ limited request acquires a token; the token is returned when the request
+ completes (full response body streamed).
+- Supports **per-route** concurrency limits via `-limit "POST /v1/chat/completions:4"`
+ and **global** concurrency via `-global-concurrency 10`.
+- Supports **route grouping** — multiple routes can share one limiter via
+ `@group` syntax.
+- Has a **circuit breaker** (trips after N failures in a window, exponential
+ backoff with phantom concurrency holds).
+- Has **concurrency protection flags**: release cooldown, cancel cooldown,
+ failure hold, adaptive headroom (reduce effective limit by 1 after a 429).
+- **429 handling**: skips retrying 429s by default (`-retry-skip-429=true`).
+- Optional **TUI dashboard** (Bubble Tea) with live metrics.
+- **In-memory only** — no persistence, no external state (single process).
+
+### Maturity
+- **1 star, 1 fork, 7 commits**, created June 10, 2026 (16 days ago).
+- **Single contributor** (Joseph Cumines).
+- **GPL-3.0 license** (copyleft — incompatible with closed-source distribution).
+- Written in **Go**, not JavaScript/TypeScript.
+
+### Critical finding: language mismatch
+This is a **Go binary**, not an npm library. It cannot be imported into a
+Bun/TypeScript application. To use it, we would need to:
+1. Run it as a **separate process** (a sidecar proxy).
+2. Point our provider's `baseURL` at the proxy's listen address.
+3. The proxy forwards to the real upstream.
+
+This fundamentally changes our architecture from "in-process concurrency
+management" to "external proxy-based concurrency management."
+
+### Feature comparison vs our implementation
+
+| Feature | Our impl | ai-concurrency-shaper |
+|---|---|---|
+| Per-provider configurable limits | ✅ (runtime API + in-memory) | ✅ (per-route CLI flags, not runtime) |
+| Oldest-agent-first scheduling | ✅ (priority queue by promptStartedAt) | ❌ (FIFO semaphore, no priority) |
+| Slot held only during token generation | ✅ (acquire/release around stream) | ✅ (token held until response body completes) |
+| 429 backoff | ✅ (pause queue, Retry-After aware) | ✅ (circuit breaker + adaptive headroom) |
+| Watchdog/deadlock recovery | ✅ (timeout-based slot reclaim) | ❌ (no watchdog; relies on queue-timeout) |
+| Runtime-configurable limits | ✅ (HTTP API, no restart) | ❌ (CLI flags only, requires restart) |
+| In-memory / no external state | ✅ | ✅ |
+| Language | TypeScript (Bun) | Go |
+| Integration model | In-process (wraps ProviderContract) | External proxy (separate process) |
+| Status reporting API | ✅ (HTTP endpoints) | TUI only (no HTTP API) |
+| License | Our code (MIT-style) | GPL-3.0 |
+
+### What would change to use it
+1. Deploy the Go binary as a sidecar process.
+2. Change each provider's `baseURL` to point at the proxy.
+3. Lose runtime-configurable limits (CLI flags only).
+4. Lose oldest-agent-first scheduling (FIFO only).
+5. Lose the in-process status API (TUI only, no HTTP).
+6. Add a process management dependency (start/stop/monitor the proxy).
+7. Accept GPL-3.0 license implications.
+
+## 2. Drawbacks of using ai-concurrency-shaper
+
+### Fatal drawbacks
+1. **No oldest-agent-first scheduling.** Our implementation prioritizes the
+ agent whose prompt started longest ago. The proxy uses a simple FIFO
+ semaphore — no priority queue. This is a core requirement.
+2. **No runtime-configurable limits.** Limits are set via CLI flags at startup.
+ Our frontend has a settings UI to add/remove/update limits at runtime.
+3. **Language mismatch.** It's a Go binary, not an npm package. We'd need to
+ run it as a separate process, adding operational complexity.
+4. **No HTTP status API.** The proxy exposes status only via TUI. Our frontend
+ polls `GET /concurrency/status` for live in-flight/queued/paused state.
+5. **GPL-3.0 license.** Copyleft — may be incompatible with our distribution
+ model.
+
+### Secondary drawbacks
+6. **No watchdog.** If a slot holder dies (e.g. the proxy's connection to the
+ client drops but the upstream request is still in flight), there's no
+ timeout-based reclaim. It relies on `-queue-timeout` for queue wait, not
+ for held slots.
+7. **Immature.** 1 star, 7 commits, 16 days old, single contributor. No
+ community, no battle-testing.
+8. **No provider-awareness.** It limits by HTTP route pattern, not by
+ "provider." Our implementation is keyed on `providerId` (e.g. "umans",
+ "openai-compat"), which maps naturally to our `ProviderContract.id`.
+9. **Operational overhead.** Running a separate Go process alongside the Bun
+ app adds deployment complexity, monitoring burden, and a failure mode
+ (proxy down = all requests fail).
+
+## 3. Other libraries/tools for AI API concurrency
+
+### p-queue (sindresorhus)
+- **npm**, TypeScript, MIT, 4.2k stars, 727k dependents.
+- General-purpose promise queue with concurrency control.
+- Supports **priority** (`{priority: number}` per task), **rate limiting**
+ (`intervalCap` + `interval`), **pause/resume**, **timeouts**, **custom
+ queue class** (for custom scheduling), AbortSignal cancellation.
+- **No AI-specific logic** — no 429 detection, no provider concept, no
+ streaming-awareness.
+- Could be used as a **building block** to replace our queue implementation,
+ but we'd still need the provider wrapper, 429 detection, watchdog, and
+ status API on top.
+- **Feature complete** (maintainer says no further development planned, but
+ accepts PRs).
+
+### p-limit (sindresorhus)
+- **npm**, TypeScript, MIT, 5.4k dependents.
+- Simpler than p-queue — just limits concurrent executions. No queue, no
+ priority, no pause. Not suitable for our use case (we need queuing).
+
+### ai-sdk-rate-limiter (piyushgupta344)
+- **npm**, TypeScript, zero dependencies.
+- Designed for the Vercel AI SDK (`@ai-sdk/*`). Wraps model objects.
+- Has **priority queuing** (high/normal/low lanes), **429 backoff with
+ Retry-After**, **cost tracking & budget enforcement**, **multi-tenant
+ scopes**, **Redis for multi-instance**, **observability** (Prometheus,
+ OpenTelemetry).
+- Built-in limits for OpenAI, Anthropic, Google, Groq, Mistral, Cohere.
+- **Closest to our use case** of any npm library found.
+- **Maturity unknown** — couldn't verify star count or maintenance activity
+ from the docs site. Appears to be a solo project.
+- **Concerns**: It's designed for the Vercel AI SDK's model-wrapping pattern.
+ Our provider architecture uses `ProviderContract.stream()` returning an
+ `AsyncIterable<ProviderEvent>`, not the Vercel AI SDK's model interface.
+ We'd need an adapter. It also doesn't expose a status API for the frontend.
+
+### LiteLLM (BerriAI)
+- **Python**, MIT, 51.7k GitHub stars, very mature (39k+ commits).
+- Full **LLM gateway/proxy** — not a library. Runs as a separate server.
+- Has `enforce_model_rate_limits` for RPM/TPM hard limits.
+- Supports **least-busy routing**, **latency-based routing**, **cost-based
+ routing**, **fallbacks**, **deployment priority**.
+- **No concurrency limiting** — it limits requests-per-minute (RPM) and
+ tokens-per-minute (TPM), not concurrent in-flight requests. RPM is a rate
+ limit, not a concurrency limit. An agent that sends 4 requests in 1 second
+ and then waits would be blocked by RPM=60 but not by concurrency=4.
+- Would require running a Python proxy server alongside our Bun app.
+- Overkill for our use case — it's a full LLM gateway with virtual keys, cost
+ tracking, guardrails, etc.
+
+### Portkey AI Gateway
+- **Hosted/self-hosted**, open-source (Apache-2.0).
+- Full AI gateway with routing, fallbacks, caching, observability.
+- Rate limiting is **Enterprise-only** and per-team/per-key, not per-provider
+ concurrency.
+- Same architectural model as LiteLLM (external proxy).
+
+### Kong AI Gateway
+- Built on Kong's API gateway. Enterprise-focused.
+- Rate limiting is plugin-based (RPM/TPM), not concurrency-based.
+- Same external-proxy model.
+
+## 4. Recommendation
+
+### **Keep the custom implementation. Do not switch to a library.**
+
+### Rationale
+
+1. **No library matches our requirements.** The core differentiators of our
+ implementation — **oldest-agent-first scheduling**, **runtime-configurable
+ per-provider limits**, **in-process status API**, and **stream-aware slot
+ lifecycle** — are not found in any library or proxy we investigated:
+ - `ai-concurrency-shaper` is a Go proxy with FIFO, no priority, no runtime
+ config, no HTTP status API, and GPL-3.0.
+ - `p-queue` has priority but no AI-specific logic, no 429 detection, no
+ streaming-awareness, no status API.
+ - `ai-sdk-rate-limiter` is close but designed for the Vercel AI SDK's model
+ interface, not our `ProviderContract.stream()` pattern.
+ - LiteLLM/Portkey/Kong are full gateway proxies with RPM/TPM rate limiting,
+ not concurrency limiting, and require running a separate server.
+
+2. **Our implementation is small and well-tested.** The core
+ `concurrency-manager.ts` is ~280 lines of pure logic with 15 tests and
+ injected timers. The `provider-wrapper.ts` is ~70 lines with 6 tests. The
+ total surface is tiny — there's no maintenance burden to justify
+ offloading to a library.
+
+3. **The architecture is a perfect fit.** Our `ProviderContract` wrapping
+ pattern (acquire before stream, release after stream) is the cleanest
+ possible integration point. An external proxy would add a network hop,
+ a process to manage, and break the direct relationship between the
+ orchestrator and the provider.
+
+4. **Oldest-agent-first scheduling is a hard requirement.** No library or
+ proxy we found supports priority queuing by turn-start time. This is the
+ core value proposition of our implementation — older agents complete
+ sooner, reducing overall wait time.
+
+5. **Runtime configurability is a hard requirement.** The frontend has a
+ settings UI for adding/removing/updating limits without restart. No
+ library supports this — they all require static configuration.
+
+### What we would lose by switching
+- Oldest-agent-first scheduling (no library supports it).
+- Runtime-configurable limits (no library supports it).
+- In-process status API (no library supports it).
+- Direct integration with `ProviderContract` (would need adapter or proxy).
+- ~280 lines of well-tested code that we own and control.
+
+### What we would gain by switching
+- Nothing we don't already have. The features libraries offer (circuit
+ breakers, adaptive headroom, TUI dashboards) are either already in our
+ implementation (429 backoff, watchdog) or not needed (TUI — we have a web
+ frontend).
+
+### When to reconsider
+- If we need **multi-instance** concurrency coordination (multiple app
+ processes sharing limits), we would need Redis-backed state. At that point,
+ `ai-sdk-rate-limiter` (which has a Redis store) or a custom Redis-backed
+ extension of our current implementation would be worth evaluating.
+- If we need **RPM/TPM rate limiting** (not just concurrency), LiteLLM or
+ Portkey could complement our concurrency limits. But these are orthogonal
+ concerns — RPM limits total requests per minute, while concurrency limits
+ in-flight requests at any moment.
+- If `ai-sdk-rate-limiter` matures and adds a status API + custom model
+ adapters, it could be worth re-evaluating as a replacement for the queue
+ layer (while keeping our provider wrapper + status API).
diff --git a/packages/host-bin/package.json b/packages/host-bin/package.json
index b5ab954..7d3b38c 100644
--- a/packages/host-bin/package.json
+++ b/packages/host-bin/package.json
@@ -13,6 +13,7 @@
"@dispatch/exec-backend": "workspace:*",
"@dispatch/heartbeat": "workspace:*",
"@dispatch/provider-openai-compat": "workspace:*",
+ "@dispatch/provider-concurrency": "workspace:*",
"@dispatch/provider-umans": "workspace:*",
"@dispatch/message-queue": "workspace:*",
"@dispatch/mcp": "workspace:*",
diff --git a/packages/host-bin/src/main.ts b/packages/host-bin/src/main.ts
index a5dabab..aa114d5 100644
--- a/packages/host-bin/src/main.ts
+++ b/packages/host-bin/src/main.ts
@@ -24,6 +24,7 @@ import {
import { extension as lspExt } from "@dispatch/lsp";
import { extension as mcpExt } from "@dispatch/mcp";
import { extension as messageQueueExt } from "@dispatch/message-queue";
+import { extension as providerConcurrencyExt } from "@dispatch/provider-concurrency";
import { extension as providerOpenaiCompatExt } from "@dispatch/provider-openai-compat";
import { extension as providerUmansExt } from "@dispatch/provider-umans";
import { extension as sessionOrchestratorExt } from "@dispatch/session-orchestrator";
@@ -80,6 +81,7 @@ const CORE_EXTENSIONS: readonly Extension[] = [
authApikeyExt,
providerOpenaiCompatExt,
providerUmansExt,
+ providerConcurrencyExt,
// exec-backend must precede the tool extensions that
// `dependsOn: ["exec-backend"]` (tool-edit-file/read/shell/write). It
// provides the ExecBackendResolver the tools resolve through; placing it
diff --git a/packages/host-bin/tsconfig.json b/packages/host-bin/tsconfig.json
index 305274c..09b87df 100644
--- a/packages/host-bin/tsconfig.json
+++ b/packages/host-bin/tsconfig.json
@@ -23,6 +23,9 @@
"path": "../message-queue"
},
{
+ "path": "../provider-concurrency"
+ },
+ {
"path": "../skills"
},
{
diff --git a/packages/provider-concurrency/package.json b/packages/provider-concurrency/package.json
new file mode 100644
index 0000000..10c522a
--- /dev/null
+++ b/packages/provider-concurrency/package.json
@@ -0,0 +1,11 @@
+{
+ "name": "@dispatch/provider-concurrency",
+ "version": "0.0.0",
+ "type": "module",
+ "private": true,
+ "main": "dist/index.js",
+ "types": "dist/index.d.ts",
+ "dependencies": {
+ "@dispatch/kernel": "workspace:*"
+ }
+}
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..62f202f
--- /dev/null
+++ b/packages/provider-concurrency/src/concurrency-manager.test.ts
@@ -0,0 +1,487 @@
+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(opts?: { releaseCooldownMs?: number }): {
+ manager: ConcurrencyService;
+ timers: ReturnType<typeof createFakeTimers>;
+} {
+ const timers = createFakeTimers();
+ const manager = createConcurrencyManager({
+ now: timers.now,
+ slotTimeoutMs: 5000,
+ watchdogIntervalMs: 1000,
+ defaultPauseMs: 30000,
+ ...(opts?.releaseCooldownMs !== undefined ? { releaseCooldownMs: opts.releaseCooldownMs } : {}),
+ 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);
+ });
+
+ it("release cooldown delays slot recycling (inFlight stays incremented during cooldown)", async () => {
+ const { manager, timers } = createManager({ releaseCooldownMs: 200 });
+ manager.setLimit("umans", 1);
+
+ const release1 = await manager.acquire("umans", "conv1", 0);
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+
+ // Queue a waiter.
+ let resolved = false;
+ const promise2 = manager.acquire("umans", "conv2", 100).then((r) => {
+ resolved = true;
+ return r;
+ });
+ await Promise.resolve();
+ await Promise.resolve();
+ expect(resolved).toBe(false);
+ expect(manager.getStatus("umans")?.queued).toBe(1);
+
+ // Release the slot — inFlight should stay 1 (cooldown active).
+ release1();
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+ expect(resolved).toBe(false); // waiter NOT granted yet
+
+ // Advance past the cooldown.
+ timers.advance(200);
+
+ // Now the slot is recycled and the waiter is granted.
+ const release2 = await promise2;
+ expect(resolved).toBe(true);
+ expect(manager.getStatus("umans")?.inFlight).toBe(1);
+ expect(manager.getStatus("umans")?.queued).toBe(0);
+ release2();
+ });
+
+ it("release cooldown is idempotent (double-release only schedules one cooldown)", async () => {
+ const { manager, timers } = createManager({ releaseCooldownMs: 200 });
+ 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(1); // still 1 (cooldown)
+
+ // Double-release should not schedule a second cooldown.
+ release();
+
+ // After cooldown, inFlight should drop by exactly 1 (to 0), not 2.
+ timers.advance(200);
+ expect(manager.getStatus("umans")?.inFlight).toBe(0);
+ });
+
+ it("destroy clears cooldown timers without error", () => {
+ const { manager } = createManager({ releaseCooldownMs: 200 });
+ manager.setLimit("umans", 1);
+ // Acquire + release to schedule a cooldown timer.
+ manager.acquire("umans", "conv1", 0).then((release) => {
+ release();
+ // Now there's a pending cooldown timer — destroy should clean it up.
+ expect(() => manager.destroy()).not.toThrow();
+ });
+ });
+
+ it("onQueued is called when the request is enqueued (not granted immediately)", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 1);
+
+ // Hold the single slot.
+ const release1 = await manager.acquire("umans", "conv1", 0);
+
+ // Second request should trigger onQueued.
+ let queuedCalled = false;
+ const promise = manager.acquire("umans", "conv2", 100, () => {
+ queuedCalled = true;
+ });
+ await Promise.resolve();
+ await Promise.resolve();
+
+ expect(queuedCalled).toBe(true);
+ expect(manager.getStatus("umans")?.queued).toBe(1);
+
+ // Release the slot — the queued request should be granted.
+ release1();
+ const release2 = await promise;
+ release2();
+ });
+
+ it("onQueued is NOT called when the slot is granted immediately", async () => {
+ const { manager } = createManager();
+ manager.setLimit("umans", 2);
+
+ let queuedCalled = false;
+ const release = await manager.acquire("umans", "conv1", 0, () => {
+ queuedCalled = true;
+ });
+
+ expect(queuedCalled).toBe(false);
+ release();
+ });
+});
diff --git a/packages/provider-concurrency/src/concurrency-manager.ts b/packages/provider-concurrency/src/concurrency-manager.ts
new file mode 100644
index 0000000..d4516b2
--- /dev/null
+++ b/packages/provider-concurrency/src/concurrency-manager.ts
@@ -0,0 +1,376 @@
+/**
+ * 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.
+ *
+ * If `onQueued` is provided and the request cannot be granted immediately
+ * (at limit or paused), it is called synchronously BEFORE the Promise is
+ * created. This lets the caller emit a "queued" status signal. If the slot
+ * is granted immediately, `onQueued` is NOT called.
+ *
+ * @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.
+ * @param onQueued Called synchronously when the request is enqueued
+ * (not granted immediately). Optional.
+ */
+ acquire(
+ providerId: string,
+ conversationId: string,
+ promptStartedAt: number,
+ onQueued?: () => void,
+ ): 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;
+ /**
+ * Delay after a slot is released before the slot is recycled (ms). During
+ * this window `inFlight` stays incremented — a new `acquire` sees the slot
+ * as still held and queues. This covers the upstream provider's accounting
+ * lag: the provider's `concurrent_sessions` counter may not decrement the
+ * instant our stream completes, so re-admitting immediately risks an N+1
+ * overshoot. 0 = instant re-admission (no cooldown). Default: 0.
+ */
+ readonly releaseCooldownMs?: 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 releaseCooldownMs = opts.releaseCooldownMs ?? 0;
+ 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>();
+ const cooldownTimers = new Set<ReturnType<typeof setTimeout>>();
+ 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);
+
+ // Recycle the slot: decrement inFlight + grant the next waiter.
+ // With a release cooldown > 0, defer this by the cooldown duration so
+ // the upstream provider has time to decrement its concurrent_sessions
+ // counter — preventing an N+1 overshoot from accounting lag. During the
+ // cooldown, inFlight stays incremented, so new acquires queue.
+ const recycle = () => {
+ state.inFlight--;
+ tryGrantNext(providerId);
+ };
+ if (releaseCooldownMs > 0) {
+ const timer = setTimeout(() => {
+ cooldownTimers.delete(timer);
+ recycle();
+ }, releaseCooldownMs);
+ cooldownTimers.add(timer);
+ } else {
+ recycle();
+ }
+ };
+ 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, onQueued) {
+ 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));
+ }
+
+ // Cannot grant immediately — the request will be queued.
+ // Notify the caller BEFORE creating the Promise so they can emit a
+ // "queued" status signal while we're still synchronous.
+ onQueued?.();
+
+ // 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 timer of cooldownTimers) {
+ clearTimeout(timer);
+ }
+ cooldownTimers.clear();
+ 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..4af4f9a
--- /dev/null
+++ b/packages/provider-concurrency/src/extension.ts
@@ -0,0 +1,140 @@
+import type { Extension, HostAPI, Logger, Manifest, StorageNamespace } from "@dispatch/kernel";
+import type { ConcurrencyManagerOpts, ConcurrencyService } from "./concurrency-manager.js";
+import { 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",
+ capabilities: { db: true },
+ 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.
+ * - `RELEASE_COOLDOWN_MS` (200ms): after a slot is released, hold it for this
+ * duration before recycling it to the next waiter. Covers the upstream
+ * provider's accounting lag — the provider's concurrent_sessions counter
+ * may not decrement the instant our stream completes, so re-admitting
+ * immediately risks an N+1 overshoot that triggers a 429. 200ms is the
+ * default most concurrency proxies use for AI/LLM APIs.
+ */
+const SLOT_TIMEOUT_MS = 5 * 60 * 1000;
+const WATCHDOG_INTERVAL_MS = 30 * 1000;
+const DEFAULT_PAUSE_MS = 30 * 1000;
+const RELEASE_COOLDOWN_MS = 200;
+
+/**
+ * Wrap a `ConcurrencyService` so `setLimit`/`removeLimit` persist to the
+ * given `StorageNamespace`. All other methods delegate directly to the inner
+ * service. Persistence is fire-and-forget — a storage write failure logs a
+ * warning but does NOT fail the API call (the in-memory limit is already set).
+ */
+function createPersistedService(
+ inner: ConcurrencyService,
+ storage: StorageNamespace,
+ logger: Logger,
+): ConcurrencyService {
+ return {
+ acquire: inner.acquire.bind(inner),
+ reportRateLimit: inner.reportRateLimit.bind(inner),
+ setLimit(providerId, limit) {
+ inner.setLimit(providerId, limit);
+ storage.set(providerId, String(limit)).catch((err) =>
+ logger.warn("provider-concurrency: failed to persist limit", {
+ providerId,
+ err: err instanceof Error ? err.message : String(err),
+ }),
+ );
+ },
+ removeLimit(providerId) {
+ inner.removeLimit(providerId);
+ storage.delete(providerId).catch((err) =>
+ logger.warn("provider-concurrency: failed to delete persisted limit", {
+ providerId,
+ err: err instanceof Error ? err.message : String(err),
+ }),
+ );
+ },
+ getLimit: inner.getLimit.bind(inner),
+ getLimits: inner.getLimits.bind(inner),
+ getStatus: inner.getStatus.bind(inner),
+ getStatusAll: inner.getStatusAll.bind(inner),
+ destroy: inner.destroy.bind(inner),
+ };
+}
+
+/**
+ * Load saved limits from storage and apply them to the manager.
+ * Called during activate, before the service is registered.
+ */
+async function loadLimits(
+ storage: StorageNamespace,
+ manager: ConcurrencyService,
+ logger: Logger,
+): Promise<void> {
+ const keys = await storage.keys();
+ for (const providerId of keys) {
+ const raw = await storage.get(providerId);
+ if (raw === null) continue;
+ const limit = Number.parseInt(raw, 10);
+ if (!Number.isNaN(limit) && limit > 0) {
+ manager.setLimit(providerId, limit);
+ logger.info(`provider-concurrency: restored limit ${limit} for "${providerId}"`);
+ }
+ }
+}
+
+export async function activate(host: HostAPI): Promise<void> {
+ const logger = host.logger;
+ const storage = host.storage("provider-concurrency");
+
+ const managerOpts: ConcurrencyManagerOpts = {
+ now: () => Date.now(),
+ slotTimeoutMs: SLOT_TIMEOUT_MS,
+ watchdogIntervalMs: WATCHDOG_INTERVAL_MS,
+ defaultPauseMs: DEFAULT_PAUSE_MS,
+ releaseCooldownMs: RELEASE_COOLDOWN_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,
+ });
+ },
+ };
+
+ const inner = createConcurrencyManager(managerOpts);
+
+ // Restore persisted limits before registering the service so the first
+ // request sees the correct configuration.
+ await loadLimits(storage, inner, logger);
+
+ const service = createPersistedService(inner, storage, logger);
+ host.provideService(concurrencyServiceHandle, service);
+ 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..e59ab39
--- /dev/null
+++ b/packages/provider-concurrency/src/provider-wrapper.test.ts
@@ -0,0 +1,246 @@
+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("calls onQueued when the request blocks and onAcquired when the slot is granted", async () => {
+ let queuedCalled = false;
+ let acquiredCalled = false;
+
+ const blockingLimiter: ConcurrencyLimiter = {
+ acquire(_providerId, _convId, _promptAt, onQueued) {
+ // Simulate a queued request: call onQueued, then resolve on next tick.
+ onQueued?.();
+ return new Promise((resolve) => {
+ setTimeout(() => {
+ resolve(() => {});
+ }, 0);
+ });
+ },
+ reportRateLimit() {},
+ };
+
+ const provider = fakeProvider([{ type: "finish", reason: "stop" }]);
+ const wrapped = wrapProviderWithConcurrency(
+ provider,
+ blockingLimiter,
+ "conv1",
+ 0,
+ () => {
+ queuedCalled = true;
+ },
+ () => {
+ acquiredCalled = true;
+ },
+ );
+
+ for await (const _e of wrapped.stream([], [])) {
+ // consume
+ }
+
+ expect(queuedCalled).toBe(true);
+ expect(acquiredCalled).toBe(true);
+ });
+
+ it("does NOT call onQueued when the slot is granted immediately", async () => {
+ let queuedCalled = false;
+ let acquiredCalled = false;
+
+ const immediateLimiter: ConcurrencyLimiter = {
+ acquire(_providerId, _convId, _promptAt, _onQueued) {
+ // Grant immediately — do NOT call onQueued.
+ return Promise.resolve(() => {});
+ },
+ reportRateLimit() {},
+ };
+
+ const provider = fakeProvider([{ type: "finish", reason: "stop" }]);
+ const wrapped = wrapProviderWithConcurrency(
+ provider,
+ immediateLimiter,
+ "conv1",
+ 0,
+ () => {
+ queuedCalled = true;
+ },
+ () => {
+ acquiredCalled = true;
+ },
+ );
+
+ for await (const _e of wrapped.stream([], [])) {
+ // consume
+ }
+
+ expect(queuedCalled).toBe(false);
+ expect(acquiredCalled).toBe(true);
+ });
+
+ 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..aa08e5b
--- /dev/null
+++ b/packages/provider-concurrency/src/provider-wrapper.ts
@@ -0,0 +1,68 @@
+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).
+ * @param onQueued Called synchronously when `acquire()` decides to
+ * queue the request (cannot grant immediately).
+ * Lets the caller emit a "queued" status signal.
+ * @param onAcquired Called when `acquire()` resolves (slot granted,
+ * whether immediately or after queueing). Lets the
+ * caller emit an "active" status signal.
+ */
+export function wrapProviderWithConcurrency(
+ provider: ProviderContract,
+ limiter: ConcurrencyLimiter,
+ conversationId: string,
+ promptStartedAt: number,
+ onQueued?: () => void,
+ onAcquired?: () => void,
+): 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, onQueued);
+ onAcquired?.();
+ 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",
+);
diff --git a/packages/provider-concurrency/tsconfig.json b/packages/provider-concurrency/tsconfig.json
new file mode 100644
index 0000000..44ed916
--- /dev/null
+++ b/packages/provider-concurrency/tsconfig.json
@@ -0,0 +1,6 @@
+{
+ "extends": "../../tsconfig.base.json",
+ "compilerOptions": { "rootDir": "src", "outDir": "dist", "composite": true },
+ "include": ["src/**/*.ts"],
+ "references": [{ "path": "../kernel" }]
+}
diff --git a/packages/session-orchestrator/package.json b/packages/session-orchestrator/package.json
index ba34c4d..b9f3d22 100644
--- a/packages/session-orchestrator/package.json
+++ b/packages/session-orchestrator/package.json
@@ -10,6 +10,7 @@
"@dispatch/conversation-store": "workspace:*",
"@dispatch/credential-store": "workspace:*",
"@dispatch/message-queue": "workspace:*",
+ "@dispatch/provider-concurrency": "workspace:*",
"@dispatch/system-prompt": "workspace:*"
}
}
diff --git a/packages/session-orchestrator/src/extension.ts b/packages/session-orchestrator/src/extension.ts
index d080e90..783d894 100644
--- a/packages/session-orchestrator/src/extension.ts
+++ b/packages/session-orchestrator/src/extension.ts
@@ -3,6 +3,7 @@ import { credentialStoreHandle } from "@dispatch/credential-store";
import type { Extension, HostAPI, Manifest } from "@dispatch/kernel";
import { runTurn } from "@dispatch/kernel";
import { messageQueueHandle } from "@dispatch/message-queue";
+import { concurrencyServiceHandle } from "@dispatch/provider-concurrency";
import { systemPromptHandle } from "@dispatch/system-prompt";
import {
cacheWarmHandle,
@@ -94,6 +95,19 @@ export function activate(host: HostAPI): void {
return undefined;
}
},
+ resolveConcurrencyLimiter: () => {
+ // Lazily resolve the concurrency limiter. Returns undefined when the
+ // provider-concurrency extension isn't loaded (no concurrency limiting —
+ // feature degrades off). Lazy so activation order with
+ // provider-concurrency doesn't matter; called per-turn, not at activate.
+ const loaded = host.getExtensions().some((m) => m.id === "provider-concurrency");
+ if (!loaded) return undefined;
+ try {
+ return host.getService(concurrencyServiceHandle);
+ } catch {
+ return undefined;
+ }
+ },
resolveVisionHandoff: () => {
// Lazily resolve the vision-handoff service. Returns undefined when the
// vision-handoff extension isn't loaded (images pass through unchanged —
diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts
index 045b88d..5c36922 100644
--- a/packages/session-orchestrator/src/orchestrator.ts
+++ b/packages/session-orchestrator/src/orchestrator.ts
@@ -21,6 +21,8 @@ import type {
} from "@dispatch/kernel";
import { defineEventHook, defineService, type ServiceHandle } from "@dispatch/kernel";
import type { MessageQueueService, QueuedMessage } from "@dispatch/message-queue";
+import type { ConcurrencyLimiter } from "@dispatch/provider-concurrency";
+import { wrapProviderWithConcurrency } from "@dispatch/provider-concurrency";
import type { SystemPromptService } from "@dispatch/system-prompt";
import { createMetricsAccumulator } from "./metrics.js";
import {
@@ -405,6 +407,14 @@ export interface SessionOrchestratorDeps {
*/
readonly resolveSystemPrompt?: () => SystemPromptService | undefined;
/**
+ * Lazily resolves the concurrency limiter, or `undefined` when the
+ * provider-concurrency extension isn't loaded (no concurrency limiting —
+ * feature degrades off). When present, each resolved provider is wrapped so
+ * that a concurrency slot is acquired before the stream starts and released
+ * when the stream completes. Lazy so activation order doesn't matter.
+ */
+ readonly resolveConcurrencyLimiter?: () => ConcurrencyLimiter | undefined;
+ /**
* Lazily resolves the vision-handoff service, or `undefined` when the
* vision-handoff extension isn't loaded. Used to transcribe image chunks to
* text for non-vision models before they reach the provider (so a text-only
@@ -520,6 +530,7 @@ export function createSessionOrchestrator(
images: readonly ImageInput[] | undefined,
): void {
const turnId = generateTurnId();
+ const promptStartedAt = deps.now?.() ?? Date.now();
const controller = new AbortController();
activeTurns.set(conversationId, { buffer: [], turnId, controller });
activeConversations.add(conversationId);
@@ -686,6 +697,42 @@ export function createSessionOrchestrator(
provider = deps.resolveProvider();
}
+ // Wrap the resolved provider with concurrency limiting when the
+ // provider-concurrency extension is loaded. The slot is acquired
+ // before the stream starts (before the HTTP request) and released
+ // when the stream completes (after all tokens are generated). The
+ // promptStartedAt (turn start time) is used for oldest-agent-first
+ // scheduling when multiple agents are queued.
+ //
+ // Status lifecycle with concurrency: "active" is emitted early (in
+ // payloadPromise.then, before this code runs). If acquire() blocks,
+ // onQueued emits "queued" (broadcast-only — persisted status stays
+ // "active"). When the slot is granted, onAcquired emits "active"
+ // again, transitioning "queued" → "active" so the FE switches from
+ // the loading ring back to dots. A request that gets a slot
+ // immediately never emits "queued" — onAcquired fires right after
+ // the early "active", which is a harmless no-op re-broadcast.
+ const limiter = deps.resolveConcurrencyLimiter?.();
+ if (limiter !== undefined) {
+ const emitStatus = (status: "queued" | "active"): void => {
+ void deps.conversationStore.getWorkspaceId(conversationId).then((workspaceId) => {
+ deps.emit?.(conversationStatusChanged, {
+ conversationId,
+ status,
+ workspaceId,
+ });
+ });
+ };
+ provider = wrapProviderWithConcurrency(
+ provider,
+ limiter,
+ conversationId,
+ promptStartedAt,
+ () => emitStatus("queued"),
+ () => emitStatus("active"),
+ );
+ }
+
const baseTools = deps.resolveTools();
const assembled = await deps.applyToolsFilter({
tools: baseTools,
@@ -1116,6 +1163,17 @@ export function createWarmService(
provider = deps.resolveProvider();
}
+ // Wrap with concurrency limiting (same as the main turn path).
+ const warmLimiter = deps.resolveConcurrencyLimiter?.();
+ if (warmLimiter !== undefined) {
+ provider = wrapProviderWithConcurrency(
+ provider,
+ warmLimiter,
+ conversationId,
+ deps.now?.() ?? Date.now(),
+ );
+ }
+
const baseTools = deps.resolveTools();
// Resolve cwd the SAME way handleMessage does — pass opts.cwd as the overrideCwd
// The tools filter is cwd-sensitive (e.g. skill discovery rewrites the
@@ -1262,6 +1320,17 @@ export function createCompactionService(
provider = deps.resolveProvider();
}
+ // Wrap with concurrency limiting (same as the main turn path).
+ const compactionLimiter = deps.resolveConcurrencyLimiter?.();
+ if (compactionLimiter !== undefined) {
+ provider = wrapProviderWithConcurrency(
+ provider,
+ compactionLimiter,
+ conversationId,
+ deps.now?.() ?? Date.now(),
+ );
+ }
+
// Build the summarization request: system prompt + conversation text + instruction
const conversationText = formatMessagesForSummary(toSummarize);
const summaryRequest: ChatMessage = {
diff --git a/packages/session-orchestrator/tsconfig.json b/packages/session-orchestrator/tsconfig.json
index bc729fc..f316387 100644
--- a/packages/session-orchestrator/tsconfig.json
+++ b/packages/session-orchestrator/tsconfig.json
@@ -7,6 +7,7 @@
{ "path": "../conversation-store" },
{ "path": "../credential-store" },
{ "path": "../message-queue" },
+ { "path": "../provider-concurrency" },
{ "path": "../system-prompt" }
]
}
diff --git a/packages/transport-contract/src/index.ts b/packages/transport-contract/src/index.ts
index 94897f7..d5f3000 100644
--- a/packages/transport-contract/src/index.ts
+++ b/packages/transport-contract/src/index.ts
@@ -1026,3 +1026,56 @@ export interface HeartbeatRunsResponse {
export interface StopHeartbeatRunResponse {
readonly ok: true;
}
+
+// ─── Provider concurrency limits ──────────────────────────────────────────────
+
+/**
+ * Response of `GET /concurrency/limits` — all providers with configured
+ * concurrency limits. Each entry pairs a provider id (e.g. "umans",
+ * "openai-compat") with its maximum concurrent in-flight requests. Providers
+ * not listed here have no limit (unlimited).
+ */
+export interface ConcurrencyLimitsResponse {
+ readonly limits: readonly { readonly providerId: string; readonly limit: number }[];
+}
+
+/**
+ * Body of `PUT /concurrency/limits/:providerId` — set or update the concurrency
+ * limit for a provider. `limit` must be a positive integer. When a limit is
+ * set, requests beyond the limit queue (oldest-agent-first) rather than being
+ * sent immediately.
+ */
+export interface SetConcurrencyLimitRequest {
+ readonly limit: number;
+}
+
+/** Response of `GET/PUT /concurrency/limits/:providerId` — the configured limit. */
+export interface ConcurrencyLimitResponse {
+ readonly providerId: string;
+ readonly limit: number;
+}
+
+/**
+ * One provider's live concurrency status.
+ *
+ * - `inFlight`: how many slots are currently held (tokens being generated).
+ * - `queued`: how many agents are waiting for a slot.
+ * - `paused`: whether the queue is paused due to a 429 backoff.
+ * - `pausedUntil`: when the pause expires (epoch-ms), present only when paused.
+ */
+export interface ConcurrencyStatusEntry {
+ readonly providerId: string;
+ readonly limit: number;
+ readonly inFlight: number;
+ readonly queued: number;
+ readonly paused: boolean;
+ readonly pausedUntil?: number;
+}
+
+/**
+ * Response of `GET /concurrency/status` — live status for every provider with a
+ * configured limit. Providers without a limit are absent (they are unlimited).
+ */
+export interface ConcurrencyStatusResponse {
+ readonly providers: readonly ConcurrencyStatusEntry[];
+}
diff --git a/packages/transport-http/package.json b/packages/transport-http/package.json
index 3c722e2..71da855 100644
--- a/packages/transport-http/package.json
+++ b/packages/transport-http/package.json
@@ -12,6 +12,7 @@
"@dispatch/kernel": "workspace:*",
"@dispatch/lsp": "workspace:*",
"@dispatch/mcp": "workspace:*",
+ "@dispatch/provider-concurrency": "workspace:*",
"@dispatch/session-orchestrator": "workspace:*",
"@dispatch/throughput-store": "workspace:*",
"@dispatch/transport-contract": "workspace:*",
diff --git a/packages/transport-http/src/app.ts b/packages/transport-http/src/app.ts
index 16c4167..0fcc8f0 100644
--- a/packages/transport-http/src/app.ts
+++ b/packages/transport-http/src/app.ts
@@ -8,6 +8,9 @@ import type {
ComputerListResponse,
ComputerResponse,
ComputerStatusResponse,
+ ConcurrencyLimitResponse,
+ ConcurrencyLimitsResponse,
+ ConcurrencyStatusResponse,
ConversationComputerResponse,
ConversationHistoryResponse,
ConversationListResponse,
@@ -28,6 +31,7 @@ import type {
QueueResponse,
ReasoningEffortResponse,
SetCompactPercentRequest,
+ SetConcurrencyLimitRequest,
SetConversationComputerRequest,
SetSystemPromptTemplateRequest,
SetWorkspaceDefaultComputerRequest,
@@ -68,6 +72,7 @@ import {
import {
type CompactionService,
type ComputerService,
+ type ConcurrencyService,
type ConversationStore,
type CredentialStore,
conversationOpened,
@@ -111,6 +116,13 @@ export interface CreateServerOptions {
readonly computerService?: ComputerService;
/** Optional — defaults to a no-op store (recording disabled, empty reports). */
readonly throughputStore?: ThroughputStore;
+ /**
+ * Optional — provider concurrency limiter service (provided by the
+ * `provider-concurrency` extension). When absent (extension not loaded),
+ * the `/concurrency/*` routes degrade: limits returns empty, status returns
+ * empty, PUT returns 503.
+ */
+ readonly concurrencyService?: ConcurrencyService;
readonly logger?: Logger;
readonly generateId?: () => string;
/** Injectable clock for sample timestamps (default Date.now). */
@@ -607,6 +619,84 @@ export function createApp(opts: CreateServerOptions): Hono {
}
});
+ // ─── Provider concurrency limits ────────────────────────────────────────────
+
+ app.get("/concurrency/limits", (c) => {
+ if (opts.concurrencyService === undefined) {
+ const body: ConcurrencyLimitsResponse = { limits: [] };
+ return c.json(body, 200);
+ }
+ const limits = opts.concurrencyService.getLimits();
+ const body: ConcurrencyLimitsResponse = { limits };
+ return c.json(body, 200);
+ });
+
+ app.get("/concurrency/limits/:providerId", (c) => {
+ const providerId = c.req.param("providerId");
+ if (opts.concurrencyService === undefined) {
+ return c.json({ error: "Concurrency service not available" }, 503);
+ }
+ const limit = opts.concurrencyService.getLimit(providerId);
+ if (limit === undefined) {
+ return c.json({ error: "No concurrency limit configured for this provider" }, 404);
+ }
+ const body: ConcurrencyLimitResponse = { providerId, limit };
+ return c.json(body, 200);
+ });
+
+ app.put("/concurrency/limits/:providerId", async (c) => {
+ const providerId = c.req.param("providerId");
+ if (opts.concurrencyService === undefined) {
+ return c.json({ error: "Concurrency service not available" }, 503);
+ }
+
+ let body: unknown;
+ try {
+ body = await c.req.json();
+ } catch {
+ log.warn("concurrency: invalid JSON body");
+ return c.json({ error: "Invalid JSON body" }, 400);
+ }
+
+ const parsed = body as SetConcurrencyLimitRequest;
+ if (
+ parsed === null ||
+ typeof parsed !== "object" ||
+ typeof parsed.limit !== "number" ||
+ !Number.isInteger(parsed.limit) ||
+ parsed.limit <= 0
+ ) {
+ return c.json({ error: "Body must be { limit: <positive integer> }" }, 400);
+ }
+
+ opts.concurrencyService.setLimit(providerId, parsed.limit);
+ const responseBody: ConcurrencyLimitResponse = { providerId, limit: parsed.limit };
+ return c.json(responseBody, 200);
+ });
+
+ app.delete("/concurrency/limits/:providerId", (c) => {
+ const providerId = c.req.param("providerId");
+ if (opts.concurrencyService === undefined) {
+ return c.json({ error: "Concurrency service not available" }, 503);
+ }
+ const existing = opts.concurrencyService.getLimit(providerId);
+ if (existing === undefined) {
+ return c.json({ error: "No concurrency limit configured for this provider" }, 404);
+ }
+ opts.concurrencyService.removeLimit(providerId);
+ return c.json({ ok: true, providerId }, 200);
+ });
+
+ app.get("/concurrency/status", (c) => {
+ if (opts.concurrencyService === undefined) {
+ const body: ConcurrencyStatusResponse = { providers: [] };
+ return c.json(body, 200);
+ }
+ const statuses = opts.concurrencyService.getStatusAll();
+ const body: ConcurrencyStatusResponse = { providers: statuses };
+ return c.json(body, 200);
+ });
+
app.post("/conversations/:id/close", (c) => {
const conversationId = c.req.param("id");
const { abortedTurn } = opts.orchestrator.closeConversation(conversationId);
diff --git a/packages/transport-http/src/extension.ts b/packages/transport-http/src/extension.ts
index 76d58b6..f46fca5 100644
--- a/packages/transport-http/src/extension.ts
+++ b/packages/transport-http/src/extension.ts
@@ -2,9 +2,11 @@ import type { Extension, HostAPI, Manifest } from "@dispatch/kernel";
import { createApp } from "./app.js";
import {
type ComputerService,
+ type ConcurrencyService,
cacheWarmHandle,
compactionHandle,
computerServiceHandle,
+ concurrencyServiceHandle,
conversationStoreHandle,
credentialStoreHandle,
heartbeatServiceHandle,
@@ -39,6 +41,9 @@ export const manifest: Manifest = {
"/computers/:alias",
"/computers/:alias/status",
"/computers/:alias/test",
+ "/concurrency/limits",
+ "/concurrency/limits/:providerId",
+ "/concurrency/status",
"/conversations",
"/conversations/:id",
"/conversations/:id/close",
@@ -105,6 +110,15 @@ export function createTransportHttpExtension(): Extension & {
} catch {
computerService = undefined;
}
+ // Optional: the `provider-concurrency` extension provides the
+ // concurrency limiter service. NOT in dependsOn (may be absent), so
+ // resolve defensively — when absent the /concurrency/* routes degrade.
+ let concurrencyService: ConcurrencyService | undefined;
+ try {
+ concurrencyService = host.getService(concurrencyServiceHandle);
+ } catch {
+ concurrencyService = undefined;
+ }
const logger = host.logger;
const app = createApp({
@@ -119,6 +133,7 @@ export function createTransportHttpExtension(): Extension & {
systemPromptService,
heartbeatService,
...(computerService !== undefined ? { computerService } : {}),
+ ...(concurrencyService !== undefined ? { concurrencyService } : {}),
logger,
emit: host.emit.bind(host),
...(process.env.DISPATCH_WEB_DIR !== undefined
diff --git a/packages/transport-http/src/logic.ts b/packages/transport-http/src/logic.ts
index a928147..c97f320 100644
--- a/packages/transport-http/src/logic.ts
+++ b/packages/transport-http/src/logic.ts
@@ -13,7 +13,7 @@ const VALID_REASONING_EFFORTS: readonly ReasoningEffort[] = [
"max",
];
-const VALID_STATUSES: readonly ConversationStatus[] = ["active", "idle", "closed"];
+const VALID_STATUSES: readonly ConversationStatus[] = ["active", "queued", "idle", "closed"];
/**
* Pure: parse a `?status=` query value into a list of valid ConversationStatus
diff --git a/packages/transport-http/src/seam.ts b/packages/transport-http/src/seam.ts
index c60edf0..dcb3f80 100644
--- a/packages/transport-http/src/seam.ts
+++ b/packages/transport-http/src/seam.ts
@@ -12,6 +12,8 @@ export type { LspServerStatus, LspService } from "@dispatch/lsp";
export { lspServiceHandle } from "@dispatch/lsp";
export type { McpServerStatus, McpService } from "@dispatch/mcp";
export { mcpServiceHandle } from "@dispatch/mcp";
+export type { ConcurrencyService } from "@dispatch/provider-concurrency";
+export { concurrencyServiceHandle } from "@dispatch/provider-concurrency";
export type {
CompactionService,
SessionOrchestrator,
diff --git a/packages/transport-http/tsconfig.json b/packages/transport-http/tsconfig.json
index 4a76434..7759bce 100644
--- a/packages/transport-http/tsconfig.json
+++ b/packages/transport-http/tsconfig.json
@@ -8,6 +8,8 @@
{ "path": "../heartbeat" },
{ "path": "../kernel" },
{ "path": "../lsp" },
+ { "path": "../mcp" },
+ { "path": "../provider-concurrency" },
{ "path": "../session-orchestrator" },
{ "path": "../system-prompt" },
{ "path": "../throughput-store" },
diff --git a/packages/wire/src/index.ts b/packages/wire/src/index.ts
index d6ea1c1..113f684 100644
--- a/packages/wire/src/index.ts
+++ b/packages/wire/src/index.ts
@@ -568,12 +568,18 @@ export interface TurnSteeringEvent {
/**
* The lifecycle status of a conversation, used for tab persistence across
- * devices. `active` = an agent is currently generating; `idle` = exists but not
+ * devices. `active` = an agent is currently generating; `queued` = the
+ * request is waiting in the per-provider concurrency queue for a slot (not yet
+ * generating — the FE shows a loading ring, not dots); `idle` = exists but not
* generating; `closed` = user dismissed the tab (hidden from the tab bar, not
* deleted). New conversations start as `idle`; transitions to `active` on
* turn-start, back to `idle` on turn done/error, and to `closed` on user close.
+ * When the concurrency extension is loaded and a request blocks on
+ * `acquire()`, `queued` is broadcast (persisted status stays `active`); when
+ * the slot is granted, `active` is re-broadcast. A request that gets a slot
+ * immediately never emits `queued`.
*/
-export type ConversationStatus = "active" | "idle" | "closed";
+export type ConversationStatus = "active" | "queued" | "idle" | "closed";
/**
* Metadata for a conversation, returned by `GET /conversations` (the list
diff --git a/tsconfig.json b/tsconfig.json
index fe5ea92..f97edde 100644
--- a/tsconfig.json
+++ b/tsconfig.json
@@ -101,6 +101,9 @@
"path": "./packages/heartbeat"
},
{
+ "path": "./packages/provider-concurrency"
+ },
+ {
"path": "./packages/system-prompt"
},
{