返回源码地图

packages/core/agent-loop/src/index.ts

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

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

1/**
2 * Concrete agent-loop plugin: creates scoped ReactLoopAgents, publishes them
3 * through the agent/session registries, and owns their ordered teardown.
4 *
5 * @module @deepseek-ai/dsh-agent-loop
6 */
7import type { Volatile } from '@deepseek-ai/cosmokit'
8
9import { Context, FiberState, Service } from '@deepseek-ai/cordis'
10import { randomUUID } from 'node:crypto'
11import z from '@deepseek-ai/schemastery'
12import { z as zod } from 'zod'
13import { brandString } from '@deepseek-ai/dsh-brand'
14import type {
15 Agent,
16 AgentFactory,
17 AgentHandle,
18 AgentOptions,
19 AgentSetup,
20 CreateAgentOptions,
21 ResumeAgentOptions,
22 SessionStartSource,
23 TurnBoundaryProjection,
24} from '@deepseek-ai/dsh-agent'
25import { errorChain, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
26import { interruptedTurnClosers, SessionLogOffset, SessionPreparation, SessionSeq } from '@deepseek-ai/dsh-session'
27import type { Session, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
28import type {} from '@deepseek-ai/dsh-system-prompt'
29import type {} from '@deepseek-ai/dsh-tools'
30import type {} from '@deepseek-ai/dsh-session-projection'
31import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
32import { SessionPersistenceNotFoundError } from '@deepseek-ai/dsh-session-persistence'
33import type { SessionHandle, SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
34import { ReactLoopAgent } from './agent.ts'
35import { inboxProjectionDefinition } from './inbox.ts'
36import { DEFAULT_MAX_PARALLEL_TOOL_CALLS } from './constants.ts'
37import type {} from './runtime-context.ts'
38
39/** Fiber states that cannot own or serve a new lifecycle. */
40const INACTIVE_STATES: ReadonlySet<FiberState> = new Set([
41 FiberState.UNLOADING,
42 FiberState.DISPOSED,
43 FiberState.FAILED,
44])
45
46const turnBoundaryProjectionSchema: zod.ZodType<TurnBoundaryProjection> = zod.object({
47 openTurnStartSeq: zod.number().int().nonnegative().transform(SessionSeq).nullable(),
48 lastStepStartSeq: zod.number().int().nonnegative().transform(SessionSeq).nullable(),
49 lastStepBoundary: zod.object({
50 kind: zod.union([zod.literal('start'), zod.literal('end')]),
51 seq: zod.number().int().nonnegative().transform(SessionSeq),
52 }).nullable(),
53 lastTurn: zod.number().int().nonnegative(),
54})
55
56/** Host projection of agent turn and step boundaries. */
57export const turnBoundaryProjectionDefinition = {
58 key: 'turnBoundary',
59 stateVersion: 2,
60 stateSchema: turnBoundaryProjectionSchema,
61 init: () => ({
62 openTurnStartSeq: null,
63 lastStepStartSeq: null,
64 lastStepBoundary: null,
65 lastTurn: 0,
66 }),
67 apply: (state, event) => {
68 switch (event.type) {
69 case 'turn/start':
70 return {
71 ...state,
72 openTurnStartSeq: event.seq,
73 lastTurn: event.data.turn,
74 }
75 case 'turn/end':
76 return {
77 ...state,
78 openTurnStartSeq: null,
79 }
80 case 'step/start':
81 return {
82 ...state,
83 lastStepStartSeq: event.seq,
84 lastStepBoundary: { kind: 'start', seq: event.seq },
85 }
86 case 'step/end':
87 return {
88 ...state,
89 lastStepBoundary: { kind: 'end', seq: event.seq },
90 }
91 default:
92 return state
93 }
94 },
95} satisfies ProjectionDefinition<'turnBoundary', TurnBoundaryProjection>
96
97/** Factory-level ownership: live agent teardowns plus config startup work. */
98class FactoryOwnership {
99 private accepting = true
100 private readonly teardown = new AbortController()
101 private readonly inactive = Promise.withResolvers<void>()
102 private readonly liveAgents = new Set<() => Promise<void>>()
103 private startupTasks = new Set<Promise<void>>()
104
105 constructor(private readonly fiber: Context['fiber']) {}
106
107 /** Aborts (reason: `agent loop is not active` error) when factory teardown begins. */
108 get signal(): AbortSignal {
109 return this.teardown.signal
110 }
111
112 isActive(): boolean {
113 return this.accepting && !INACTIVE_STATES.has(this.fiber.state)
114 }
115
116 /** Track one live agent's shared teardown until it has run. */
117 track(dispose: () => Promise<void>): () => void {
118 this.liveAgents.add(dispose)
119 return () => { this.liveAgents.delete(dispose) }
120 }
121
122 /** Join config startup work that begins before an agent exists. */
123 trackStartup(job: Promise<void>): void {
124 this.startupTasks.add(job)
125 const forget = () => { this.startupTasks.delete(job) }
126 void job.then(forget, forget)
127 }
128
129 /** Join one public create/resume continuation; factory dispose awaits its settlement. */
130 trackWrapper(job: Promise<unknown>): void {
131 this.trackStartup(job.then(() => undefined, () => undefined))
132 }
133
134 /** Resolve `task`, or stop waiting when factory teardown begins. */
135 async waitWhileActive(job: Promise<void>): Promise<void> {
136 await Promise.race([job, this.inactive.promise])
137 }
138
139 async dispose(): Promise<void> {
140 this.accepting = false
141 this.teardown.abort(new Error('agent loop is not active'))
142 this.inactive.resolve()
143 await Promise.all([
144 ...[...this.liveAgents].map(dispose => dispose()),
145 ...this.startupTasks,
146 ])
147 }
148}
149
150/** Await `operation`, or throw the signal's reason as soon as it aborts. */
151async function raceAbort<T>(operation: PromiseLike<T> | T, signal: AbortSignal, id: SessionId): Promise<T> {
152 const toAbortError = (): Error => signal.reason instanceof Error
153 ? signal.reason
154 : new Error(`agent "${id}" creation aborted`, { cause: signal.reason })
155 if (signal.aborted) throw toAbortError()
156 const aborted = Promise.withResolvers<never>()
157 const listener = (): void => { aborted.reject(toAbortError()) }
158 signal.addEventListener('abort', listener, { once: true })
159 try {
160 return await Promise.race([Promise.resolve(operation), aborted.promise])
161 } finally {
162 signal.removeEventListener('abort', listener)
163 }
164}
165
166/** Start an abortable operation and release a value that arrives after cancellation. */
167async function raceAbortCall<T>(
168 operation: () => PromiseLike<T> | T,
169 signal: AbortSignal,
170 id: SessionId,
171 releaseAbandoned?: (value: T) => void,
172): Promise<T> {
173 if (signal.aborted) {
174 throw signal.reason instanceof Error
175 ? signal.reason
176 : new Error(`agent "${id}" creation aborted`, { cause: signal.reason })
177 }
178 const pending = Promise.resolve().then(operation)
179 try {
180 return await raceAbort(pending, signal, id)
181 } catch (error: unknown) {
182 // oxlint-disable-next-line typescript/no-unnecessary-condition -- the signal can abort while the operation is awaited.
183 if (signal.aborted && releaseAbandoned !== undefined) {
184 void pending.then(releaseAbandoned, () => undefined)
185 }
186 throw error
187 }
188}
189
190/** Reject an output-token cap that cannot be represented exactly on the request wire. */
191function assertAgentOptions(options: AgentOptions): void {
192 if (options.maxTokens !== undefined
193 && (!Number.isSafeInteger(options.maxTokens) || options.maxTokens <= 0)) {
194 throw new TypeError('agent maxTokens must be a positive safe integer')
195 }
196}
197
198/** One session's owned write handle plus the count of events already stored through it. */
199interface StoredSession {
200 readonly handle: SessionHandle
201 storedCount: number
202}
203
204/** Prepared-but-unpublished agent resources sharing one memoized teardown. */
205interface PreparedAgent {
206 agent: ReactLoopAgent
207 /** Aborts when the factory unloads, the caller cancels, or teardown begins — ends any setup await. */
208 signal: AbortSignal
209 /** Enter both registries and await creation listeners. */
210 publish(source: SessionStartSource): Promise<AgentHandle>
211 /** Reverse teardown: stop the machine, unregister, unwind the scope. Memoized. */
212 dispose(): Promise<void>
213}
214
215declare module '@deepseek-ai/cordis' {
216 interface Context {
217 agentLoop: AgentLoop
218 /**
219 * Launcher-owned exact session identities for configured agents, keyed by
220 * the agent's config `id` and set with `ctx.provide()` before any Loader
221 * entry mounts (see {@link CONFIGURED_AGENT_IDENTITIES_KEY}). A launcher
222 * owns identity because only it knows whether the session already exists,
223 * while the `cordis.yml` row keeps the model route as ordinary patchable
224 * config. An entry with no matching key keeps its configured identity.
225 */
226 configuredAgentIdentities?: ConfiguredAgentIdentities
227 }
228 interface Events {
229 /**
230 * A declarative agent entry failed before it could publish a live agent.
231 * Consumers that buffer work for the configured identity use this
232 * transient signal to reject that work instead of waiting forever. Normal
233 * factory teardown suppresses failures from the cancelled startup attempt.
234 * @param payload.sessionId - exact shared agent/session identity that failed startup.
235 * @param payload.error - persistence, setup, or publication failure.
236 * @mode emit
237 */
238 'agent-loop/config-start-failed'(payload: { sessionId: SessionId; error: unknown }): void
239 }
240}
241
242export { DEFAULT_MAX_PARALLEL_TOOL_CALLS }
243
244/**
245 * One launcher-selected session identity for a configured agent. `resume`
246 * distinguishes rehydrating existing persisted history from creating the
247 * session fresh under that exact id, which the two config keys express as
248 * `resumeSessionId` and `sessionId`.
249 */
250export interface LauncherAgentIdentity {
251 /** Exact session id to create fresh or resume. */
252 id: SessionId
253 /** Resume existing persisted history instead of creating the session fresh. */
254 resume: boolean
255}
256
257/** Launcher-selected identities keyed by the configured agent's `id`. */
258export interface ConfiguredAgentIdentities extends Readonly<Record<string, LauncherAgentIdentity>> {}
259
260/**
261 * Context key a launcher sets before any Loader entry mounts
262 * (`ctx.provide(CONFIGURED_AGENT_IDENTITIES_KEY, identities)`) to fix
263 * configured agents' session identities without a config key, so an overlay
264 * repointing the row's model route cannot drop them.
265 */
266export const CONFIGURED_AGENT_IDENTITIES_KEY = 'configuredAgentIdentities'
267
268/**
269 * Apply launcher-owned identities over the configured agents, replacing both
270 * identity keys for every entry the launcher named so a config-supplied
271 * identity can never survive alongside a launcher-supplied one.
272 * @param agents - the configured agent entries.
273 * @param identities - launcher identities keyed by configured agent `id`, or `undefined`.
274 * @returns the entries with launcher-owned identities applied.
275 */
276function applyLauncherIdentities(
277 agents: Config['agents'],
278 identities: ConfiguredAgentIdentities | undefined,
279): Config['agents'] {
280 if (identities === undefined) return agents
281 return agents.map((agent) => {
282 const identity = identities[agent.id]
283 if (identity === undefined) return agent
284 const { sessionId: _sessionId, resumeSessionId: _resumeSessionId, ...rest } = agent
285 return identity.resume
286 ? { ...rest, resumeSessionId: identity.id }
287 : { ...rest, sessionId: identity.id }
288 })
289}
290
291/** Agent-loop plugin configuration. */
292export interface Config {
293 /**
294 * Maximum parallel-safe calls in flight per agent step. `1` is serial;
295 * omission defaults to {@link DEFAULT_MAX_PARALLEL_TOOL_CALLS}.
296 */
297 maxParallelToolCalls: Volatile<number>
298 /** Agents created or resumed at plugin startup. */
299 agents: (AgentOptions & {
300 /** Stable config label used in logs and as the fresh combined-id prefix. */
301 id: string
302 /** Optional stable identity; remounts resume its materialized history, while first use creates it fresh. */
303 sessionId?: SessionId
304 /** Optional workspace for a fresh session. */
305 cwd?: string
306 /** Persisted session to resume instead of creating a fresh session. */
307 resumeSessionId?: SessionId
308 })[]
309}
310
311/** Reject self-contained identity conflicts before any configured agent starts. */
312function validateConfiguredAgents(agents: Config['agents']): void {
313 const exactIdentities = new Map<SessionId, string>()
314 for (const { id, sessionId, resumeSessionId } of agents) {
315 const hasResumeId = resumeSessionId !== undefined && resumeSessionId !== ''
316 if (sessionId !== undefined && hasResumeId) {
317 throw new Error(`agent "${id}": sessionId and resumeSessionId are mutually exclusive`)
318 }
319 const exactIdentity = hasResumeId ? resumeSessionId : sessionId
320 if (exactIdentity === undefined) continue
321 const firstId = exactIdentities.get(exactIdentity)
322 if (firstId !== undefined) {
323 throw new Error(`agents "${firstId}" and "${id}" use duplicate exact session identity "${exactIdentity}"`)
324 }
325 exactIdentities.set(exactIdentity, id)
326 }
327}
328
329/** Concrete agent factory and driver service. */
330export class AgentLoop extends Service implements AgentFactory {
331 static inject = ['agents', 'sessions', 'llm', 'tools', 'systemPrompt', 'sessionProjections']
332
333 /** Runtime schema for declarative agents. */
334 static Config: z<{ agents?: Config['agents']; maxParallelToolCalls?: number }, Config> = z.object({
335 maxParallelToolCalls: z.number().step(1).min(1).default(DEFAULT_MAX_PARALLEL_TOOL_CALLS).volatile(),
336 agents: z.array(z.object({
337 id: z.string().required(),
338 sessionId: z.string().min(1),
339 provider: z.string(),
340 model: z.string(),
341 reasoningEffort: z.string().min(1) as z<ReturnType<typeof ReasoningEffortId>>,
342 maxTokens: z.number().step(1).min(1).max(Number.MAX_SAFE_INTEGER),
343 cwd: z.string(),
344 resumeSessionId: z.string(),
345 })).default([]),
346 }) as z<{ agents?: Config['agents']; maxParallelToolCalls?: number }, Config>
347
348 /** Validated configuration owned by the agent-loop service. */
349 readonly config: Config
350 private readonly ownership: FactoryOwnership
351 /** Plain holder prevents Cordis from re-tracing the factory's dependency context through a caller shadow. */
352 private readonly runtime: { ctx: Context }
353
354 constructor(ctx: Context, config: Config) {
355 super(ctx, 'agentLoop')
356
357 this.config = {
358 agents: applyLauncherIdentities(config.agents, ctx.get(CONFIGURED_AGENT_IDENTITIES_KEY)),
359 maxParallelToolCalls: config.maxParallelToolCalls,
360 }
361 validateConfiguredAgents(this.config.agents)
362 // Register only after every config validation above has passed, so a
363 // rejected constructor leaves no projection unit behind.
364 ctx.sessionProjections.register(turnBoundaryProjectionDefinition)
365 ctx.sessionProjections.register(inboxProjectionDefinition)
366 this.ownership = new FactoryOwnership(ctx.fiber)
367 this.runtime = { ctx }
368 ctx.effect(() => () => this.ownership.dispose(), 'agentLoop.transactions()')
369 ctx.effect(() => ctx.agents.setFactory(this), 'agentLoop.setFactory()')
370 ctx.systemPrompt.variable('provider', context => context.agent?.options.provider)
371 ctx.systemPrompt.variable('model', context => context.agent?.options.model)
372 ctx.systemPrompt.variable('cwd', context => context.agent?.session.header.cwd)
373
374 for (const { id, sessionId, cwd, resumeSessionId, ...options } of this.config.agents) {
375 const meta = cwd === undefined ? {} : { cwd }
376 if (resumeSessionId === undefined || resumeSessionId === '') {
377 const configuredId = sessionId ?? brandString<SessionId>(`${id}-session-${randomUUID()}`)
378 const persistence = sessionId === undefined ? undefined : ctx.get('sessionPersistence')
379 if (persistence === undefined) {
380 const startup = this.create(configuredId, options, meta).then(() => undefined, (error: unknown) => {
381 this.reportConfiguredStartupFailure(id, 'restore', configuredId, error)
382 })
383 this.ownership.trackStartup(startup)
384 } else {
385 const startup = this.restoreOrCreateConfigured(ctx, persistence, configuredId, options, meta).catch((error: unknown) => {
386 this.reportConfiguredStartupFailure(id, 'restore', configuredId, error)
387 })
388 this.ownership.trackStartup(startup)
389 }
390 continue
391 }
392 ctx.effect(() => {
393 const fiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {
394 void this.resumeWith(ctx, childCtx.sessionPersistence, {
395 resumeSessionId,
396 agentOptions: options,
397 }).catch((error: unknown) => {
398 this.reportConfiguredStartupFailure(id, 'resume', resumeSessionId, error)
399 })
400 })
401 return fiber.dispose
402 }, `agentLoop.resume(${id})`)
403 }
404 }
405
406 /** Report a contained declarative-start failure to identity-bound consumers. */
407 private reportConfiguredStartupFailure(
408 configId: string,
409 action: 'restore' | 'resume',
410 sessionId: SessionId,
411 error: unknown,
412 ): void {
413 if (!this.ownership.isActive()) return
414 this.ctx.logger.warn(`agent "${configId}": config-driven ${action} of "${sessionId}" failed: ${errorChain(error)}`)
415 const args: unknown[] = ['agent-loop/config-start-failed', { sessionId, error }]
416 for (const callback of this.ctx.events.dispatch('emit', args)) {
417 try {
418 const returned: unknown = callback(...args)
419 void Promise.resolve(returned).catch((listenerError: unknown) => {
420 this.ctx.logger.warn(`agent "${configId}": config-start-failed listener rejected: ${errorChain(listenerError)}`)
421 })
422 } catch (listenerError: unknown) {
423 this.ctx.logger.warn(`agent "${configId}": config-start-failed listener threw: ${errorChain(listenerError)}`)
424 }
425 }
426 }
427
428 /** Restore a materialized exact config identity on remount, or create it on first use. */
429 private async restoreOrCreateConfigured(
430 ownerCtx: Context,
431 persistence: SessionPersistence,
432 sessionId: SessionId,
433 agentOptions: AgentOptions,
434 meta: Pick<SessionHeader, 'cwd'>,
435 ): Promise<void> {
436 await this.waitForDrainingConfiguredIdentity(ownerCtx, sessionId)
437 if (!this.ownership.isActive()) return
438 try {
439 await this.resumeWith(ownerCtx, persistence, { resumeSessionId: sessionId, agentOptions })
440 return
441 } catch (error: unknown) {
442 if (!this.ownership.isActive()) return
443 // Only a genuinely absent stored session falls back to first creation;
444 // corruption, ownership conflicts, and backend failures stay loud.
445 if (!(error instanceof SessionPersistenceNotFoundError)) throw error
446 }
447 await this.create(sessionId, agentOptions, meta)
448 }
449
450 /** Wait for a draining same-id lifecycle to finish registry teardown. */
451 private async waitForDrainingConfiguredIdentity(ownerCtx: Context, sessionId: SessionId): Promise<void> {
452 // Only an id still occupying a registry needs waiting for; a live healthy
453 // occupant is a collision the create/resume below will surface itself.
454 if (ownerCtx.agents.get(sessionId) === undefined && ownerCtx.sessions.get(sessionId) === undefined) return
455
456 const released = Promise.withResolvers<void>()
457 const checkReleased = (): void => {
458 if (ownerCtx.agents.get(sessionId) === undefined && ownerCtx.sessions.get(sessionId) === undefined) {
459 released.resolve()
460 }
461 }
462 const disposeAgentListener = ownerCtx.on('agent/disposed', () => { checkReleased() })
463 const disposeSessionListener = ownerCtx.on('session/disposed', checkReleased)
464 try {
465 checkReleased()
466 await this.ownership.waitWhileActive(released.promise)
467 } finally {
468 disposeAgentListener()
469 disposeSessionListener()
470 }
471 }
472
473 /**
474 * Construct the driver, scope, and one memoized reverse teardown for a new
475 * agent. The teardown is registered with the factory and the owner fiber
476 * BEFORE publication, so a mid-setup unload rolls everything back; `signal`
477 * fuses caller cancellation with lifecycle teardown for setup awaits.
478 */
479 private prepare(
480 ownerCtx: Context,
481 id: SessionId,
482 options: AgentOptions,
483 session: Session,
484 callerSignal?: AbortSignal,
485 handle?: SessionHandle,
486 parentAgent?: Agent,
487 ): PreparedAgent {
488 assertAgentOptions(options)
489 ownerCtx.fiber.assertActive()
490 // Every caller reaches prepare() synchronously from a service method
491 // whose Cordis dispatch already requires the live factory fiber, or
492 // re-checks ownership itself after its awaits (resume's load barrier).
493 /* v8 ignore next -- unreachable backstop, see above */
494 if (!this.ownership.isActive()) throw new Error('agent loop is not active')
495 if (callerSignal?.aborted) {
496 throw callerSignal.reason instanceof Error
497 ? callerSignal.reason
498 : new Error(`agent "${id}" creation aborted`, { cause: callerSignal.reason })
499 }
500 const loopCtx = this.runtime.ctx
501
502 // Deactivation fuses three owners, each with its own reason: the caller's
503 // cancellation signal, the owner fiber's unload, and factory teardown.
504 // It is registered BEFORE any resource exists, over mutable slots, so an
505 // unload arriving while the scope is still minting finds a working
506 // disposer instead of a leak.
507 const abort = new AbortController()
508 const onCallerAbort = (): void => {
509 abort.abort(callerSignal?.reason instanceof Error
510 ? callerSignal.reason
511 : new Error(`agent "${id}" creation aborted`, { cause: callerSignal?.reason }))
512 }
513 const onFactoryTeardown = (): void => { abort.abort(this.ownership.signal.reason) }
514 callerSignal?.addEventListener('abort', onCallerAbort, { once: true })
515 this.ownership.signal.addEventListener('abort', onFactoryTeardown, { once: true })
516
517 let machine: ReactLoopAgent | undefined
518 let detachSession: (() => void) | undefined
519 let detachAgent: (() => void) | undefined
520 let disposing: Promise<void> | undefined
521 let publication: ReturnType<typeof Promise.withResolvers<void>> | undefined
522 const machineReady = Promise.withResolvers<void>()
523 // Reverse teardown, memoized so every racing owner awaits one quiescence:
524 // stop the machine, drain and close the session's write path, leave the
525 // registries, unwind the scope, release bookkeeping.
526 const dispose = (ownerTriggered = false): Promise<void> => (disposing ??= (async () => {
527 abort.abort(new Error(`agent "${id}" lifecycle disposed`))
528 callerSignal?.removeEventListener('abort', onCallerAbort)
529 this.ownership.signal.removeEventListener('abort', onFactoryTeardown)
530 // Teardown failures are collected, never swallowed: registry, scope,
531 // and ownership cleanup always run to quiescence, then the memoized
532 // disposal rejects with what failed so every racing owner observes it.
533 const failures: unknown[] = []
534 try {
535 // Creation listeners retain the session and scope through their awaits.
536 if (publication !== undefined) await publication.promise
537 // Disposal IS a disposed-cause cancel followed by quiescence. New work
538 // sent after this point is the sender's bug — the registries are about
539 // to drop the agent, so nothing should still hold it.
540 /* v8 ignore next -- Cordis effect teardown waits for synchronous setup before observing the machine slot. */
541 if (machine === undefined) await machineReady.promise
542 /* v8 ignore next -- setup failure untracks this disposer before resolving without a machine. */
543 if (machine !== undefined) {
544 machine.cancel({ kind: 'disposed' })
545 await machine.whenIdle()
546 await machine.scope.dispose()
547 }
548 } catch (error: unknown) {
549 failures.push(error)
550 }
551 // The loop above committed its closing events synchronously into the
552 // session; handle close drains them durably before releasing the write
553 // path. The close drain can be the first operation that surfaces a
554 // durability failure, so its error is retained, not logged away.
555 try {
556 await handle?.close()
557 } catch (error: unknown) {
558 failures.push(error)
559 }
560 try {
561 detachAgent?.()
562 detachSession?.()
563 } finally {
564 untrack()
565 if (!ownerTriggered) await unfollowOwner()
566 }
567 if (failures.length === 1) throw failures[0]
568 if (failures.length > 1) {
569 throw new AggregateError(failures, `agent "${id}" disposal failed`)
570 }
571 })())
572 const untrack = this.ownership.track(dispose)
573 let unfollowOwner: () => Promise<void> | void
574 try {
575 unfollowOwner = ownerCtx.effect(function* () {
576 machine = new ReactLoopAgent(loopCtx, id, options, session)
577 machineReady.resolve()
578 yield machine.scope.rawDispose
579 yield () => {
580 // Owner disposal owns the same quiescence boundary. Its teardown skips
581 // unregistering this already-running owner effect from inside itself.
582 if (disposing !== undefined) return
583 abort.abort(new Error(`agent "${id}" setup aborted: owner disposed during setup`))
584 return dispose(true)
585 }
586 }, `agentLoop.lifecycle(${id})`)
587 /* v8 ignore start -- ctx.effect throws only on an inactive fiber, which assertActive() above already rejected */
588 } catch (error: unknown) {
589 machineReady.resolve()
590 untrack()
591 callerSignal?.removeEventListener('abort', onCallerAbort)
592 this.ownership.signal.removeEventListener('abort', onFactoryTeardown)
593 throw error
594 }
595 /* v8 ignore stop */
596
597 const assertLive = (): void => {
598 if (!abort.signal.aborted) return
599 // Every fused abort source carries an Error reason: onCallerAbort and
600 // raceAbort wrap non-Error caller reasons, and the factory/lifecycle
601 // owners abort with constructed Errors.
602 /* v8 ignore next -- unreachable String() arm, see above */
603 throw abort.signal.reason instanceof Error ? abort.signal.reason : new Error(String(abort.signal.reason))
604 }
605 try {
606 /* v8 ignore next -- a synchronous effect exhausts the generator before returning */
607 if (machine === undefined) throw new Error(`agent "${id}" lifecycle did not construct its driver`)
608 const agent = machine
609 assertLive()
610
611 return {
612 agent,
613 signal: abort.signal,
614 publish: async (source) => {
615 publication = Promise.withResolvers<void>()
616 try {
617 assertLive()
618 detachSession = agent.ctx.sessions.enter(session)
619 // The mounted backend routes announced live events into the active
620 // write handle by session id; the loop only owns the handle itself.
621 detachAgent = loopCtx.agents.enter(agent, parentAgent)
622 agent.ctx.sessions.announce(session)
623 assertLive()
624 await loopCtx.agents.announce(agent, source, abort.signal)
625 assertLive()
626 return { agent, dispose }
627 } finally {
628 publication.resolve()
629 publication = undefined
630 }
631 },
632 dispose,
633 }
634 } catch (error: unknown) {
635 machineReady.resolve()
636 // Rollback swallows a disposal rejection: the setup failure is primary.
637 void dispose().catch(() => {})
638 throw error
639 }
640 }
641
642 /**
643 * Create an agent and session under one caller-supplied identity, owned by
644 * the accessing fiber. Constructor-driven config calls mint a fresh combined
645 * id before entering this boundary. When a persistence backend is mounted,
646 * the session's durable identity and any seed are stored before publication.
647 * @param id - shared agent/session identity.
648 * @param options - concrete loop options.
649 * @param meta - optional fresh-session workspace metadata.
650 * @returns the published running agent.
651 */
652 async create(id: SessionId, options: AgentOptions = {}, meta: Pick<SessionHeader, 'cwd'> = {}): Promise<Agent> {
653 using preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(id, { meta }))
654 const stored = await this.createStoredSession(preparation.session)
655 let prepared: PreparedAgent
656 try {
657 prepared = this.prepare(this.ctx, id, options, preparation.session, undefined, stored?.handle)
658 } catch (error: unknown) {
659 await stored?.handle.close().catch(() => {})
660 throw error
661 }
662 return (await this.initializeAgent(prepared, async () => {
663 await this.appendUnstoredSuffix(stored, preparation.session)
664 return await prepared.publish('startup')
665 })).agent
666 }
667
668 /**
669 * Take a fresh session's write ownership when persistence is mounted.
670 * Nothing is appended here: the constructor seed (which never re-emits
671 * through `session/event`) is stored by `appendUnstoredSuffix` at the
672 * publication commit point, so a failed or cancelled validation or setup
673 * closes an unmaterialized handle and leaves no stored residue — the same
674 * id can be created again.
675 * @param session - the unpublished session to store.
676 * @param signal - optional cancellation forwarded to the backend create.
677 * @returns the owned handle and stored cursor, or `undefined` without a backend.
678 */
679 private async createStoredSession(session: Session, signal?: AbortSignal): Promise<StoredSession | undefined> {
680 const persistence = this.runtime.ctx.get('sessionPersistence')
681 if (persistence === undefined) return undefined
682 const handle = await persistence.create(session.header, {
683 inheritedEventCount: session.inheritedEventCount,
684 ...signal === undefined ? {} : { signal },
685 })
686 return { handle, storedCount: 0 }
687 }
688
689 /**
690 * Durably store the session events appended since the last stored cursor.
691 * Pre-publication appends (constructor seed markers, setup-window events
692 * such as delegation policy records) never re-emit through `session/event`,
693 * so publication must flush them through the handle before live events
694 * start routing into it.
695 * @param stored - the session's owned handle and stored cursor, if any.
696 * @param session - the unpublished session whose suffix is stored.
697 */
698 private async appendUnstoredSuffix(stored: StoredSession | undefined, session: Session): Promise<void> {
699 if (stored === undefined) return
700 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
701 const suffix = session.snapshotEvents(SessionLogOffset(stored.storedCount))
702 if (suffix.length > 0) await stored.handle.append(suffix)
703 // Advance by what was stored, not to `session.seq`: an event appended
704 // during the await must stay unstored for the next flush.
705 stored.storedCount += suffix.length
706 }
707
708 /**
709 * Create an owned agent on a caller-supplied session id.
710 * @param ownerCtx - caller context that structurally owns the lifecycle.
711 * @param options - identities, optional live parent, session seed/metadata, loop options, setup, and cancellation.
712 * @returns the published handle.
713 */
714 async createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> {
715 const preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(options.sessionId, {
716 ...options.seed === undefined ? {} : { seed: options.seed },
717 ...options.meta === undefined ? {} : { meta: options.meta },
718 ...options.inheritedEventCount === undefined ? {} : { inheritedEventCount: options.inheritedEventCount },
719 }))
720 const published = (async () => {
721 let stored: StoredSession | undefined
722 try {
723 // raceAbortCall normalizes a pre-aborted or mid-create abort and
724 // closes a handle that finishes creating after abandonment.
725 stored = options.signal === undefined
726 ? await this.createStoredSession(preparation.session)
727 : await raceAbortCall(
728 () => this.createStoredSession(preparation.session, options.signal),
729 options.signal,
730 options.sessionId,
731 (abandoned) => { void abandoned?.handle.close().catch(() => {}) },
732 )
733 } catch (error: unknown) {
734 preparation[Symbol.dispose]()
735 throw error
736 }
737 return this.setupAndPublish(
738 ownerCtx,
739 options.sessionId,
740 preparation,
741 options.agentOptions ?? {},
742 options.setup,
743 options.signal,
744 'startup',
745 stored,
746 options.parentAgent,
747 )
748 })()
749 this.ownership.trackWrapper(published)
750 return published
751 }
752
753 /** Prepare one Agent around an acquired Session, run setup, and publish it. */
754 private async setupAndPublish(
755 ownerCtx: Context,
756 id: SessionId,
757 preparation: SessionPreparation,
758 agentOptions: AgentOptions,
759 setup: AgentSetup | undefined,
760 signal: AbortSignal | undefined,
761 source: SessionStartSource,
762 stored?: StoredSession,
763 parentAgent?: Agent,
764 ): Promise<AgentHandle> {
765 using ownedPreparation = preparation
766 const session = ownedPreparation.session
767 let prepared: PreparedAgent
768 try {
769 prepared = this.prepare(ownerCtx, id, agentOptions, session, signal, stored?.handle, parentAgent)
770 } catch (error: unknown) {
771 await stored?.handle.close().catch(() => {})
772 throw error
773 }
774 return await this.initializeAgent(prepared, async () => {
775 const setupCommit = await raceAbort(setup?.(prepared.agent.ctx, prepared.agent), prepared.signal, id)
776 setupCommit?.commit()
777 await this.appendUnstoredSuffix(stored, session)
778 return await prepared.publish(source)
779 })
780 }
781
782 private async initializeAgent(prepared: PreparedAgent, initialize: () => Promise<AgentHandle>): Promise<AgentHandle> {
783 try {
784 return await prepared.agent.runMaintenance(async () => {
785 try {
786 return await initialize()
787 } catch (error: unknown) {
788 // Teardown owns inbox cleanup and may already have removed its projection.
789 prepared.agent.cancel({ kind: 'disposed' }, { keepInbox: true })
790 throw error
791 }
792 })
793 } catch (error: unknown) {
794 // Rollback swallows a disposal rejection (a failing final handle close):
795 // the setup failure is the primary error the caller must see.
796 await prepared.dispose().catch(() => {})
797 throw error
798 }
799 }
800
801 /**
802 * Resume an owned agent from the configured persistence service.
803 * @param ownerCtx - caller context that owns load, setup, and the live lifecycle.
804 * @param options - persisted identity, optional live parent, loop options, setup, and cancellation.
805 * @returns the published handle.
806 */
807 async resume(ownerCtx: Context, options: ResumeAgentOptions): Promise<AgentHandle> {
808 const persistence = this.runtime.ctx.get('sessionPersistence')
809 if (persistence === undefined) {
810 throw new Error('cannot resume: session persistence is not configured (load a dsh-session-persistence backend)')
811 }
812 return this.resumeWith(ownerCtx, persistence, options)
813 }
814
815 /** Resume through an explicit persistence handle used by the deferred config path. */
816 private resumeWith(
817 ownerCtx: Context,
818 persistence: SessionPersistence,
819 options: ResumeAgentOptions,
820 ): Promise<AgentHandle> {
821 const id = options.resumeSessionId
822 const published = (async () => {
823 // The open and read may outlive their owner: race them against caller
824 // cancellation, owner-fiber unload, and factory teardown so a
825 // never-settling backend cannot pin the identity.
826 const ownerAbort = new AbortController()
827 const unfollowOwner = ownerCtx.effect(() => () => {
828 ownerAbort.abort(new Error(`agent "${id}" setup aborted: owner disposed during setup`))
829 }, `agentLoop.resume-load(${id})`)
830 const fused = AbortSignal.any([
831 ...options.signal === undefined ? [] : [options.signal],
832 ownerAbort.signal,
833 this.ownership.signal,
834 ])
835 let handle: SessionHandle | undefined
836 let stored: StoredSession | undefined
837 let preparation: SessionPreparation | undefined
838 try {
839 try {
840 // Taking write ownership FIRST excludes a concurrent resume of the
841 // same id (in this process, a live agent's handle holds the claim).
842 handle = await raceAbortCall(
843 () => persistence.open(id, 'write', { signal: fused }),
844 fused,
845 id,
846 (abandoned) => { void abandoned.close() },
847 )
848 // Semantic crash repair is the agent layer's job: persistence hands
849 // back the physically valid log; an interrupted final turn receives
850 // synthetic closers (missing tool errors, step/end, turn/end) that
851 // are appended through the same handle as an ordinary batch.
852 const coldRead = await handle.read(0, undefined, { signal: fused })
853 fused.throwIfAborted()
854 const persisted = coldRead.events
855 const closers = interruptedTurnClosers(persisted)
856 if (closers.length > 0) await handle.append(closers)
857 preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(id, {
858 seed: [...persisted, ...closers],
859 meta: structuredClone(handle.header),
860 inheritedEventCount: handle.inheritedEventCount,
861 eventState: coldRead.eventState,
862 }))
863 stored = { handle, storedCount: persisted.length + closers.length }
864 await this.appendUnstoredSuffix(stored, preparation.session)
865 } finally {
866 await unfollowOwner()
867 }
868 ownerCtx.fiber.assertActive()
869 if (!this.ownership.isActive()) throw new Error('agent loop is not active')
870 const owned = stored
871 handle = undefined // ownership passes to setupAndPublish/prepare
872 return await this.setupAndPublish(
873 ownerCtx,
874 id,
875 preparation,
876 options.agentOptions ?? {},
877 options.setup,
878 options.signal,
879 'resume',
880 owned,
881 options.parentAgent,
882 )
883 } finally {
884 preparation?.[Symbol.dispose]()
885 await handle?.close().catch(() => {})
886 }
887 })()
888 this.ownership.trackWrapper(published)
889 return published
890 }
891}
892
893export default AgentLoop