summaryrefslogtreecommitdiffhomepage
path: root/packages/transport-http/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-27 03:03:53 +0900
committerAdam Malczewski <[email protected]>2026-06-27 03:03:53 +0900
commita6b95188a110464b6ffa0334c8af58463f2a36f2 (patch)
treeeb6ef57909e164be4ae721ea1fb25585354d351e /packages/transport-http/src
parentad9d135e583c99a0d93327115defa43187cde1c3 (diff)
downloaddispatch-a6b95188a110464b6ffa0334c8af58463f2a36f2.tar.gz
dispatch-a6b95188a110464b6ffa0334c8af58463f2a36f2.zip
feat(provider-concurrency): implement per-provider in-memory concurrency limits with oldest-agent-first scheduling
Diffstat (limited to 'packages/transport-http/src')
-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/seam.ts2
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,