summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src/extension.ts
diff options
context:
space:
mode:
Diffstat (limited to 'packages/session-orchestrator/src/extension.ts')
-rw-r--r--packages/session-orchestrator/src/extension.ts336
1 files changed, 189 insertions, 147 deletions
diff --git a/packages/session-orchestrator/src/extension.ts b/packages/session-orchestrator/src/extension.ts
index 1a57cc3..777feaa 100644
--- a/packages/session-orchestrator/src/extension.ts
+++ b/packages/session-orchestrator/src/extension.ts
@@ -3,171 +3,213 @@ 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,
- compactionHandle,
- createCompactionService,
- createSessionOrchestrator,
- createWarmService,
- sessionOrchestratorHandle,
+ cacheWarmHandle,
+ compactionHandle,
+ createCompactionService,
+ createSessionOrchestrator,
+ createWarmService,
+ sessionOrchestratorHandle,
+ visionHandoffLocalHandle,
} from "./orchestrator.js";
import { selectFirstProvider } from "./pure.js";
import { filterRemoteIncompatibleTools, toolsFilter } from "./tools-filter.js";
export const manifest: Manifest = {
- id: "session-orchestrator",
- name: "Session Orchestrator",
- version: "0.0.0",
- apiVersion: "^0.1.0",
- trust: "bundled",
- dependsOn: ["conversation-store", "credential-store"],
- activation: "eager",
- contributes: {
- services: [
- "session-orchestrator/orchestrator",
- "session-orchestrator/warm",
- "session-orchestrator/compaction",
- ],
- hooks: [
- "session-orchestrator/turn-started",
- "session-orchestrator/turn-settled",
- "session-orchestrator/warm-completed",
- "session-orchestrator/conversation-closed",
- "session-orchestrator/conversation-status-changed",
- "session-orchestrator/conversation-compacted",
- ],
- },
+ id: "session-orchestrator",
+ name: "Session Orchestrator",
+ version: "0.0.0",
+ apiVersion: "^0.1.0",
+ trust: "bundled",
+ dependsOn: ["conversation-store", "credential-store"],
+ activation: "eager",
+ contributes: {
+ services: [
+ "session-orchestrator/orchestrator",
+ "session-orchestrator/warm",
+ "session-orchestrator/compaction",
+ ],
+ hooks: [
+ "session-orchestrator/turn-started",
+ "session-orchestrator/turn-settled",
+ "session-orchestrator/warm-completed",
+ "session-orchestrator/conversation-closed",
+ "session-orchestrator/conversation-status-changed",
+ "session-orchestrator/conversation-compacted",
+ ],
+ },
};
export function activate(host: HostAPI): void {
- const conversationStore = host.getService(conversationStoreHandle);
+ const conversationStore = host.getService(conversationStoreHandle);
- const { orchestrator, activeConversations } = createSessionOrchestrator({
- conversationStore,
- resolveProvider: () => selectFirstProvider(host.getProviders()),
- resolveTools: () => [...host.getTools().values()],
- resolveModel: (modelName: string) => {
- const store = host.getService(credentialStoreHandle);
- const r = store.resolve(modelName);
- if (r === undefined) return undefined;
- const provider = host.getProviders().get(r.providerId);
- return provider ? { provider, model: r.model } : undefined;
- },
- resolveModelInfo: async (modelName: string) => {
- const store = host.getService(credentialStoreHandle);
- return store.getModelInfo(modelName);
- },
- applyToolsFilter: (assembly) => host.applyFilters(toolsFilter, assembly),
- runTurn,
- logger: host.logger,
- now: () => Date.now(),
- emit: (hook, payload) => host.emit(hook, payload),
- resolveQueue: () => {
- // Lazily resolve the message-queue service. Returns undefined when the
- // extension isn't loaded (feature degrades off) — checked via the
- // activated-manifests list so `host.getService` is only called when the
- // service is registered. Lazy so activation order with message-queue
- // doesn't matter; called per-turn / per-enqueue, not at activate time.
- const loaded = host.getExtensions().some((m) => m.id === "message-queue");
- return loaded ? host.getService(messageQueueHandle) : undefined;
- },
- resolveCompaction: () => {
- // Lazily resolve the compaction service (registered below after
- // the orchestrator). By the time this is called at runtime
- // (after a turn settles), the service is registered.
- try {
- return host.getService(compactionHandle);
- } catch {
- return undefined;
- }
- },
- resolveSystemPrompt: () => {
- // Lazily resolve the system-prompt service. Returns undefined when
- // the system-prompt extension isn't loaded (no system prompt sent —
- // current behavior). Lazy so activation order with system-prompt
- // doesn't matter; called per-turn / per-compaction, not at activate.
- try {
- return host.getService(systemPromptHandle);
- } catch {
- return undefined;
- }
- },
- });
+ const { orchestrator, activeConversations } = createSessionOrchestrator({
+ conversationStore,
+ resolveProvider: () => selectFirstProvider(host.getProviders()),
+ resolveTools: () => [...host.getTools().values()],
+ resolveModel: (modelName: string) => {
+ const store = host.getService(credentialStoreHandle);
+ const r = store.resolve(modelName);
+ if (r === undefined) return undefined;
+ const provider = host.getProviders().get(r.providerId);
+ return provider ? { provider, model: r.model } : undefined;
+ },
+ resolveModelInfo: async (modelName: string) => {
+ const store = host.getService(credentialStoreHandle);
+ return store.getModelInfo(modelName);
+ },
+ applyToolsFilter: (assembly) => host.applyFilters(toolsFilter, assembly),
+ runTurn,
+ logger: host.logger,
+ now: () => Date.now(),
+ emit: (hook, payload) => host.emit(hook, payload),
+ // Injected process.memoryUsage() sampler — the production edge. Tests
+ // inject a fake to assert per-turn before/after telemetry. Wired in the
+ // shell (like `now: () => Date.now()`); pure decision logic is untouched.
+ sampleMemory: () => {
+ const m = process.memoryUsage();
+ return {
+ rss: m.rss,
+ heapUsed: m.heapUsed,
+ heapTotal: m.heapTotal,
+ external: m.external,
+ arrayBuffers: m.arrayBuffers,
+ };
+ },
+ resolveQueue: () => {
+ // Lazily resolve the message-queue service. Returns undefined when the
+ // extension isn't loaded (feature degrades off) — checked via the
+ // activated-manifests list so `host.getService` is only called when the
+ // service is registered. Lazy so activation order with message-queue
+ // doesn't matter; called per-turn / per-enqueue, not at activate time.
+ const loaded = host.getExtensions().some((m) => m.id === "message-queue");
+ return loaded ? host.getService(messageQueueHandle) : undefined;
+ },
+ resolveCompaction: () => {
+ // Lazily resolve the compaction service (registered below after
+ // the orchestrator). By the time this is called at runtime
+ // (after a turn settles), the service is registered.
+ try {
+ return host.getService(compactionHandle);
+ } catch {
+ return undefined;
+ }
+ },
+ resolveSystemPrompt: () => {
+ // Lazily resolve the system-prompt service. Returns undefined when
+ // the system-prompt extension isn't loaded (no system prompt sent —
+ // current behavior). Lazy so activation order with system-prompt
+ // doesn't matter; called per-turn / per-compaction, not at activate.
+ try {
+ return host.getService(systemPromptHandle);
+ } catch {
+ 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 —
+ // correct for vision-capable models; the feature degrades off cleanly for
+ // text-only turns). Lazy so activation order doesn't matter; the
+ // activated-manifests guard avoids a getService throw when absent.
+ const loaded = host.getExtensions().some((m) => m.id === "vision-handoff");
+ if (!loaded) return undefined;
+ try {
+ return host.getService(visionHandoffLocalHandle);
+ } catch {
+ return undefined;
+ }
+ },
+ });
- host.provideService(sessionOrchestratorHandle, orchestrator);
+ host.provideService(sessionOrchestratorHandle, orchestrator);
- // Remote-degradation rule (plan §6): when a turn is REMOTE
- // (`assembly.computerId !== undefined`), drop tools that spawn local
- // processes and cannot run over SFTP — the `lsp` tool (local LSP servers)
- // and MCP-namespaced tools (`<serverId>__<toolName>`, local MCP servers).
- // When LOCAL (`computerId === undefined`), the filter is a passthrough —
- // byte-identical to today. Registered at default priority (0) with
- // activation-order tie-breaking: session-orchestrator activates before
- // MCP (which dependsOn it), so this runs FIRST in the chain — the drops
- // happen before MCP's filter connects/registers servers. Mirrors how MCP
- // adds its own filter via host.addFilter.
- host.addFilter(toolsFilter, filterRemoteIncompatibleTools);
+ // Remote-degradation rule (plan §6): when a turn is REMOTE
+ // (`assembly.computerId !== undefined`), drop tools that spawn local
+ // processes and cannot run over SFTP — the `lsp` tool (local LSP servers)
+ // and MCP-namespaced tools (`<serverId>__<toolName>`, local MCP servers).
+ // When LOCAL (`computerId === undefined`), the filter is a passthrough —
+ // byte-identical to today. Registered at default priority (0) with
+ // activation-order tie-breaking: session-orchestrator activates before
+ // MCP (which dependsOn it), so this runs FIRST in the chain — the drops
+ // happen before MCP's filter connects/registers servers. Mirrors how MCP
+ // adds its own filter via host.addFilter.
+ host.addFilter(toolsFilter, filterRemoteIncompatibleTools);
- const warmService = createWarmService(
- {
- conversationStore,
- resolveProvider: () => selectFirstProvider(host.getProviders()),
- resolveTools: () => [...host.getTools().values()],
- resolveModel: (modelName: string) => {
- const store = host.getService(credentialStoreHandle);
- const r = store.resolve(modelName);
- if (r === undefined) return undefined;
- const provider = host.getProviders().get(r.providerId);
- return provider ? { provider, model: r.model } : undefined;
- },
- applyToolsFilter: (assembly) => host.applyFilters(toolsFilter, assembly),
- runTurn,
- logger: host.logger,
- now: () => Date.now(),
- emit: (hook, payload) => host.emit(hook, payload),
- },
- activeConversations,
- );
+ const warmService = createWarmService(
+ {
+ conversationStore,
+ resolveProvider: () => selectFirstProvider(host.getProviders()),
+ resolveTools: () => [...host.getTools().values()],
+ resolveModel: (modelName: string) => {
+ const store = host.getService(credentialStoreHandle);
+ const r = store.resolve(modelName);
+ if (r === undefined) return undefined;
+ const provider = host.getProviders().get(r.providerId);
+ return provider ? { provider, model: r.model } : undefined;
+ },
+ applyToolsFilter: (assembly) => host.applyFilters(toolsFilter, assembly),
+ runTurn,
+ logger: host.logger,
+ now: () => Date.now(),
+ emit: (hook, payload) => host.emit(hook, payload),
+ },
+ activeConversations,
+ );
- host.provideService(cacheWarmHandle, warmService);
+ host.provideService(cacheWarmHandle, warmService);
- const compactionService = createCompactionService(
- {
- conversationStore,
- resolveProvider: () => selectFirstProvider(host.getProviders()),
- resolveTools: () => [...host.getTools().values()],
- resolveModel: (modelName: string) => {
- const store = host.getService(credentialStoreHandle);
- const r = store.resolve(modelName);
- if (r === undefined) return undefined;
- const provider = host.getProviders().get(r.providerId);
- return provider ? { provider, model: r.model } : undefined;
- },
- resolveModelInfo: async (modelName: string) => {
- const store = host.getService(credentialStoreHandle);
- return store.getModelInfo(modelName);
- },
- resolveSystemPrompt: () => {
- try {
- return host.getService(systemPromptHandle);
- } catch {
- return undefined;
- }
- },
- applyToolsFilter: (assembly) => host.applyFilters(toolsFilter, assembly),
- runTurn,
- logger: host.logger,
- now: () => Date.now(),
- emit: (hook, payload) => host.emit(hook, payload),
- },
- activeConversations,
- );
+ const compactionService = createCompactionService(
+ {
+ conversationStore,
+ resolveProvider: () => selectFirstProvider(host.getProviders()),
+ resolveTools: () => [...host.getTools().values()],
+ resolveModel: (modelName: string) => {
+ const store = host.getService(credentialStoreHandle);
+ const r = store.resolve(modelName);
+ if (r === undefined) return undefined;
+ const provider = host.getProviders().get(r.providerId);
+ return provider ? { provider, model: r.model } : undefined;
+ },
+ resolveModelInfo: async (modelName: string) => {
+ const store = host.getService(credentialStoreHandle);
+ return store.getModelInfo(modelName);
+ },
+ resolveSystemPrompt: () => {
+ try {
+ return host.getService(systemPromptHandle);
+ } catch {
+ return undefined;
+ }
+ },
+ applyToolsFilter: (assembly) => host.applyFilters(toolsFilter, assembly),
+ runTurn,
+ logger: host.logger,
+ now: () => Date.now(),
+ emit: (hook, payload) => host.emit(hook, payload),
+ },
+ activeConversations,
+ );
- host.provideService(compactionHandle, compactionService);
+ host.provideService(compactionHandle, compactionService);
}
export const extension: Extension = {
- manifest,
- activate,
+ manifest,
+ activate,
};