返回 CodeWhale
root.ts
根目录 / crates / tui / extension-host / src / root.ts
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
688 lines TYPESCRIPT