返回源码地图

packages/core/agent/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 02 / 07 / 07;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Agent service: live registry, factory delegation, and process-local
3 * initiator scope. Concrete creation and driving belong to the loop.
4 *
5 * @module @deepseek-ai/dsh-agent
6 */
7
8import { Context, FiberState, getTraceable, Service, symbols } from '@deepseek-ai/cordis'
9import type { Fiber } from '@deepseek-ai/cordis'
10import { AsyncLocalStorage } from 'node:async_hooks'
11import { isPromise } from 'node:util/types'
12import { scopeTarget } from '@deepseek-ai/dsh-scope'
13import type { Scoped } from '@deepseek-ai/dsh-scope'
14import type { SessionEvent, SessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'
15import { installTurnArchiveAdmission } from './archive-admission.ts'
16import type { Agent } from './types.ts'
17import type { AgentOptions, SessionStartSource } from './runtime-types.ts'
18
19export * from './runtime-types.ts'
20export * from './types.ts'
21export type * from './projection.ts'
22export * from './consumed-work.ts'
23export * from './model-selection.ts'
24export { agentCarrier, agentEvents, assembleContextFor, emitAgentEvent } from './dispatch.ts'
25export type { AgentEventDispatch, AgentSubjectEvent } from './dispatch.ts'
26
27declare module '@deepseek-ai/cordis' {
28 interface Context {
29 agents: AgentRegistry
30 }
31}
32
33/**
34 * Synchronous finalizer returned by unpublished Agent setup when its
35 * contributions need validation at the exact publication commit point.
36 */
37export interface AgentSetupCommit {
38 /**
39 * Validate and commit the prepared setup immediately before publication.
40 * @throws when publication must roll the unpublished Agent back.
41 */
42 commit(): void
43}
44
45/**
46 * Compose an unpublished Agent scope and optionally return its publication commit.
47 * @param agentCtx - unpublished Agent scope.
48 * @param agent - unpublished Agent being composed.
49 * @returns an optional synchronous commit invoked after setup awaits settle and immediately before publication.
50 */
51export type AgentSetup = (
52 agentCtx: Context,
53 agent: Agent,
54) => AgentSetupCommit | Promise<AgentSetupCommit | void> | void
55
56/**
57 * Options for programmatically creating an agent through the registry factory
58 * ({@link AgentRegistry.create}). The caller supplies the single live
59 * `sessionId` shared by the agent registry and session log (e.g. an
60 * ACP-generated id), plus optional session metadata (the validated `cwd`, fork
61 * lineage); the factory creates the session and agent under that identity.
62 */
63export interface CreateAgentOptions {
64 /** The live agent/session identity. */
65 readonly sessionId: SessionId
66 /** Live parent Agent for runtime ownership; omit for a root Agent. */
67 readonly parentAgent?: Agent
68 /**
69 * Session creation metadata: validated absolute `cwd`, `parentSession`
70 * fork lineage, the `isSeeded` fork marker, the coarse `origin`
71 * classification, and the `delegationDepth` recursion budget. Mirrors the
72 * `cwd`/`parentSession`/`isSeeded`/`origin`/`delegationDepth` fields of
73 * {@link CreateSessionOptions.meta} in dsh-session (the internal-only
74 * `createdAt`, used when reconstructing a persisted session, is deliberately
75 * excluded — a factory caller never sets it). This is durable session data,
76 * so the session boundary validates and snapshots it before asynchronous
77 * setup begins.
78 */
79 readonly meta?: {
80 readonly cwd?: string
81 readonly parentSession?: SessionId
82 readonly isSeeded?: boolean
83 readonly origin?: 'subagent'
84 readonly delegationDepth?: number
85 readonly agentPreset?: string
86 }
87 /** Exact fork-inherited prefix length when the session metadata sets `isSeeded`. */
88 readonly inheritedEventCount?: SessionLogOffset
89 /**
90 * Initial replay/fork history, contiguous from seq 0 with lossless-JSON data.
91 * A fork supplies an exact parent prefix, its inherited marker, and closers
92 * for the open tail. Previously closed steps and turns remain unchanged.
93 * The factory validates and snapshots the seed before publication.
94 */
95 readonly seed?: readonly SessionEvent[]
96 /** Per-agent options (model, …). */
97 readonly agentOptions?: AgentOptions
98 /** Optional creation-only cancellation signal; detached before the returned handle becomes visible. */
99 readonly signal?: AbortSignal
100 /**
101 * Creation-time composition of the agent's scoped world. The factory awaits
102 * setup after minting `agentCtx` but BEFORE inserting or announcing either
103 * the session or agent, so observers can never see a partially configured
104 * world. Setup may return an {@link AgentSetupCommit}; the factory invokes its
105 * synchronous `commit()` after every setup await settles and immediately
106 * before registry publication. This lets mutable provisioning revalidate at
107 * the exact publication boundary. Everything registered through `agentCtx`
108 * (scoped tools, prompt sections/variables, `restrict()`, listeners, awaited
109 * child plugins) exists before `session/created`, `agent/created`,
110 * and the first prompt assembly. A setup
111 * throw/rejection, commit throw, or owner disposal rolls the scope back
112 * without publishing either id.
113 *
114 * **Setup composes, it never drives**: the callback is trusted same-process
115 * code and receives the full scoped context, so this is a contract rather
116 * than a runtime restriction. Drive the agent only after creation resolves.
117 */
118 readonly setup?: AgentSetup
119}
120
121/**
122 * Options for resuming an agent on a persisted session
123 * ({@link AgentRegistry.resume}).
124 */
125export interface ResumeAgentOptions {
126 /** The persisted session id to load and use as the live agent/session identity. */
127 readonly resumeSessionId: SessionId
128 /** Live parent Agent for runtime ownership; omit for a root Agent. */
129 readonly parentAgent?: Agent
130 /** Per-agent options (model, …). */
131 readonly agentOptions?: AgentOptions
132 /** Optional creation-only cancellation signal for persistence load/setup; detached before return. */
133 readonly signal?: AbortSignal
134 /**
135 * Resume-time composition of the agent's fresh scoped world. Persistence is
136 * loaded first; the factory then mints `agentCtx` and awaits setup while the
137 * reconstructed session and agent remain unpublished. The callback has the
138 * same trusted composition-only contract and optional synchronous
139 * publication commit as {@link CreateAgentOptions.setup}: all registrations
140 * exist before either creation announcement, and rejection, commit failure,
141 * or owner disposal rolls the transaction back without publishing either id.
142 */
143 readonly setup?: AgentSetup
144}
145
146/**
147 * An owned agent plus its disposer, returned by {@link AgentRegistry.create} /
148 * {@link AgentRegistry.resume}. The disposer is a CAPABILITY: among consumers,
149 * only the holder can tear this agent down. The registered factory provider is
150 * also a structural owner because the scoped agent depends on that provider's
151 * service API; provider unload stops and drains every live handle it made.
152 * `dispose()` stops the loop, awaits its exit, unregisters the agent, removes
153 * its session from the store, and finally unwinds its scoped world.
154 *
155 * `ctx.agents.get(id)` still returns a bare {@link Agent} — the handle is
156 * exposed only to the consumer owner that created it; the structural provider
157 * reaches the same teardown internally. Config-created agents (the loop's own
158 * startup) are owned by the loop fiber and never need a handle.
159 */
160export interface AgentHandle {
161 agent: Agent
162 dispose(): Promise<void>
163}
164
165/**
166 * The agent-creation factory the loop implementation provides to the registry
167 * via {@link AgentRegistry.setFactory}. Kept on the `dsh-agent` interface so
168 * consumers (e.g. the ACP bridge) program against `ctx.agents` without
169 * depending on the concrete `dsh-agent-loop` package.
170 */
171export interface AgentFactory {
172 /**
173 * Create a new agent on a caller-supplied session id. Async because creation
174 * awaits unpublished setup, invokes its optional synchronous commit, inserts
175 * both session and agent, announces session creation, and awaits serial
176 * `agent/created` listeners before releasing queued work. The sequence is
177 * rollback-covered, but notifications delivered before a later listener
178 * failure remain observable; every agent or session creation announcement
179 * that began is paired by `agent/disposed` or `session/disposed` during
180 * rollback. The owner disposes the resolved handle to stop/drain,
181 * unregister, remove the session, and unwind the scope.
182 * The registry passes a context carrying the `create()` caller's fiber and
183 * scope as `ownerCtx`. The implementation attaches the unpublished
184 * transaction and resulting lifecycle to that owner; it must not infer
185 * ownership from the factory object's registration context.
186 * @param ownerCtx - caller-bound context that owns the transaction and live handle.
187 * @param options - agent/session identity, configuration, optional live parent, and setup.
188 * @returns the owned handle after setup and both creation announcements complete.
189 */
190 createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle>
191 /**
192 * Resume an agent on a persisted session. Async because it opens the
193 * persisted session for write, reads and repairs the log, publishes it, and
194 * awaits the optional unpublished setup transaction; must be called after
195 * `ctx.sessionPersistence` exists (consumers inject `sessionPersistence`).
196 * Publication follows the same setup-commit and ordered boundary as
197 * {@link createAgent}.
198 * @param ownerCtx - caller-bound context that owns load, setup, and the live handle.
199 * @param options - persisted identity, configuration, optional live parent, and setup.
200 * @returns the owned handle after setup, both announcements, and loop start complete.
201 */
202 resume(ownerCtx: Context, options: ResumeAgentOptions): Promise<AgentHandle>
203}
204
205/** Thrown when create/resume is called before an agent factory is registered. */
206const NO_FACTORY_MESSAGE = 'no agent factory registered (load an agent-loop plugin)'
207const NO_INITIATOR_MESSAGE = 'no initiating agent is active'
208const DISPOSED_INITIATOR_MESSAGE = 'agent initiator scope is disposed'
209
210/** All mutable lifecycle state for one exact registry entry. */
211interface AgentEntry {
212 readonly id: SessionId
213 readonly agent: Agent
214 /** Runtime creator-agent ownership; independent of durable session lineage. */
215 readonly owner: Agent | undefined
216 readonly carrier: Scoped<Agent>
217 announced: boolean
218 announcing: boolean
219 detachRequested: boolean
220}
221
222/** One tracked boundary plus its inherited nesting chain. */
223interface InitiatorRun {
224 active: boolean
225 readonly parent: InitiatorRun | undefined
226}
227
228/** Plain holder prevents Cordis from tracing the factory field before the caller context is known. */
229interface FactorySlot {
230 readonly target: AgentFactory
231}
232
233/**
234 * Agent service (`ctx.agents`): tracks live agents and carries the initiating
235 * Agent through one process-local asynchronous driver chain. Agent *creation*
236 * is provided by whichever plugin implements the {@link AgentFactory}
237 * (`@deepseek-ai/dsh-agent-loop`), registered via {@link setFactory}.
238 *
239 * Initiator methods provide same-process causal attribution only. Ambient
240 * presence is neither liveness proof nor authorization; subjects and owners
241 * remain explicit, as does identity at worker, process, persistence, and wire
242 * boundaries. Returned Promise boundaries drain during teardown, except a
243 * nested lineage that starts an owning-fiber unload is excluded from its own drain.
244 */
245export class AgentRegistry extends Service {
246 private store = new Map<SessionId, AgentEntry>()
247 private factory: FactorySlot | undefined
248 private readonly initiators = new AsyncLocalStorage<Agent | undefined>()
249 private readonly initiatorRuns = new AsyncLocalStorage<InitiatorRun>()
250 private initiatorState: 'active' | 'closing' | 'disposed' = 'active'
251 private activeInitiatorRuns = 0
252 private initiatorDrain: PromiseWithResolvers<void> | undefined
253 private initiatorDisposal: Promise<void> | undefined
254
255 constructor(ctx: Context) {
256 super(ctx, 'agents')
257 ctx.inject(['typert'], (typeCtx) => {
258 typeCtx.typert.lookups.register('agent', {
259 parameter: 'agent',
260 wire: 'agentId',
261 hostTypeSymbol: '@deepseek-ai/dsh-agent#Agent',
262 wireTypeSymbol: '@deepseek-ai/dsh-session/types#SessionId',
263 resolve: sessionId => this.get(sessionId),
264 })
265 typeCtx.typert.contexts.registerHost('agent', {
266 wire: 'agentId',
267 wireTypeSymbol: '@deepseek-ai/dsh-session/types#SessionId',
268 resolve: sessionId => this.get(sessionId)?.ctx,
269 })
270 })
271 ctx.on('internal/status', (fiber) => {
272 if (fiber.state === FiberState.UNLOADING && this.hasLifecycleAncestor(fiber)) {
273 this.closeInitiators()
274 }
275 })
276 ctx.effect(function* (this: AgentRegistry) {
277 yield () => this.disposeInitiators()
278 yield () => { this.closeInitiators() }
279 }.bind(this), 'agents.initiatorLifecycle()')
280 // Archive admission: the Workspace registry asks what still runs for a
281 // Session before hiding it; a running turn answers here, for every Agent.
282 installTurnArchiveAdmission(ctx, sessionId => this.get(sessionId))
283 }
284
285 /**
286 * Read the Agent that initiated the inherited asynchronous driver chain.
287 * Use this optional form for logging, tracing, metrics, or host attribution
288 * that also supports agentless calls. When a parent creates a child, setup
289 * reports the causal parent while the setup callback's Agent parameter
290 * identifies the child.
291 * @returns the inherited Agent, or `undefined` outside an initiator boundary
292 * and inside an explicit clearing boundary.
293 * @throws when this service instance has been disposed.
294 */
295 currentInitiator(): Agent | undefined {
296 this.assertInitiatorsReadable()
297 return this.initiators.getStore()
298 }
299
300 /**
301 * Read the initiating Agent and fail when no initiator boundary is active.
302 * Use this for private helpers contractually below a driver, or for a
303 * deployment-owned outbound request whose contract forbids agentless calls.
304 * Generic or direct-call paths use optional lookup or explicit request fields.
305 * @returns the inherited Agent.
306 * @throws when no initiator is active or this service instance has been disposed.
307 */
308 requireInitiator(): Agent {
309 const agent = this.currentInitiator()
310 if (agent === undefined) throw new Error(NO_INITIATOR_MESSAGE)
311 return agent
312 }
313
314 /**
315 * Run an operation with one exact Agent as its process-local initiator. The
316 * exact synchronous value or Promise returned by the operation is preserved.
317 * Custom drivers and test harnesses wrap their complete returned foreground
318 * lifetime.
319 * A queue or wire receiver may establish this boundary only after validating
320 * explicit identity and resolving the exact live Agent; this method does neither.
321 * Detached work remains owned by the subsystem that starts it.
322 * @param agent - initiating Agent to inherit; presence is neither liveness proof nor authorization.
323 * @param operation - synchronous or asynchronous operation to invoke.
324 * @returns the exact value returned by `operation`.
325 * @throws when the initiator scope is closing/disposed, or when `operation` throws.
326 */
327 withInitiator<T>(agent: Agent, operation: () => T): T {
328 return this.runWithInitiator(agent, operation)
329 }
330
331 /**
332 * Run an operation inside a boundary that hides any inherited initiating
333 * Agent. The exact synchronous value or Promise is preserved.
334 * Use this while creating lazy shared timers, queue pumps, pool maintenance,
335 * watchers, or exporters so they do not inherit the first Agent that happens
336 * to initialize them. It clears only initiator attribution, not explicit
337 * fields, and does not own or drain detached resources.
338 * @param operation - synchronous or asynchronous operation to invoke without an initiator.
339 * @returns the exact value returned by `operation`.
340 * @throws when the initiator scope is closing/disposed, or when `operation` throws.
341 */
342 withoutInitiator<T>(operation: () => T): T {
343 return this.runWithInitiator(undefined, operation)
344 }
345
346 /**
347 * Register the agent-creation factory (the loop calls this on construction,
348 * effect-scoped). A traced Cordis service is canonicalized to its concrete
349 * target; each create/resume call is then traced through that caller's
350 * context so ownership follows the caller without stacking proxy layers.
351 * Throws if a factory is already registered. Returns the disposer; on
352 * dispose the factory slot is cleared.
353 * @param factory - the loop-owned factory {@link create}/{@link resume} delegate to.
354 * @returns the disposer that clears the factory slot. The exact
355 * Cordis effect disposer (single-shot): composite (generator) effects may
356 * yield it directly — exact identity nests the teardown in order.
357 */
358 setFactory(factory: AgentFactory): () => void {
359 const dispose = this.ctx.effect(() => {
360 if (this.factory !== undefined) throw new Error('an agent factory is already registered')
361 // Avoid stacking two Cordis shadow layers when a caller passes a Service
362 // already read through a context. Calls are re-traced through their
363 // actual owner context below.
364 const target = (factory as AgentFactory & { [symbols.original]?: AgentFactory })[symbols.original] ?? factory
365 this.factory = { target }
366 return () => { this.factory = undefined }
367 }, 'agents.setFactory()')
368 // The exact cordis effect disposer (the agents.register() convention): a
369 // caller's composite effect can yield it for in-order teardown; the
370 // loop's constructor effect returns it directly, identity-nesting the
371 // registration under that effect.
372 // oxlint-disable-next-line typescript/no-misused-promises -- synchronous cleanup; direct return preserves disposer identity
373 return dispose
374 }
375
376 /** Return the active creation factory. */
377 private requireFactory(): FactorySlot {
378 if (this.factory === undefined) throw new Error(NO_FACTORY_MESSAGE)
379 return this.factory
380 }
381
382 /**
383 * Create and publish a new agent through the registered factory.
384 * Distinct from {@link register} (which records an already-constructed
385 * agent): this constructs the agent and its session. Rejects if no factory is
386 * registered or creation/setup fails. The resolved {@link AgentHandle} lets
387 * the owner tear down exactly this agent.
388 * @param options - shared identity, optional live parent, session seed/metadata, and agent options.
389 * @returns the handle after setup, rollback-covered publication, and loop start complete.
390 */
391 async create(options: CreateAgentOptions): Promise<AgentHandle> {
392 const ownerCtx = this.ctx
393 // Re-trace a Service-backed factory through the accessing context
394 // explicitly. This preserves AgentLoop's dependency origin while binding
395 // its effects to ownerCtx; plain factories receive ownerCtx as an explicit
396 // capability and need no Cordis tracker magic.
397 const { target } = this.requireFactory()
398 const receiver = getTraceable(ownerCtx, target)
399 // oxlint-disable-next-line typescript/unbound-method -- Reflect.apply intentionally supplies the caller-traced receiver
400 return Reflect.apply(target.createAgent, receiver, [ownerCtx, options])
401 }
402
403 /**
404 * Load a persisted session and resume an agent on it through the registered
405 * factory. Rejects if no factory is registered; the factory rejects if
406 * session persistence is not configured or persistence/setup fails.
407 * @param options - persisted identity, optional live parent, configuration, and setup.
408 * @returns the handle after setup, rollback-covered publication, and loop start complete.
409 */
410 async resume(options: ResumeAgentOptions): Promise<AgentHandle> {
411 const ownerCtx = this.ctx
412 const { target } = this.requireFactory()
413 const receiver = getTraceable(ownerCtx, target)
414 // oxlint-disable-next-line typescript/unbound-method -- Reflect.apply intentionally supplies the caller-traced receiver
415 return Reflect.apply(target.resume, receiver, [ownerCtx, options])
416 }
417
418 /**
419 * Register a live agent with source `startup`. Rejects if the id is already registered or a
420 * serial `agent/created` listener fails. Emits `agent/disposed`
421 * when the calling fiber is disposed — both with the agent's scope carrier
422 * (`scopeTarget(agent, agent)`): the subject is the agent in hand, so the
423 * emits are scope-filtered regardless of which context invoked `register`
424 * (calling through `agent.ctx` scopes EFFECTS; dispatch scoping always
425 * requires passing the carrier). The entry is a runtime root; factory-backed
426 * creation uses `options.parentAgent` for child ownership. Await the registration before using the agent.
427 * @param agent - the already-constructed agent to record in the store.
428 * @returns the awaitable Cordis effect disposer (single-shot; a repeat call
429 * returns undefined without awaiting an in-flight teardown). Exact
430 * identity is load-bearing: a composite (generator) effect that owns a
431 * teardown ORDER — the agent factory's lifecycle chain — must yield THIS
432 * function so Cordis nests the unregistration at that yield position;
433 * yielding a wrapper would leave it disposing as a concurrent sibling on
434 * owner unload, unregistering the agent (and emitting `agent/disposed`)
435 * while its final turn is still draining.
436 */
437 register(agent: Agent): ReturnType<Context['effect']> {
438 return this.ctx.effect(async function* (this: AgentRegistry) {
439 yield this.enter(agent, undefined)
440 await this.announce(agent, 'startup')
441 }.bind(this), 'agents.register()')
442 }
443
444 /**
445 * Insert an already-constructed agent without announcing it. This is the
446 * advanced ordered-lifecycle primitive used by the async agent factory: it
447 * first completes setup while the agent is unpublished, then assigns the
448 * returned detach closure into its pre-installed composite teardown before
449 * calling {@link announce}. Ordinary callers use {@link register}.
450 * @param agent - the prepared, unpublished agent.
451 * @param owner - explicitly supplied live runtime owner, or
452 * undefined for a top-level runtime root. This is runtime ownership, not
453 * the resumed session's durable parent lineage.
454 * @returns an idempotent closure that removes this exact entry and emits
455 * `agent/disposed` with listener failures contained. When called from a
456 * `agent/created` listener, removal and disposal wait until the serial
457 * creation dispatch settles.
458 */
459 enter(agent: Agent, owner: Agent | undefined): () => void {
460 const id = agent.id
461 if (id !== agent.session.id) {
462 throw new Error(`agent id "${id}" does not match session id "${agent.session.id}"`)
463 }
464 const carrier = scopeTarget(agent, agent)
465 // This is the authoritative collision boundary. Concurrent create/resume
466 // operations may both prepare, but only one exact entry can publish.
467 if (this.store.has(id)) throw new Error(`agent "${id}" is already registered`)
468 const entry: AgentEntry = {
469 id,
470 agent,
471 owner,
472 carrier,
473 announced: false,
474 announcing: false,
475 detachRequested: false,
476 }
477 this.store.set(id, entry)
478 let entered = true
479 const detach = (): void => {
480 if (!entered) return
481 entered = false
482 // Every callback reached by this creation dispatch must observe the same
483 // live entry, and disposal must follow creation. A listener may own
484 // the advanced detach capability, so make that ordering structural:
485 // visibility and the paired disposal are deferred until announce()'s
486 // serial dispatch has settled.
487 if (entry.announcing) {
488 entry.detachRequested = true
489 return
490 }
491 this.detachEntered(entry)
492 }
493 return detach
494 }
495
496 /** Remove one exact entered agent and emit its paired disposal when announced. */
497 private detachEntered(entry: AgentEntry): void {
498 entry.detachRequested = false
499 // A stale capability can never delete a later same-id lifecycle. The
500 // captured entry identity is the final boundary.
501 /* v8 ignore next -- enter() rejects replacement while this single-shot detach capability is live. */
502 if (this.store.get(entry.id) !== entry) return
503 this.store.delete(entry.id)
504 // An insertion rolled back before announce was never externally created,
505 // so emitting disposed would invent an impossible lifecycle edge. Marking
506 // happens before the created emit: if a later created listener throws,
507 // earlier listeners may already have observed it and must see disposal.
508 if (!entry.announced) return
509 this.emitDisposed(entry)
510 }
511
512 /** Emit the paired disposal edge through the entry's stable carrier. */
513 private emitDisposed(entry: AgentEntry): void {
514 const args: unknown[] = [entry.carrier, 'agent/disposed', { agent: entry.agent }]
515 for (const callback of this.ctx.events.dispatch('emit', args)) {
516 try {
517 const returned: unknown = callback(...args)
518 void Promise.resolve(returned).catch((error: unknown) => {
519 this.ctx.logger.warn(`agent "${entry.id}": agent/disposed listener rejected: ${String(error)}`)
520 })
521 } catch (error: unknown) {
522 this.ctx.logger.warn(`agent "${entry.id}": agent/disposed listener threw: ${String(error)}`)
523 }
524 }
525 }
526
527 /**
528 * Announce an agent previously inserted with {@link enter}.
529 * @param agent - the live inserted agent to announce.
530 * @param source - fresh creation, resume, clear, or compaction source.
531 * @param signal - optional factory initialization cancellation signal passed to listeners.
532 * @returns completion of the serial creation listeners; a listener failure rejects.
533 * @throws if `agent` is not the exact live registry entry for its id, or its
534 * creation announcement already began (including a reentrant call from a
535 * creation listener).
536 */
537 async announce(agent: Agent, source: SessionStartSource, signal?: AbortSignal): Promise<void> {
538 const entry = this.store.get(agent.id)
539 if (entry === undefined || entry.agent !== agent) {
540 throw new Error(`agent "${agent.id}" is not live in this registry`)
541 }
542 if (entry.announced || entry.announcing) {
543 throw new Error(`agent "${entry.id}" was already announced`)
544 }
545 // Mark before dispatch so a listener cannot recursively create a second
546 // lifecycle edge; detach still pairs a partially delivered first edge.
547 entry.announcing = true
548 entry.announced = true
549 try {
550 await this.ctx.serial(entry.carrier, 'agent/created', {
551 agent: entry.agent,
552 source,
553 ...signal === undefined ? {} : { signal },
554 })
555 } finally {
556 entry.announcing = false
557 if (entry.detachRequested) this.detachEntered(entry)
558 }
559 }
560
561 /**
562 * Look up a live agent.
563 * @param id - the shared agent/session id to look up.
564 * @returns the agent, or undefined when no live agent has that id.
565 */
566 get(id: SessionId): Agent | undefined {
567 return this.store.get(id)?.agent
568 }
569
570 /**
571 * Test whether a live agent was created through one exact parent agent's
572 * scoped context. Runtime ownership is independent of durable session
573 * lineage and remains unambiguous when unrelated providers reuse an id.
574 * @param id - the candidate child agent's shared agent/session id.
575 * @param owner - the expected runtime creator agent.
576 * @returns true only while the exact child entry is live under that owner.
577 */
578 isOwnedBy(id: SessionId, owner: Agent): boolean {
579 return this.store.get(id)?.owner === owner
580 }
581
582 /**
583 * All live agents, in registration order.
584 * @returns a fresh array; mutating it does not affect the registry.
585 */
586 list(): Agent[] {
587 return [...this.store.values()].map(entry => entry.agent)
588 }
589
590 /**
591 * All live top-level agents in registration order. A top-level agent was
592 * created without an owning agent context; durable session lineage does not
593 * affect this runtime relation, so a resumed fork may still be a root.
594 * @returns a fresh array; mutating it does not affect the registry.
595 */
596 roots(): Agent[] {
597 return [...this.store.values()]
598 .filter(entry => entry.owner === undefined)
599 .map(entry => entry.agent)
600 }
601
602 /** Reject new initiator boundaries while inherited continuations drain. */
603 private closeInitiators(): void {
604 if (this.initiatorState === 'active') this.initiatorState = 'closing'
605 }
606
607 /** Wait for returned-Promise boundaries, then invalidate retained references. */
608 private disposeInitiators(): Promise<void> {
609 return (this.initiatorDisposal ??= (async () => {
610 this.closeInitiators()
611 this.releaseReentrantInitiatorRuns()
612 if (this.activeInitiatorRuns !== 0) {
613 this.initiatorDrain ??= Promise.withResolvers<void>()
614 await this.initiatorDrain.promise
615 }
616 this.initiatorState = 'disposed'
617 this.initiators.disable()
618 this.initiatorRuns.disable()
619 })())
620 }
621
622 /** Establish one tracked initiator or clearing boundary. */
623 private runWithInitiator<T>(agent: Agent | undefined, operation: () => T): T {
624 if (this.initiatorState !== 'active') throw new Error(DISPOSED_INITIATOR_MESSAGE)
625 const run: InitiatorRun = {
626 active: true,
627 parent: this.initiatorRuns.getStore(),
628 }
629 this.activeInitiatorRuns += 1
630 let result: T
631 try {
632 result = this.initiatorRuns.run(run, () => this.initiators.run(agent, operation))
633 } catch (error: unknown) {
634 this.releaseInitiatorRun(run)
635 throw error
636 }
637 if (isPromise(result)) {
638 try {
639 void Promise.prototype.then.call(
640 result,
641 () => { this.releaseInitiatorRun(run) },
642 () => { this.releaseInitiatorRun(run) },
643 )
644 } catch {
645 // A branded Promise may expose a failing @@species. Observer setup did
646 // not attach, so preserve the exact return without leaking the run.
647 this.releaseInitiatorRun(run)
648 }
649 } else {
650 this.releaseInitiatorRun(run)
651 }
652 return result
653 }
654
655 /** Whether one unloading fiber owns this service's lifecycle. */
656 private hasLifecycleAncestor(candidate: Fiber): boolean {
657 let fiber = this.ctx.fiber
658 while (true) {
659 if (fiber === candidate) return true
660 const parent = fiber.parent.fiber
661 if (parent === fiber) return false
662 fiber = parent
663 }
664 }
665
666 private assertInitiatorsReadable(): void {
667 if (this.initiatorState === 'disposed') throw new Error(DISPOSED_INITIATOR_MESSAGE)
668 }
669
670 /** Exclude the boundary chain that initiated this teardown from its own drain. */
671 private releaseReentrantInitiatorRuns(): void {
672 let run = this.initiatorRuns.getStore()
673 while (run !== undefined) {
674 this.releaseInitiatorRun(run)
675 run = run.parent
676 }
677 }
678
679 private releaseInitiatorRun(run: InitiatorRun): void {
680 if (!run.active) return
681 run.active = false
682 this.activeInitiatorRuns -= 1
683 if (this.activeInitiatorRuns !== 0) return
684 this.initiatorDrain?.resolve()
685 this.initiatorDrain = undefined
686 }
687}
688
689export default AgentRegistry