summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/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/session-orchestrator/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/session-orchestrator/src')
-rw-r--r--packages/session-orchestrator/src/extension.ts14
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts49
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 = {