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/session-orchestrator/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/session-orchestrator/src')
| -rw-r--r-- | packages/session-orchestrator/src/extension.ts | 14 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.ts | 49 |
2 files changed, 63 insertions, 0 deletions
diff --git a/packages/session-orchestrator/src/extension.ts b/packages/session-orchestrator/src/extension.ts index 5afffd8..0cd83ef 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, @@ -93,6 +94,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; + } + }, }); host.provideService(sessionOrchestratorHandle, orchestrator); diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts index 96cd3a3..b73647d 100644 --- a/packages/session-orchestrator/src/orchestrator.ts +++ b/packages/session-orchestrator/src/orchestrator.ts @@ -20,6 +20,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 { @@ -335,6 +337,14 @@ export interface SessionOrchestratorDeps { * order doesn't matter. */ 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; /** Apply the per-turn tools filter chain. Injected for testability. */ readonly applyToolsFilter: (assembly: ToolAssembly) => Promise<ToolAssembly>; /** Base logger (auto-scoped to this extension); childed per turn for span capture. */ @@ -439,6 +449,7 @@ export function createSessionOrchestrator( systemPromptOverride: string | undefined, ): void { const turnId = generateTurnId(); + const promptStartedAt = deps.now?.() ?? Date.now(); const controller = new AbortController(); activeTurns.set(conversationId, { buffer: [], turnId, controller }); activeConversations.add(conversationId); @@ -594,6 +605,22 @@ 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. + const limiter = deps.resolveConcurrencyLimiter?.(); + if (limiter !== undefined) { + provider = wrapProviderWithConcurrency( + provider, + limiter, + conversationId, + promptStartedAt, + ); + } + const baseTools = deps.resolveTools(); const assembled = await deps.applyToolsFilter({ tools: baseTools, @@ -990,6 +1017,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 @@ -1136,6 +1174,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 = { |
