| 1 | import { ReviewedLoader, installCompositionLoader } from './dsh/composition.ts' |
| 2 | /** |
| 3 | * The Cordis root that plugin fibers run under. |
| 4 | * |
| 5 | * The host has no turn loop, store, prompt authority or approval. Plugins |
| 6 | * reach the core only through the shim services here, and every one of them |
| 7 | * becomes a `registry/*` request that Rust admits or refuses. Service names |
| 8 | * that belong to the core cannot be provided by a plugin at all. |
| 9 | */ |
| 10 | import { createHash } from 'node:crypto' |
| 11 | import { readFile } from 'node:fs/promises' |
| 12 | import { AsyncLocalStorage } from 'node:async_hooks' |
| 13 | import { pathToFileURL } from 'node:url' |
| 14 | import { Context, Inject, Service } from '@deepseek-ai/cordis' |
| 15 | import { |
| 16 | ErrorCode, |
| 17 | type ActivateParams, |
| 18 | type ActivateResult, |
| 19 | type CommandResultWire, |
| 20 | type ContentBlockWire, |
| 21 | type DeactivateResult, |
| 22 | type HookEvaluateParams, |
| 23 | type HookVerdictWire, |
| 24 | type Json, |
| 25 | type OwnerRef, |
| 26 | type EntryRef, |
| 27 | type HarnessRunParams, |
| 28 | type McpOpenParams, |
| 29 | type McpRequestParams, |
| 30 | type McpCloseParams, |
| 31 | type ToolResultWire, |
| 32 | } from './protocol.ts' |
| 33 | import { RpcError, type RpcPeer } from './rpc.ts' |
| 34 | import { explainImportError } from './dsh/resolve-hooks.ts' |
| 35 | import { isJson } from './json.ts' |
| 36 | import { makeCoreApi } from './shims/core.ts' |
| 37 | import { OwnedRegistrations } from './shims/owned.ts' |
| 38 | import { defineShellHooksService,type LocalShellHook } from './shims/shell-hooks.ts' |
| 39 | import { hookExecution, hookVerdict, type LocalHook } from './shims/hooks.ts' |
| 40 | import { PromptSections, definePromptService, type LocalPromptSection } from './shims/prompt.ts' |
| 41 | import { createStorage, type PluginStorage } from './shims/storage.ts' |
| 42 | import {McpDefinitions,defineMcpService,type LocalMcp} from './shims/mcp.ts' |
| 43 | import { SkillRoots, defineSkillsService, type LocalSkillRoot } from './shims/skills.ts' |
| 44 | import { ownerTier, type HostTier } from './tier.ts' |
| 45 | import { |
| 46 | commandSpec, |
| 47 | defineCommandsService, |
| 48 | makeInvocation, |
| 49 | normalizeResult, |
| 50 | type LocalCommand, |
| 51 | type NormalizedCommand, |
| 52 | } from './shims/commands.ts' |
| 53 | |
| 54 | /** Context key carrying the owner record; inherited by every nested fiber. */ |
| 55 | export const OWNER = Symbol.for('codewhale.extension-host.owner') |
| 56 | |
| 57 | /** |
| 58 | * Service names a plugin may never provide: each is one authority the Rust |
| 59 | * core owns (§4.3 of the design). `tools`, `commands` and `logger` are |
| 60 | * provided by the host root as shims and are refused to plugins for the same |
| 61 | * reason. |
| 62 | */ |
| 63 | export const REFUSED_SERVICES = new Set(['shellHooks', |
| 64 | 'loader', // one host-owned composition loader; plugins may not replace it |
| 65 | 'approval', |
| 66 | 'agents', |
| 67 | 'sessions', |
| 68 | 'llm', |
| 69 | 'sandboxPolicy', |
| 70 | 'credentials', |
| 71 | 'fs', |
| 72 | 'subprocess', |
| 73 | 'systemPrompt', |
| 74 | 'runtimeLoop', // the one Rust Engine owns scheduling and execution |
| 75 | 'tools', |
| 76 | 'commands', |
| 77 | 'prompt', |
| 78 | 'storage', |
| 79 | 'skills', |
| 80 | 'mcp', |
| 81 | 'logger', |
| 82 | ]) |
| 83 | |
| 84 | /** Services the root provides; `inject` of anything else fails activation. */ |
| 85 | const PROVIDED_SERVICES = new Set(['tools', 'commands', 'prompt', 'storage', 'skills', 'mcp', 'logger', 'events', 'reflect', 'registry']) |
| 86 | |
| 87 | const ACTIVATE_DEADLINE_MS = 5_000 |
| 88 | const DISPOSE_DEADLINE_MS = 2_000 |
| 89 | |
| 90 | interface HarnessBuiltin { run(params: HarnessRunParams, signal: AbortSignal): Promise<unknown>; dispose(): Promise<void> } |
| 91 | |
| 92 | interface McpBuiltin { |
| 93 | open(params: McpOpenParams, signal: AbortSignal): Promise<unknown> |
| 94 | request(params: McpRequestParams, signal: AbortSignal): Promise<unknown> |
| 95 | close(owner: OwnerRef, sessionId: string): Promise<void> |
| 96 | dispose(): Promise<void> |
| 97 | } |
| 98 | |
| 99 | export interface OwnerRecord { |
| 100 | mcp?: McpBuiltin; |
| 101 | harness?: HarnessBuiltin; |
| 102 | ref: OwnerRef |
| 103 | scope?: EntryRef |
| 104 | views?: Map<string, OwnerRecord> |
| 105 | pluginName: string |
| 106 | fibers: any[] |
| 107 | /** In-flight `registry/register` requests, awaited before activation acks. */ |
| 108 | pendingRegistrations: Set<Promise<void>> |
| 109 | refusals: string[] |
| 110 | tools: Map<number, LocalTool> |
| 111 | commands: Map<number, LocalCommand<OwnerRecord>> |
| 112 | shellHooks:Map<number,LocalShellHook<OwnerRecord>> |
| 113 | hooks: Map<number, LocalHook<OwnerRecord>> |
| 114 | promptSections: Map<number, LocalPromptSection<OwnerRecord>> |
| 115 | skillRoots: Map<number, LocalSkillRoot<OwnerRecord>> |
| 116 | mcpDefinitions:Map<number,LocalMcp<OwnerRecord>> |
| 117 | storage?: PluginStorage |
| 118 | warnedAllow?: boolean |
| 119 | /** Entry modules activated under this owner so far (a plugin may declare several). */ |
| 120 | entries: Set<string> |
| 121 | /** The plugin's own writable directory, as the core named it at activation. */ |
| 122 | dataDir?: string |
| 123 | disposing?: Promise<void> |
| 124 | state: 'activating' | 'active' | 'failed' | 'disposed' |
| 125 | } |
| 126 | |
| 127 | interface LocalTool { |
| 128 | owner: OwnerRecord |
| 129 | name: string |
| 130 | handle?: number |
| 131 | definition: any |
| 132 | disposed: boolean |
| 133 | } |
| 134 | |
| 135 | export const ownerStorage = new AsyncLocalStorage<OwnerRecord>() |
| 136 | |
| 137 | function describeError(error: unknown): string { |
| 138 | if (error instanceof Error) return `${error.name}: ${error.message}` |
| 139 | return String(error) |
| 140 | } |
| 141 | |
| 142 | function withDeadline<T>(promise: Promise<T>, ms: number, label: string): Promise<T> { |
| 143 | let timer: NodeJS.Timeout |
| 144 | const deadline = new Promise<never>((_, reject) => { |
| 145 | timer = setTimeout(() => reject(new Error(`${label} exceeded ${ms} ms`)), ms) |
| 146 | timer.unref() |
| 147 | }) |
| 148 | return Promise.race([promise, deadline]).finally(() => clearTimeout(timer)) |
| 149 | } |
| 150 | |
| 151 | export class HostRoot { |
| 152 | readonly root: any |
| 153 | readonly owners = new Map<string, OwnerRecord>() |
| 154 | private readonly toolRegistrations: OwnedRegistrations<OwnerRecord, LocalTool> |
| 155 | private readonly commandRegistrations: OwnedRegistrations<OwnerRecord, LocalCommand<OwnerRecord>> |
| 156 | private readonly shellRegistrations:OwnedRegistrations<OwnerRecord,LocalShellHook<OwnerRecord>> |
| 157 | private readonly hookRegistrations: OwnedRegistrations<OwnerRecord, LocalHook<OwnerRecord>> |
| 158 | private readonly promptSections: PromptSections<OwnerRecord> |
| 159 | private readonly skillRoots: SkillRoots<OwnerRecord> |
| 160 | private readonly mcpDefinitions:McpDefinitions<OwnerRecord> |
| 161 | |
| 162 | constructor( |
| 163 | private readonly rpc: RpcPeer, |
| 164 | /** The trust tier this process serves; it activates only owners of that tier. */ |
| 165 | readonly tier: HostTier, |
| 166 | ) { |
| 167 | const root: any = new Context() |
| 168 | this.root = root |
| 169 | const host = this |
| 170 | this.toolRegistrations = new OwnedRegistrations( |
| 171 | rpc, |
| 172 | 'tool', |
| 173 | (owner) => owner.tools, |
| 174 | (message, owner) => this.log('warn', message, owner), |
| 175 | ) |
| 176 | this.commandRegistrations = new OwnedRegistrations( |
| 177 | rpc, |
| 178 | 'command', |
| 179 | (owner) => owner.commands, |
| 180 | (message, owner) => this.log('warn', message, owner), |
| 181 | ) |
| 182 | this.shellRegistrations=new OwnedRegistrations(rpc,'shell_hook',(owner)=>owner.shellHooks,(message,owner)=>this.log('warn',message,owner)) |
| 183 | this.hookRegistrations = new OwnedRegistrations(rpc, 'hook', (owner) => owner.hooks, |
| 184 | (message, owner) => this.log('warn', message, owner)) |
| 185 | this.promptSections = new PromptSections(rpc, (owner) => owner.promptSections, |
| 186 | (message, owner) => this.log('warn', message, owner)) |
| 187 | this.mcpDefinitions=new McpDefinitions(rpc,(owner)=>owner.mcpDefinitions,(message,owner)=>this.log('warn',message,owner)) |
| 188 | this.skillRoots = new SkillRoots(rpc, (owner) => owner.skillRoots, |
| 189 | (message, owner) => this.log('warn', message, owner)) |
| 190 | |
| 191 | // Cordis already owns listener effects and teardown. Intercept this one |
| 192 | // event using its supported extension point instead of replacing ctx.on. |
| 193 | // Actual upstream bridge seams for which this checkpoint has no core |
| 194 | // projection. Refuse instead of accepting listeners that never fire. |
| 195 | const unsupportedDshLifecycle = new Set([ |
| 196 | 'agent/created', 'agent/pre-step', 'agent/turn-stopping', |
| 197 | 'tools/post-execute', 'subagent/start', 'subagent/end', |
| 198 | ]) |
| 199 | root.on('internal/listener', function (this: any, name: string, callback: any, options: any) { |
| 200 | if (this[OWNER] && unsupportedDshLifecycle.has(name)) throw new Error(`DSH lifecycle ${name} is not projected by this host checkpoint`) |
| 201 | if (name !== 'tools/pre-execute') return |
| 202 | const owner: OwnerRecord | undefined = this[OWNER] |
| 203 | if (!owner) throw new Error('pre-execute listener registered outside an extension owner') |
| 204 | if (options?.prepend || options?.global) throw new Error('pre-execute listeners use core registration order; prepend/global are unsupported') |
| 205 | return this.effect(() => host.hookRegistrations.add( |
| 206 | { owner, name, callback, disposed: false }, |
| 207 | { name, description: 'Programmable tool admission listener' }, |
| 208 | ), 'ctx.on("tools/pre-execute")') |
| 209 | }) |
| 210 | |
| 211 | // Refuse core service names before any plugin can run. The refusal does |
| 212 | // not depend on who calls: `ctx.root.provide(...)` runs with the root as |
| 213 | // its context, so an owner check alone could be sidestepped. The one |
| 214 | // exception is each host shim of its own name, provided once, below. |
| 215 | const shimClasses = new Map<string, Function>() |
| 216 | const shimsProvided = new Set<string>() |
| 217 | const reflect = root.reflect |
| 218 | const originalProvide = reflect.provide |
| 219 | reflect.provide = function (this: any, name: string, value: unknown, ...rest: unknown[]) { |
| 220 | if (REFUSED_SERVICES.has(name)) { |
| 221 | const shim = shimClasses.get(name) |
| 222 | const hostShim = shim !== undefined && !shimsProvided.has(name) && value instanceof (shim as any) |
| 223 | if (!hostShim) { |
| 224 | throw new Error(`extension may not provide core service \`${name}\`: the Codewhale core owns it`) |
| 225 | } |
| 226 | shimsProvided.add(name) |
| 227 | } |
| 228 | return originalProvide.call(this, name, value, ...rest) |
| 229 | } |
| 230 | |
| 231 | // Logger shim: every Cordis log line becomes a `log` notification. |
| 232 | root.logger.exporter({ |
| 233 | colors: false, |
| 234 | export: (message: any) => { |
| 235 | const fiber = message.fiber?.deref?.() |
| 236 | const owner: OwnerRecord | undefined = fiber?.ctx?.[OWNER] |
| 237 | const text = (message.args ?? []).map((arg: unknown) => (arg instanceof Error ? describeError(arg) : typeof arg === 'string' ? arg : safeStringify(arg))).join(' ') |
| 238 | host.log(message.type === 'error' ? 'error' : message.type === 'warn' ? 'warn' : 'info', `[${message.name}] ${text}`, owner) |
| 239 | }, |
| 240 | }) |
| 241 | |
| 242 | class ToolsShim extends Service { |
| 243 | constructor(ctx: any) { |
| 244 | super(ctx, 'tools') |
| 245 | } |
| 246 | |
| 247 | /** DSH `ctx.tools.register(defineTool(...))`: returns an idempotent disposer. */ |
| 248 | register(definition: any) { |
| 249 | const ctx: any = this.ctx |
| 250 | const owner: OwnerRecord | undefined = ctx[OWNER] |
| 251 | if (!owner) throw new Error('tools.register called outside an extension owner') |
| 252 | validateDefinition(definition) |
| 253 | return ctx.effect(() => host.addTool(owner, definition), `tools.register(${JSON.stringify(definition.name)})`) |
| 254 | } |
| 255 | } |
| 256 | // Plugins share one process, so the owner token is not a boundary |
| 257 | // between them (design §4.4, threat 3). Freezing a shim at least stops |
| 258 | // the direct route of one plugin rewriting `register` for every other |
| 259 | // plugin; shared globals remain, and the approval card says so. |
| 260 | Object.freeze(ToolsShim.prototype) |
| 261 | const CommandsShim = defineCommandsService<OwnerRecord>({ |
| 262 | ownerOf: (ctx) => ctx[OWNER], |
| 263 | addCommand: (owner, command) => host.addCommand(owner, command), |
| 264 | }) |
| 265 | const PromptShim = definePromptService<OwnerRecord>({ |
| 266 | ownerOf: (ctx) => ctx[OWNER], |
| 267 | promptSections: this.promptSections, |
| 268 | }) |
| 269 | const ShellHooksShim=defineShellHooksService<OwnerRecord>({ownerOf:(ctx)=>ctx[OWNER],registrations:this.shellRegistrations}) |
| 270 | const McpShim=defineMcpService<OwnerRecord>({ownerOf:(ctx)=>(ctx as Context & { [OWNER]?: OwnerRecord })[OWNER],definitions:this.mcpDefinitions}) |
| 271 | const SkillsShim = defineSkillsService<OwnerRecord>({ ownerOf: (ctx) => ctx[OWNER], skillRoots: this.skillRoots }) |
| 272 | class StorageShim extends Service { |
| 273 | constructor(ctx: any) { super(ctx, 'storage') } |
| 274 | private api(): PluginStorage { |
| 275 | const owner: OwnerRecord | undefined = (this.ctx as any)[OWNER] |
| 276 | if (!owner || !owner.dataDir) throw new Error('storage requires an active extension owner and data directory') |
| 277 | owner.storage ??= createStorage({ |
| 278 | dataDir: owner.dataDir, |
| 279 | isActive: () => (owner.state === 'activating' || owner.state === 'active') && !owner.disposing, |
| 280 | }) |
| 281 | return owner.storage |
| 282 | } |
| 283 | get(key: string) { return this.api().get(key) } |
| 284 | set(key: string, value: Json) { return this.api().set(key, value) } |
| 285 | delete(key: string) { return this.api().delete(key) } |
| 286 | } |
| 287 | Object.freeze(StorageShim.prototype) |
| 288 | shimClasses.set('tools', ToolsShim) |
| 289 | shimClasses.set('commands', CommandsShim) |
| 290 | shimClasses.set('prompt', PromptShim) |
| 291 | shimClasses.set('storage', StorageShim) |
| 292 | shimClasses.set('skills', SkillsShim) |
| 293 | shimClasses.set('mcp',McpShim) |
| 294 | shimClasses.set('shellHooks',ShellHooksShim) |
| 295 | root.plugin(ToolsShim) |
| 296 | root.plugin(CommandsShim) |
| 297 | root.plugin(PromptShim) |
| 298 | root.plugin(StorageShim) |
| 299 | root.plugin(SkillsShim) |
| 300 | root.plugin(McpShim) |
| 301 | root.plugin(ShellHooksShim) |
| 302 | shimClasses.set('loader', ReviewedLoader) |
| 303 | installCompositionLoader(root) |
| 304 | } |
| 305 | |
| 306 | log(level: string, msg: string, owner?: OwnerRecord) { |
| 307 | const params: Record<string, string> = { level, msg: msg.slice(0, 8192) } |
| 308 | if (owner) params.plugin_id = owner.ref.plugin_id |
| 309 | this.rpc.notify('log', params) |
| 310 | } |
| 311 | |
| 312 | /** Called inside the owner's effect; returns the effect's cleanup. */ |
| 313 | private addTool(owner: OwnerRecord, definition: any): () => void { |
| 314 | const local: LocalTool = { owner, name: definition.name, definition, disposed: false } |
| 315 | return this.toolRegistrations.add(local, { |
| 316 | name: definition.name, |
| 317 | description: String(definition.description ?? ''), |
| 318 | input_schema: definition.parameters ?? { type: 'object', properties: {} }, |
| 319 | }) |
| 320 | } |
| 321 | |
| 322 | private addCommand(owner: OwnerRecord, definition: NormalizedCommand): () => void { |
| 323 | const local: LocalCommand<OwnerRecord> = { owner, name: definition.name, definition, disposed: false } |
| 324 | return this.commandRegistrations.add(local, commandSpec(definition)) |
| 325 | } |
| 326 | |
| 327 | /** |
| 328 | * `ext/activate`: load one `native` entry of a plugin under its owner. A |
| 329 | * manifest may declare several entries; the core sends one `ext/activate` |
| 330 | * per entry, in order, under the same owner token, and every entry becomes a |
| 331 | * fiber of that one owner. A further entry is accepted only after the |
| 332 | * previous one finished activating, for the same plugin, and only once per |
| 333 | * path. Core-selected Native scopes have independent fibers; failure withdraws |
| 334 | * that entry while siblings remain active. An unscoped owner stays atomic. |
| 335 | */ |
| 336 | async activate(params: ActivateParams): Promise<ActivateResult> { |
| 337 | // This process serves one tier. An owner of the other tier is refused |
| 338 | // outright, before anything of it is read or loaded. |
| 339 | const wanted = ownerTier(params.owner.plugin_id) |
| 340 | if (wanted !== this.tier) { |
| 341 | throw new RpcError( |
| 342 | ErrorCode.InvalidParams, |
| 343 | `owner ${JSON.stringify(params.owner.plugin_id)} belongs to the ${wanted} tier, but this host serves the ${this.tier} tier`, |
| 344 | ) |
| 345 | } |
| 346 | const key = params.owner.owner_token |
| 347 | let parent=this.owners.get(key) |
| 348 | const scopeKey=params.scope===undefined ? undefined : `${params.scope.path}\0${params.scope.sha256}` |
| 349 | if (scopeKey!==undefined) { |
| 350 | if (params.scope?.path!==params.entry.path || params.scope?.sha256!==params.entry.sha256) return {status:'failed',diagnostic:'scope does not match the core-selected entry'} |
| 351 | parent ??= {ref:params.owner,pluginName:params.plugin_name,fibers:[],pendingRegistrations:new Set(),refusals:[],tools:new Map(),commands:new Map(),hooks:new Map(),shellHooks:new Map(),promptSections:new Map(),skillRoots:new Map(),mcpDefinitions:new Map(),entries:new Set(),views:new Map(),state:'active',...(params.data_dir===undefined?{}:{dataDir:params.data_dir})} |
| 352 | if (parent.state!=='active' || parent.ref.plugin_id!==params.owner.plugin_id || parent.ref.generation!==params.owner.generation || parent.pluginName!==params.plugin_name) return {status:'failed',diagnostic:'scope owner was withdrawn'} |
| 353 | parent.views ??=new Map() |
| 354 | this.owners.set(key,parent) |
| 355 | } |
| 356 | const existing = scopeKey===undefined ? this.owners.get(key) : parent!.views!.get(scopeKey) |
| 357 | if (existing) { |
| 358 | if (existing.state !== 'active') return { status: 'failed', diagnostic: 'owner token already active' } |
| 359 | if (existing.ref.plugin_id !== params.owner.plugin_id || existing.pluginName !== params.plugin_name) { |
| 360 | return { status: 'failed', diagnostic: 'owner token already active for another plugin' } |
| 361 | } |
| 362 | if (existing.entries.has(params.entry.path)) { |
| 363 | return { status: 'failed', diagnostic: `entry ${params.entry.path} is already activated under this owner` } |
| 364 | } |
| 365 | existing.state = 'activating' |
| 366 | } |
| 367 | const owner: OwnerRecord = existing ?? { |
| 368 | ref: params.owner, |
| 369 | ...(params.scope===undefined ? {} : {scope:params.scope}), |
| 370 | pluginName: params.plugin_name, |
| 371 | fibers: [], |
| 372 | pendingRegistrations: new Set(), |
| 373 | refusals: [], |
| 374 | tools: new Map(), |
| 375 | commands: new Map(), |
| 376 | hooks: new Map(), |
| 377 | shellHooks:new Map(), |
| 378 | promptSections: new Map(), |
| 379 | skillRoots: new Map(), |
| 380 | mcpDefinitions:new Map(), |
| 381 | entries: new Set(), |
| 382 | ...(params.data_dir === undefined ? {} : { dataDir: params.data_dir }), |
| 383 | state: 'activating', |
| 384 | } |
| 385 | if (!existing) {if(scopeKey===undefined)this.owners.set(key,owner);else parent!.views!.set(scopeKey,owner)} |
| 386 | owner.entries.add(params.entry.path) |
| 387 | try { |
| 388 | const bytes = await readFile(params.entry.path) |
| 389 | const digest = createHash('sha256').update(bytes).digest('hex') |
| 390 | if (digest !== params.entry.sha256) { |
| 391 | throw new Error(`entry ${params.entry.path} changed after review (sha256 ${digest.slice(0, 12)}…)`) |
| 392 | } |
| 393 | const module = await ownerStorage.run(owner, () => import(pathToFileURL(params.entry.path).href)).catch((error) => { |
| 394 | throw explainImportError(error) |
| 395 | }) |
| 396 | if (this.tier === 'builtin' && owner.ref.plugin_id === 'host:harness') { |
| 397 | if (existing || typeof module.createHarnessModule !== 'function') throw new Error('invalid harness builtin module') |
| 398 | owner.harness = module.createHarnessModule(this.rpc, Object.freeze(structuredClone(owner.ref))) as HarnessBuiltin |
| 399 | owner.state = 'active' |
| 400 | return { status: 'ok', tools: [], commands: [] } |
| 401 | } |
| 402 | if (this.tier === 'builtin' && owner.ref.plugin_id === 'host:mcp') { |
| 403 | if (existing || typeof module.createMcpModule !== 'function') throw new Error('invalid MCP builtin module') |
| 404 | owner.mcp = module.createMcpModule(this.rpc, Object.freeze(structuredClone(owner.ref))) as McpBuiltin |
| 405 | owner.state = 'active' |
| 406 | return { status: 'ok', tools: [], commands: [] } |
| 407 | } |
| 408 | const plugin = pickPlugin(module) |
| 409 | const missing = Object.keys(Inject.resolve(plugin.inject)).filter((name) => !PROVIDED_SERVICES.has(name)) |
| 410 | if (missing.length > 0) { |
| 411 | throw new Error(`requires ${missing.map((name) => `\`${name}\``).join(', ')}, not provided by the Codewhale extension host in this phase`) |
| 412 | } |
| 413 | const ownerCtx = this.root.extend({ [OWNER]: owner }) |
| 414 | const fiber = ownerStorage.run(owner, () => ownerCtx.plugin(plugin, params.config ?? {})) |
| 415 | owner.fibers.push(fiber) |
| 416 | await withDeadline(Promise.resolve(fiber.await()), ACTIVATE_DEADLINE_MS, 'activation') |
| 417 | const pending = pendingFibers(fiber) |
| 418 | if (pending.length > 0) { |
| 419 | throw new Error(`plugin fibers are waiting on services the host does not provide: ${pending.join(', ')}`) |
| 420 | } |
| 421 | while (owner.pendingRegistrations.size > 0) { |
| 422 | await Promise.allSettled([...owner.pendingRegistrations]) |
| 423 | } |
| 424 | if (owner.refusals.length > 0) throw new Error(owner.refusals.join('; ')) |
| 425 | if (owner.state !== 'activating' || owner.disposing || (scopeKey !== undefined && parent?.state !== 'active')) throw new Error('entry was withdrawn while activating') |
| 426 | owner.state = 'active' |
| 427 | return { |
| 428 | status: 'ok', |
| 429 | tools: [...owner.tools.values()].map((tool) => tool.name).sort(), |
| 430 | commands: [...owner.commands.values()].map((command) => command.name).sort(), |
| 431 | } |
| 432 | } catch (error) { |
| 433 | owner.state = 'failed' |
| 434 | // Return the activation cause promptly, while retaining a dirty owner |
| 435 | // for Rust's bounded ext/deactivate receipt. Never hide an unfinished |
| 436 | // disposer by deleting the record before Core can observe it. |
| 437 | const cleanup = this.disposeOwner(owner) |
| 438 | void cleanup.then(() => { |
| 439 | this.toolRegistrations.forget(owner) |
| 440 | this.commandRegistrations.forget(owner) |
| 441 | this.hookRegistrations.forget(owner);this.shellRegistrations.forget(owner) |
| 442 | this.promptSections.forget(owner) |
| 443 | this.skillRoots.forget(owner);this.mcpDefinitions.forget(owner) |
| 444 | if (scopeKey === undefined) { |
| 445 | if (this.owners.get(key) === owner) this.owners.delete(key) |
| 446 | } else if (parent?.views?.get(scopeKey) === owner) parent.views.delete(scopeKey) |
| 447 | }, () => undefined) |
| 448 | return { status: 'failed', diagnostic: describeError(error) } |
| 449 | } |
| 450 | } |
| 451 | |
| 452 | async harnessRun(params: HarnessRunParams, signal: AbortSignal): Promise<unknown> { |
| 453 | const owner = this.owners.get(params.owner.owner_token) |
| 454 | if (this.tier !== 'builtin' || params.owner.plugin_id !== 'host:harness' || owner?.state !== 'active' || owner.ref.generation !== params.owner.generation || !owner.harness || owner.disposing) throw new RpcError(ErrorCode.NotAvailable, 'harness builtin owner is no longer live') |
| 455 | return owner.harness.run(params, signal) |
| 456 | } |
| 457 | |
| 458 | private mcpOwner(ref: OwnerRef): McpBuiltin { |
| 459 | const owner = this.owners.get(ref.owner_token) |
| 460 | if (this.tier !== 'builtin' || ref.plugin_id !== 'host:mcp' || owner?.state !== 'active' || owner.ref.generation !== ref.generation || !owner.mcp || owner.disposing) { |
| 461 | throw new RpcError(ErrorCode.NotAvailable, 'MCP builtin owner is no longer live') |
| 462 | } |
| 463 | return owner.mcp |
| 464 | } |
| 465 | async mcpOpen(params: McpOpenParams, signal: AbortSignal): Promise<unknown> { return this.mcpOwner(params.owner).open(params, signal) } |
| 466 | async mcpRequest(params: McpRequestParams, signal: AbortSignal): Promise<unknown> { return this.mcpOwner(params.owner).request(params, signal) } |
| 467 | async mcpClose(params: McpCloseParams): Promise<unknown> { await this.mcpOwner(params.owner).close(params.owner, params.session_id); return {} } |
| 468 | |
| 469 | /** Dispose one owner's fibers (reverse order, async disposers awaited). Memoised. */ |
| 470 | disposeOwner(owner: OwnerRecord): Promise<void> { |
| 471 | owner.state = 'disposed' |
| 472 | owner.disposing ??= (async () => { |
| 473 | for (const view of owner.views?.values() ?? []) await this.disposeOwner(view) |
| 474 | await owner.mcp?.dispose() |
| 475 | await owner.harness?.dispose() |
| 476 | for (const fiber of [...owner.fibers].reverse()) { |
| 477 | await fiber.dispose() |
| 478 | } |
| 479 | // Admission may have answered after a failed fiber's disposer ran. |
| 480 | // OwnedRegistrations compensates that admission and tracks its cleanup; |
| 481 | // acknowledge withdrawal only once both phases have settled. |
| 482 | while (owner.pendingRegistrations.size > 0) { |
| 483 | await Promise.allSettled([...owner.pendingRegistrations]) |
| 484 | } |
| 485 | owner.state = 'disposed' |
| 486 | })() |
| 487 | return owner.disposing |
| 488 | } |
| 489 | |
| 490 | async deactivate(ref: OwnerRef, entry?: EntryRef): Promise<DeactivateResult> { |
| 491 | const parent=this.owners.get(ref.owner_token) |
| 492 | if (parent && (parent.ref.plugin_id!==ref.plugin_id || parent.ref.generation!==ref.generation)) throw new Error('deactivation owner differs') |
| 493 | if (entry===undefined && parent?.views) { |
| 494 | parent.state = 'disposed' |
| 495 | const results=await Promise.all([...parent.views.values()].map(view=>this.deactivate(ref,view.scope))) |
| 496 | parent.views.clear() |
| 497 | this.owners.delete(ref.owner_token) |
| 498 | return {disposed:results.every(result=>result.disposed),leaked:results.flatMap(result=>result.leaked)} |
| 499 | } |
| 500 | const scopeKey=entry===undefined?undefined:`${entry.path}\0${entry.sha256}` |
| 501 | const owner = scopeKey===undefined?parent:parent?.views?.get(scopeKey) |
| 502 | if (!owner) return { disposed: true, leaked: [] } |
| 503 | let disposed = true |
| 504 | try { |
| 505 | await withDeadline(this.disposeOwner(owner), DISPOSE_DEADLINE_MS, 'dispose') |
| 506 | } catch { |
| 507 | disposed = false |
| 508 | } |
| 509 | const leaked = [ |
| 510 | ...[...owner.tools.values()].map((tool) => `tool:${tool.name}`), |
| 511 | ...[...owner.commands.values()].map((command) => `command:${command.name}`), |
| 512 | ...[...owner.shellHooks.values()].map((hook)=>`shell_hook:${hook.name}`), |
| 513 | ...[...owner.hooks.values()].map((hook) => `hook:${hook.name}`), |
| 514 | ...[...owner.promptSections.values()].map((section) => `prompt_section:${section.name}`), |
| 515 | ...[...owner.mcpDefinitions.values()].map(server=>`mcp_server:${server.name}`), |
| 516 | ...[...owner.skillRoots.values()].map((root) => `skill_root:${root.name}`), |
| 517 | ] |
| 518 | for (const fiber of owner.fibers) { |
| 519 | for (const effect of fiber.getEffects?.() ?? []) leaked.push(`effect:${effect.label}`) |
| 520 | } |
| 521 | this.toolRegistrations.forget(owner) |
| 522 | this.commandRegistrations.forget(owner) |
| 523 | this.hookRegistrations.forget(owner);this.shellRegistrations.forget(owner) |
| 524 | this.promptSections.forget(owner) |
| 525 | this.skillRoots.forget(owner);this.mcpDefinitions.forget(owner) |
| 526 | if(scopeKey===undefined)this.owners.delete(ref.owner_token);else parent?.views?.delete(scopeKey) |
| 527 | return { disposed, leaked } |
| 528 | } |
| 529 | |
| 530 | async deactivateAll(deadlineMs: number) { |
| 531 | await withDeadline( |
| 532 | Promise.allSettled([...this.owners.values()].map((owner) => this.deactivate(owner.ref))), |
| 533 | deadlineMs, |
| 534 | 'shutdown', |
| 535 | ).catch(() => undefined) |
| 536 | } |
| 537 | |
| 538 | async callTool( |
| 539 | handle: number, |
| 540 | input: unknown, |
| 541 | callId: string, |
| 542 | signal: AbortSignal, |
| 543 | workspace?: string, |
| 544 | /** The core's invocation ticket: present only when this call runs under the turn loop's permission gate. */ |
| 545 | ticket?: string, |
| 546 | identity?: InvocationIdentityWire, |
| 547 | ): Promise<ToolResultWire> { |
| 548 | const local = this.toolRegistrations.byHandle.get(handle) |
| 549 | if (!local || local.disposed || local.owner.state !== 'active') { |
| 550 | throw new RpcError(ErrorCode.NotAvailable, `tool handle ${handle} is not live`) |
| 551 | } |
| 552 | const definition = local.definition |
| 553 | // `workspace` is the calling session's workspace, per call; `dataDir` is |
| 554 | // this plugin's own directory. Both are read-only strings. |
| 555 | // `core` exists only when the core gave this call a ticket (`exec.core` in shims/core.ts). |
| 556 | const core = ticket === undefined ? {} : { core: makeCoreApi(this.rpc, local.owner.ref, ticket, signal) } |
| 557 | const exec = Object.freeze({ signal, callId, args: input, ...callContext(local.owner, workspace, identity), ...core }) |
| 558 | const run = ownerStorage.run(local.owner, async () => { |
| 559 | const value = await definition.execute(input, exec) |
| 560 | return renderResult(definition, input, value) |
| 561 | }) |
| 562 | return Promise.race([run, abortedBy(signal)]) |
| 563 | } |
| 564 | |
| 565 | /** `command/run`: only for a handle the core admitted, and only on a user's own invocation. */ |
| 566 | async callCommand(handle: number, rawInput: string, commandId: string, signal: AbortSignal, workspace?: string, identity?: InvocationIdentityWire): Promise<CommandResultWire> { |
| 567 | const local = this.commandRegistrations.byHandle.get(handle) |
| 568 | if (!local || local.disposed || local.owner.state !== 'active') { |
| 569 | throw new RpcError(ErrorCode.NotAvailable, `command handle ${handle} is not live`) |
| 570 | } |
| 571 | const { definition } = local |
| 572 | const run = ownerStorage.run(local.owner, async () => { |
| 573 | const value = await definition.handler(makeInvocation(rawInput, commandId, signal, callContext(local.owner, workspace, identity))) |
| 574 | return normalizeResult(definition.name, value) |
| 575 | }) |
| 576 | return Promise.race([run, abortedBy(signal)]) |
| 577 | } |
| 578 | |
| 579 | async evaluateHook(params: HookEvaluateParams, signal: AbortSignal): Promise<HookVerdictWire> { |
| 580 | const local = this.hookRegistrations.byHandle.get(params.handle) |
| 581 | if (!local || local.disposed || local.owner.state !== 'active' || params.event !== local.name) { |
| 582 | throw new RpcError(ErrorCode.NotAvailable, 'hook handle is not live for this event') |
| 583 | } |
| 584 | const exec = hookExecution(params.payload, signal) |
| 585 | const run = ownerStorage.run(local.owner, async () => hookVerdict( |
| 586 | await local.callback(exec, async () => ({ kind: 'abstain' })), |
| 587 | () => { |
| 588 | if (local.owner.warnedAllow) return |
| 589 | local.owner.warnedAllow = true |
| 590 | this.log('warn', 'pre-execute allow is an abstention; Rust still evaluates all admission gates', local.owner) |
| 591 | }, |
| 592 | )) |
| 593 | return Promise.race([run, abortedBy(signal)]) |
| 594 | } |
| 595 | } |
| 596 | |
| 597 | /** The read-only context a call carries beyond its input; a field the core did not send is absent. */ |
| 598 | interface InvocationIdentityWire { |
| 599 | session_id?: string |
| 600 | agent_id?: string |
| 601 | origin_turn_id?: string |
| 602 | } |
| 603 | |
| 604 | function callContext(owner: OwnerRecord, workspace: string | undefined, identity?: InvocationIdentityWire) { |
| 605 | return { |
| 606 | ...(workspace === undefined ? {} : { workspace }), |
| 607 | ...(owner.dataDir === undefined ? {} : { dataDir: owner.dataDir }), |
| 608 | ...(identity?.session_id === undefined ? {} : { sessionId: identity.session_id }), |
| 609 | ...(identity?.agent_id === undefined ? {} : { agentId: identity.agent_id }), |
| 610 | ...(identity?.origin_turn_id === undefined ? {} : { originTurnId: identity.origin_turn_id }), |
| 611 | } |
| 612 | } |
| 613 | |
| 614 | /** Rejects as cancelled when `signal` aborts. */ |
| 615 | function abortedBy(signal: AbortSignal): Promise<never> { |
| 616 | return new Promise<never>((_, reject) => { |
| 617 | const onAbort = () => reject(new RpcError(ErrorCode.Cancelled, 'cancelled')) |
| 618 | if (signal.aborted) onAbort() |
| 619 | else signal.addEventListener('abort', onAbort, { once: true }) |
| 620 | }) |
| 621 | } |
| 622 | |
| 623 | function renderResult(definition: any, input: unknown, value: unknown): ToolResultWire { |
| 624 | let blocks: unknown = undefined |
| 625 | if (typeof definition.output?.render === 'function') { |
| 626 | blocks = definition.output.render(input, value) |
| 627 | } |
| 628 | const content: ContentBlockWire[] = [] |
| 629 | if (Array.isArray(blocks)) { |
| 630 | for (const block of blocks) { |
| 631 | if (block && typeof block === 'object' && (block as any).type === 'text' && typeof (block as any).text === 'string') { |
| 632 | content.push({ type: 'text', text: (block as any).text }) |
| 633 | } |
| 634 | } |
| 635 | } else { |
| 636 | content.push({ type: 'text', text: typeof value === 'string' ? value : safeStringify(value) }) |
| 637 | } |
| 638 | const result: ToolResultWire = { content, is_error: false } |
| 639 | if (value !== undefined && isJson(value)) result.structured = value |
| 640 | return result |
| 641 | } |
| 642 | |
| 643 | function safeStringify(value: unknown): string { |
| 644 | try { |
| 645 | return JSON.stringify(value) ?? String(value) |
| 646 | } catch { |
| 647 | return String(value) |
| 648 | } |
| 649 | } |
| 650 | |
| 651 | function validateDefinition(definition: any) { |
| 652 | if (!definition || typeof definition !== 'object') throw new TypeError('tool definition must be an object') |
| 653 | if (typeof definition.name !== 'string' || definition.name.length === 0) throw new TypeError('tool definition needs a name') |
| 654 | if (typeof definition.execute !== 'function') throw new TypeError(`tool \`${definition.name}\` needs an execute function`) |
| 655 | if (definition.parameters !== undefined && (typeof definition.parameters !== 'object' || definition.parameters === null)) { |
| 656 | throw new TypeError(`tool \`${definition.name}\` parameters must be a JSON schema object`) |
| 657 | } |
| 658 | } |
| 659 | |
| 660 | function pickPlugin(module: any): any { |
| 661 | const candidate = module?.default ?? module |
| 662 | if (typeof candidate === 'function') return candidate |
| 663 | if (candidate && typeof candidate.apply === 'function') return candidate |
| 664 | if (module && typeof module.apply === 'function') return module |
| 665 | throw new Error('entry module exports no Cordis plugin (a function, or an object with `apply`)') |
| 666 | } |
| 667 | |
| 668 | /** Names of services that nested fibers under `fiber` are still waiting on. */ |
| 669 | function pendingFibers(fiber: any): string[] { |
| 670 | const out: string[] = [] |
| 671 | const seen = new Set<any>() |
| 672 | const visit = (current: any) => { |
| 673 | if (!current || seen.has(current)) return |
| 674 | seen.add(current) |
| 675 | if (current.state === 0 /* PENDING */) { |
| 676 | const names = Object.keys(current.inject ?? {}) |
| 677 | out.push(`${current.name ?? 'plugin'} (${names.join(', ')})`) |
| 678 | } |
| 679 | } |
| 680 | visit(fiber) |
| 681 | for (const runtime of fiber.ctx?.registry?.values?.() ?? []) { |
| 682 | for (const child of runtime.fibers ?? []) { |
| 683 | if (child.parent?.[OWNER] === fiber.ctx?.[OWNER]) visit(child) |
| 684 | } |
| 685 | } |
| 686 | return out |
| 687 | } |
| 688 |