diff options
| author | Adam Malczewski <[email protected]> | 2026-06-27 03:03:53 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-27 03:03:53 +0900 |
| commit | a6b95188a110464b6ffa0334c8af58463f2a36f2 (patch) | |
| tree | eb6ef57909e164be4ae721ea1fb25585354d351e /packages/transport-http/src | |
| parent | ad9d135e583c99a0d93327115defa43187cde1c3 (diff) | |
| download | dispatch-a6b95188a110464b6ffa0334c8af58463f2a36f2.tar.gz dispatch-a6b95188a110464b6ffa0334c8af58463f2a36f2.zip | |
feat(provider-concurrency): implement per-provider in-memory concurrency limits with oldest-agent-first scheduling
Diffstat (limited to 'packages/transport-http/src')
| -rw-r--r-- | packages/transport-http/src/app.ts | 90 | ||||
| -rw-r--r-- | packages/transport-http/src/extension.ts | 15 | ||||
| -rw-r--r-- | packages/transport-http/src/seam.ts | 2 |
3 files changed, 107 insertions, 0 deletions
diff --git a/packages/transport-http/src/app.ts b/packages/transport-http/src/app.ts index 4fb295e..23f8dde 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, @@ -67,6 +71,7 @@ import { import { type CompactionService, type ComputerService, + type ConcurrencyService, type ConversationStore, type CredentialStore, conversationOpened, @@ -110,6 +115,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). */ @@ -562,6 +574,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/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, |
