diff options
27 files changed, 1861 insertions, 3 deletions
@@ -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" }, { |
