diff options
Diffstat (limited to 'packages/session-orchestrator/src/extension.ts')
| -rw-r--r-- | packages/session-orchestrator/src/extension.ts | 336 |
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, }; |
