index.ts 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776
  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. */
  7. import { Context, FiberState, Service } from '@deepseek-ai/cordis'
  8. import { randomUUID } from 'node:crypto'
  9. import z from '@deepseek-ai/schemastery'
  10. import { z as zod } from 'zod'
  11. import { brandString } from '@deepseek-ai/dsh-brand'
  12. import { emitAgentEvent } from '@deepseek-ai/dsh-agent'
  13. import type {
  14. Agent,
  15. AgentFactory,
  16. AgentHandle,
  17. AgentOptions,
  18. AgentSetup,
  19. CreateAgentOptions,
  20. ResumeAgentOptions,
  21. SessionStartSource,
  22. TurnBoundaryProjection,
  23. } from '@deepseek-ai/dsh-agent'
  24. import { errorChain, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  25. import type {} from '@deepseek-ai/dsh-settings'
  26. import { SessionPreparation } from '@deepseek-ai/dsh-session'
  27. import type { Session, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  28. import type {} from '@deepseek-ai/dsh-system-prompt'
  29. import type {} from '@deepseek-ai/dsh-tools'
  30. import type {} from '@deepseek-ai/dsh-session-projection'
  31. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  32. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  33. import { ReactLoopAgent } from './agent.ts'
  34. import { DEFAULT_MAX_PARALLEL_TOOL_CALLS } from './constants.ts'
  35. /** Fiber states that cannot own or serve a new lifecycle. */
  36. const INACTIVE_STATES: ReadonlySet<FiberState> = new Set([
  37. FiberState.UNLOADING,
  38. FiberState.DISPOSED,
  39. FiberState.FAILED,
  40. ])
  41. const turnBoundaryProjectionSchema: zod.ZodType<TurnBoundaryProjection> = zod.object({
  42. openTurnStartSeq: zod.number().int().nonnegative().nullable(),
  43. lastStepStartSeq: zod.number().int().nonnegative().nullable(),
  44. lastStepBoundary: zod.object({
  45. kind: zod.union([zod.literal('start'), zod.literal('end')]),
  46. seq: zod.number().int().nonnegative(),
  47. }).nullable(),
  48. lastTurn: zod.number().int().nonnegative(),
  49. })
  50. /** Host projection of agent turn and step boundaries. */
  51. export const turnBoundaryProjectionDefinition = {
  52. key: 'turnBoundary',
  53. stateVersion: 2,
  54. stateSchema: turnBoundaryProjectionSchema,
  55. init: () => ({
  56. openTurnStartSeq: null,
  57. lastStepStartSeq: null,
  58. lastStepBoundary: null,
  59. lastTurn: 0,
  60. }),
  61. apply: (state, event) => {
  62. switch (event.type) {
  63. case 'turn/start':
  64. return {
  65. ...state,
  66. openTurnStartSeq: event.seq,
  67. lastTurn: event.data.turn,
  68. }
  69. case 'turn/end':
  70. return {
  71. ...state,
  72. openTurnStartSeq: null,
  73. }
  74. case 'step/start':
  75. return {
  76. ...state,
  77. lastStepStartSeq: event.seq,
  78. lastStepBoundary: { kind: 'start', seq: event.seq },
  79. }
  80. case 'step/end':
  81. return {
  82. ...state,
  83. lastStepBoundary: { kind: 'end', seq: event.seq },
  84. }
  85. default:
  86. return state
  87. }
  88. },
  89. } satisfies ProjectionDefinition<'turnBoundary', TurnBoundaryProjection>
  90. /** Factory-level ownership: live agent teardowns plus config startup work. */
  91. class FactoryOwnership {
  92. private accepting = true
  93. private readonly teardown = new AbortController()
  94. private readonly inactive = Promise.withResolvers<void>()
  95. private readonly liveAgents = new Set<() => Promise<void>>()
  96. private startupTasks = new Set<Promise<void>>()
  97. constructor(private readonly fiber: Context['fiber']) {}
  98. /** Aborts (reason: `agent loop is not active` error) when factory teardown begins. */
  99. get signal(): AbortSignal {
  100. return this.teardown.signal
  101. }
  102. isActive(): boolean {
  103. return this.accepting && !INACTIVE_STATES.has(this.fiber.state)
  104. }
  105. /** Track one live agent's shared teardown until it has run. */
  106. track(dispose: () => Promise<void>): () => void {
  107. this.liveAgents.add(dispose)
  108. return () => { this.liveAgents.delete(dispose) }
  109. }
  110. /** Join config startup work that begins before an agent exists. */
  111. trackStartup(job: Promise<void>): void {
  112. this.startupTasks.add(job)
  113. const forget = () => { this.startupTasks.delete(job) }
  114. void job.then(forget, forget)
  115. }
  116. /** Join one public create/resume continuation; factory dispose awaits its settlement. */
  117. trackWrapper(job: Promise<unknown>): void {
  118. this.trackStartup(job.then(() => undefined, () => undefined))
  119. }
  120. /** Resolve `task`, or stop waiting when factory teardown begins. */
  121. async waitWhileActive(job: Promise<void>): Promise<void> {
  122. await Promise.race([job, this.inactive.promise])
  123. }
  124. async dispose(): Promise<void> {
  125. this.accepting = false
  126. this.teardown.abort(new Error('agent loop is not active'))
  127. this.inactive.resolve()
  128. await Promise.all([
  129. ...[...this.liveAgents].map(dispose => dispose()),
  130. ...this.startupTasks,
  131. ])
  132. }
  133. }
  134. /** Await `operation`, or throw the signal's reason as soon as it aborts. */
  135. async function raceAbort<T>(operation: PromiseLike<T> | T, signal: AbortSignal, id: SessionId): Promise<T> {
  136. const toAbortError = (): Error => signal.reason instanceof Error
  137. ? signal.reason
  138. : new Error(`agent "${id}" creation aborted`, { cause: signal.reason })
  139. if (signal.aborted) throw toAbortError()
  140. const aborted = Promise.withResolvers<never>()
  141. const listener = (): void => { aborted.reject(toAbortError()) }
  142. signal.addEventListener('abort', listener, { once: true })
  143. try {
  144. return await Promise.race([Promise.resolve(operation), aborted.promise])
  145. } finally {
  146. signal.removeEventListener('abort', listener)
  147. }
  148. }
  149. /** Start an abortable operation and release a value that arrives after cancellation. */
  150. async function raceAbortCall<T>(
  151. operation: () => PromiseLike<T> | T,
  152. signal: AbortSignal,
  153. id: SessionId,
  154. releaseAbandoned?: (value: T) => void,
  155. ): Promise<T> {
  156. if (signal.aborted) {
  157. throw signal.reason instanceof Error
  158. ? signal.reason
  159. : new Error(`agent "${id}" creation aborted`, { cause: signal.reason })
  160. }
  161. const pending = Promise.resolve().then(operation)
  162. try {
  163. return await raceAbort(pending, signal, id)
  164. } catch (error: unknown) {
  165. // oxlint-disable-next-line typescript/no-unnecessary-condition -- the signal can abort while the operation is awaited.
  166. if (signal.aborted && releaseAbandoned !== undefined) {
  167. void pending.then(releaseAbandoned, () => undefined)
  168. }
  169. throw error
  170. }
  171. }
  172. /** Resolve the deployment-wide scheduler cap at the owning config boundary. */
  173. function resolveMaxParallelToolCalls(value: number | undefined): number {
  174. const maxParallelToolCalls = value ?? DEFAULT_MAX_PARALLEL_TOOL_CALLS
  175. if (!Number.isInteger(maxParallelToolCalls) || maxParallelToolCalls < 1) {
  176. throw new Error('maxParallelToolCalls must be a positive integer')
  177. }
  178. return maxParallelToolCalls
  179. }
  180. /** Reject an output-token cap that cannot be represented exactly on the request wire. */
  181. function assertAgentOptions(options: AgentOptions): void {
  182. if (options.maxTokens !== undefined
  183. && (!Number.isSafeInteger(options.maxTokens) || options.maxTokens <= 0)) {
  184. throw new TypeError('agent maxTokens must be a positive safe integer')
  185. }
  186. }
  187. /** Prepared-but-unpublished agent resources sharing one memoized teardown. */
  188. interface PreparedAgent {
  189. agent: ReactLoopAgent
  190. /** Aborts when the factory unloads, the caller cancels, or teardown begins — ends any setup await. */
  191. signal: AbortSignal
  192. /** Enter registries, announce, notify session-start, and start the machine. */
  193. publish(source: SessionStartSource): AgentHandle
  194. /** Reverse teardown: stop the machine, unregister, unwind the scope. Memoized. */
  195. dispose(): Promise<void>
  196. }
  197. declare module '@deepseek-ai/cordis' {
  198. interface Context {
  199. agentLoop: AgentLoop
  200. /**
  201. * Launcher-owned exact session identities for configured agents, keyed by
  202. * the agent's config `id` and set with `ctx.provide()` before any Loader
  203. * entry mounts (see {@link CONFIGURED_AGENT_IDENTITIES_KEY}). A launcher
  204. * owns identity because only it knows whether the session already exists,
  205. * while the `cordis.yml` row keeps the model route as ordinary patchable
  206. * config. An entry with no matching key keeps its configured identity.
  207. */
  208. configuredAgentIdentities?: ConfiguredAgentIdentities
  209. }
  210. interface Events {
  211. /**
  212. * A declarative agent entry failed before it could publish a live agent.
  213. * Consumers that buffer work for the configured identity use this
  214. * transient signal to reject that work instead of waiting forever. Normal
  215. * factory teardown suppresses failures from the cancelled startup attempt.
  216. * @param payload.sessionId - exact shared agent/session identity that failed startup.
  217. * @param payload.error - persistence, setup, or publication failure.
  218. * @mode emit
  219. */
  220. 'agent-loop/config-start-failed'(payload: { sessionId: SessionId; error: unknown }): void
  221. }
  222. }
  223. export { DEFAULT_MAX_PARALLEL_TOOL_CALLS }
  224. /**
  225. * One launcher-selected session identity for a configured agent. `resume`
  226. * distinguishes rehydrating existing persisted history from creating the
  227. * session fresh under that exact id, which the two config keys express as
  228. * `resumeSessionId` and `sessionId`.
  229. */
  230. export interface LauncherAgentIdentity {
  231. /** Exact session id to create fresh or resume. */
  232. id: SessionId
  233. /** Resume existing persisted history instead of creating the session fresh. */
  234. resume: boolean
  235. }
  236. /** Launcher-selected identities keyed by the configured agent's `id`. */
  237. export interface ConfiguredAgentIdentities extends Readonly<Record<string, LauncherAgentIdentity>> {}
  238. /**
  239. * Context key a launcher sets before any Loader entry mounts
  240. * (`ctx.provide(CONFIGURED_AGENT_IDENTITIES_KEY, identities)`) to fix
  241. * configured agents' session identities without a config key, so an overlay
  242. * repointing the row's model route cannot drop them.
  243. */
  244. export const CONFIGURED_AGENT_IDENTITIES_KEY = 'configuredAgentIdentities'
  245. /**
  246. * Apply launcher-owned identities over the configured agents, replacing both
  247. * identity keys for every entry the launcher named so a config-supplied
  248. * identity can never survive alongside a launcher-supplied one.
  249. * @param agents - the configured agent entries.
  250. * @param identities - launcher identities keyed by configured agent `id`, or `undefined`.
  251. * @returns the entries with launcher-owned identities applied.
  252. */
  253. function applyLauncherIdentities(
  254. agents: Config['agents'],
  255. identities: ConfiguredAgentIdentities | undefined,
  256. ): Config['agents'] {
  257. if (identities === undefined) return agents
  258. return agents.map((agent) => {
  259. const identity = identities[agent.id]
  260. if (identity === undefined) return agent
  261. const { sessionId: _sessionId, resumeSessionId: _resumeSessionId, ...rest } = agent
  262. return identity.resume
  263. ? { ...rest, resumeSessionId: identity.id }
  264. : { ...rest, sessionId: identity.id }
  265. })
  266. }
  267. /** Settings namespace carrying the tool-call parallelism a user owns. */
  268. export const AGENT_LOOP_SETTINGS_NAMESPACE = 'agent-loop'
  269. /**
  270. * The agent-loop fields a user owns. Deliberately a strict subset of
  271. * {@link Config}: `agents` is a boot-time composition array consumed once when
  272. * the service starts, so a stored change could only look like it had an effect.
  273. */
  274. export interface AgentLoopSettings {
  275. /** Maximum parallel-safe calls in flight per agent step. */
  276. maxParallelToolCalls: number
  277. }
  278. /** Schema of the agent-loop settings section. */
  279. export const AGENT_LOOP_SETTINGS_SCHEMA: z<AgentLoopSettings> = z.object({
  280. maxParallelToolCalls: z.number().step(1).min(1).default(DEFAULT_MAX_PARALLEL_TOOL_CALLS),
  281. })
  282. /** Agent-loop plugin configuration. */
  283. export interface Config {
  284. /**
  285. * Maximum parallel-safe calls in flight per agent step. `1` is serial;
  286. * omission defaults to {@link DEFAULT_MAX_PARALLEL_TOOL_CALLS}.
  287. */
  288. maxParallelToolCalls?: number
  289. /** Agents created or resumed at plugin startup. */
  290. agents: (AgentOptions & {
  291. /** Stable config label used in logs and as the fresh combined-id prefix. */
  292. id: string
  293. /** Optional stable identity; remounts resume its materialized history, while first use creates it fresh. */
  294. sessionId?: SessionId
  295. /** Optional workspace for a fresh session. */
  296. cwd?: string
  297. /** Persisted session to resume instead of creating a fresh session. */
  298. resumeSessionId?: SessionId
  299. })[]
  300. }
  301. /** Agent-loop configuration after defaults and load-time validation. */
  302. type ResolvedConfig = Config & { maxParallelToolCalls: number }
  303. /** Reject self-contained identity conflicts before any configured agent starts. */
  304. function validateConfiguredAgents(agents: Config['agents']): void {
  305. const exactIdentities = new Map<SessionId, string>()
  306. for (const { id, sessionId, resumeSessionId } of agents) {
  307. const hasResumeId = resumeSessionId !== undefined && resumeSessionId !== ''
  308. if (sessionId !== undefined && hasResumeId) {
  309. throw new Error(`agent "${id}": sessionId and resumeSessionId are mutually exclusive`)
  310. }
  311. const exactIdentity = hasResumeId ? resumeSessionId : sessionId
  312. if (exactIdentity === undefined) continue
  313. const firstId = exactIdentities.get(exactIdentity)
  314. if (firstId !== undefined) {
  315. throw new Error(`agents "${firstId}" and "${id}" use duplicate exact session identity "${exactIdentity}"`)
  316. }
  317. exactIdentities.set(exactIdentity, id)
  318. }
  319. }
  320. /** Concrete agent factory and driver service. */
  321. export class AgentLoop extends Service implements AgentFactory {
  322. static inject = ['agents', 'sessions', 'llm', 'tools', 'systemPrompt', 'sessionProjections']
  323. /** Runtime schema for declarative agents. */
  324. static Config = z.object({
  325. maxParallelToolCalls: z.number().step(1).min(1).default(DEFAULT_MAX_PARALLEL_TOOL_CALLS),
  326. agents: z.array(z.object({
  327. id: z.string().required(),
  328. sessionId: z.string().min(1),
  329. provider: z.string(),
  330. model: z.string(),
  331. reasoningEffort: z.string().min(1) as z<ReturnType<typeof ReasoningEffortId>>,
  332. maxTokens: z.number().step(1).min(1).max(Number.MAX_SAFE_INTEGER),
  333. cwd: z.string(),
  334. resumeSessionId: z.string(),
  335. })).default([]),
  336. }) as z<Config>
  337. /** Validated configuration owned by the agent-loop service. */
  338. readonly config: ResolvedConfig
  339. private readonly ownership: FactoryOwnership
  340. /** Plain holder prevents Cordis from re-tracing the factory's dependency context through a caller shadow. */
  341. private readonly runtime: { ctx: Context }
  342. constructor(ctx: Context, config: Config) {
  343. super(ctx, 'agentLoop')
  344. const entry: AgentLoopSettings = {
  345. maxParallelToolCalls: resolveMaxParallelToolCalls(config.maxParallelToolCalls),
  346. }
  347. let source: () => AgentLoopSettings = () => entry
  348. this.config = {
  349. ...config,
  350. agents: applyLauncherIdentities(config.agents, ctx.get(CONFIGURED_AGENT_IDENTITIES_KEY)),
  351. // Read through on every scheduler decision: `tool-calls.ts` destructures
  352. // this at the start of each group, so a committed change caps the next
  353. // group without disturbing the one in flight.
  354. get maxParallelToolCalls() {
  355. return source().maxParallelToolCalls
  356. },
  357. }
  358. ctx.inject(['settings'], (settingsCtx) => {
  359. settingsCtx.settings.installSection(ctx, AGENT_LOOP_SETTINGS_NAMESPACE, AGENT_LOOP_SETTINGS_SCHEMA, entry, {
  360. // The schema admits any integer above zero; `resolveMaxParallelToolCalls`
  361. // owns the whole rule, so refusing here keeps the running scheduler on
  362. // its last good cap instead of failing at the next tool group.
  363. validate: value => void resolveMaxParallelToolCalls(value.maxParallelToolCalls),
  364. setSource: (current) => {
  365. source = current
  366. },
  367. // Nothing is derived from the cap: the getter above is the only reader.
  368. onChange: () => {},
  369. })
  370. })
  371. validateConfiguredAgents(this.config.agents)
  372. // Register only after every config validation above has passed, so a
  373. // rejected constructor leaves no projection unit behind.
  374. ctx.sessionProjections.register(turnBoundaryProjectionDefinition)
  375. this.ownership = new FactoryOwnership(ctx.fiber)
  376. this.runtime = { ctx }
  377. ctx.effect(() => () => this.ownership.dispose(), 'agentLoop.transactions()')
  378. ctx.effect(() => ctx.agents.setFactory(this), 'agentLoop.setFactory()')
  379. ctx.systemPrompt.variable('provider', context => context.agent?.options.provider)
  380. ctx.systemPrompt.variable('model', context => context.agent?.options.model)
  381. ctx.systemPrompt.variable('cwd', context => context.agent?.session.header.cwd)
  382. for (const { id, sessionId, cwd, resumeSessionId, ...options } of this.config.agents) {
  383. const meta = cwd === undefined ? {} : { cwd }
  384. if (resumeSessionId === undefined || resumeSessionId === '') {
  385. const configuredId = sessionId ?? brandString<SessionId>(`${id}-session-${randomUUID()}`)
  386. const persistence = sessionId === undefined ? undefined : ctx.get('sessionPersistence')
  387. if (persistence === undefined) {
  388. this.create(configuredId, options, meta)
  389. } else {
  390. const startup = this.restoreOrCreateConfigured(ctx, persistence, configuredId, options, meta).catch((error: unknown) => {
  391. this.reportConfiguredStartupFailure(id, 'restore', configuredId, error)
  392. })
  393. this.ownership.trackStartup(startup)
  394. }
  395. continue
  396. }
  397. ctx.effect(() => {
  398. const fiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {
  399. void this.resumeWith(ctx, childCtx.sessionPersistence, {
  400. resumeSessionId,
  401. agentOptions: options,
  402. }).catch((error: unknown) => {
  403. this.reportConfiguredStartupFailure(id, 'resume', resumeSessionId, error)
  404. })
  405. })
  406. return fiber.dispose
  407. }, `agentLoop.resume(${id})`)
  408. }
  409. }
  410. /** Report a contained declarative-start failure to identity-bound consumers. */
  411. private reportConfiguredStartupFailure(
  412. configId: string,
  413. action: 'restore' | 'resume',
  414. sessionId: SessionId,
  415. error: unknown,
  416. ): void {
  417. if (!this.ownership.isActive()) return
  418. this.ctx.logger.warn(`agent "${configId}": config-driven ${action} of "${sessionId}" failed: ${errorChain(error)}`)
  419. const args: unknown[] = ['agent-loop/config-start-failed', { sessionId, error }]
  420. for (const callback of this.ctx.events.dispatch('emit', args)) {
  421. try {
  422. const returned: unknown = callback(...args)
  423. void Promise.resolve(returned).catch((listenerError: unknown) => {
  424. this.ctx.logger.warn(`agent "${configId}": config-start-failed listener rejected: ${errorChain(listenerError)}`)
  425. })
  426. } catch (listenerError: unknown) {
  427. this.ctx.logger.warn(`agent "${configId}": config-start-failed listener threw: ${errorChain(listenerError)}`)
  428. }
  429. }
  430. }
  431. /** Restore a materialized exact config identity on remount, or create it on first use. */
  432. private async restoreOrCreateConfigured(
  433. ownerCtx: Context,
  434. persistence: SessionPersistence,
  435. sessionId: SessionId,
  436. agentOptions: AgentOptions,
  437. meta: Pick<SessionHeader, 'cwd'>,
  438. ): Promise<void> {
  439. await this.waitForDrainingConfiguredIdentity(ownerCtx, sessionId)
  440. if (!this.ownership.isActive()) return
  441. try {
  442. await this.resumeWith(ownerCtx, persistence, { resumeSessionId: sessionId, agentOptions })
  443. return
  444. } catch (error: unknown) {
  445. if (!this.ownership.isActive()) return
  446. // A load is the per-id serialization barrier for eager write-behind and
  447. // lifecycle retirement. Only a genuinely absent artifact falls back to
  448. // first creation; corruption and backend failures stay loud.
  449. const exists = (await persistence.list()).some(header => header.id === sessionId)
  450. if (exists) throw error
  451. }
  452. this.create(sessionId, agentOptions, meta)
  453. }
  454. /** Wait for a draining same-id lifecycle to finish registry teardown. */
  455. private async waitForDrainingConfiguredIdentity(ownerCtx: Context, sessionId: SessionId): Promise<void> {
  456. // Only an id still occupying a registry needs waiting for; a live healthy
  457. // occupant is a collision the create/resume below will surface itself.
  458. if (ownerCtx.agents.get(sessionId) === undefined && ownerCtx.sessions.get(sessionId) === undefined) return
  459. const released = Promise.withResolvers<void>()
  460. const checkReleased = (): void => {
  461. if (ownerCtx.agents.get(sessionId) === undefined && ownerCtx.sessions.get(sessionId) === undefined) {
  462. released.resolve()
  463. }
  464. }
  465. const disposeAgentListener = ownerCtx.on('agent/disposed', () => { checkReleased() })
  466. const disposeSessionListener = ownerCtx.on('session/disposed', checkReleased)
  467. try {
  468. checkReleased()
  469. await this.ownership.waitWhileActive(released.promise)
  470. } finally {
  471. disposeAgentListener()
  472. disposeSessionListener()
  473. }
  474. }
  475. /**
  476. * Construct the driver, scope, and one memoized reverse teardown for a new
  477. * agent. The teardown is registered with the factory and the owner fiber
  478. * BEFORE publication, so a mid-setup unload rolls everything back; `signal`
  479. * fuses caller cancellation with lifecycle teardown for setup awaits.
  480. */
  481. private prepare(ownerCtx: Context, id: SessionId, options: AgentOptions, session: Session, callerSignal?: AbortSignal): PreparedAgent {
  482. assertAgentOptions(options)
  483. ownerCtx.fiber.assertActive()
  484. // Every caller reaches prepare() synchronously from a service method
  485. // whose Cordis dispatch already requires the live factory fiber, or
  486. // re-checks ownership itself after its awaits (resume's load barrier).
  487. /* v8 ignore next -- unreachable backstop, see above */
  488. if (!this.ownership.isActive()) throw new Error('agent loop is not active')
  489. if (callerSignal?.aborted) {
  490. throw callerSignal.reason instanceof Error
  491. ? callerSignal.reason
  492. : new Error(`agent "${id}" creation aborted`, { cause: callerSignal.reason })
  493. }
  494. const loopCtx = this.runtime.ctx
  495. // Deactivation fuses three owners, each with its own reason: the caller's
  496. // cancellation signal, the owner fiber's unload, and factory teardown.
  497. // It is registered BEFORE any resource exists, over mutable slots, so an
  498. // unload arriving while the scope is still minting finds a working
  499. // disposer instead of a leak.
  500. const abort = new AbortController()
  501. const onCallerAbort = (): void => {
  502. abort.abort(callerSignal?.reason instanceof Error
  503. ? callerSignal.reason
  504. : new Error(`agent "${id}" creation aborted`, { cause: callerSignal?.reason }))
  505. }
  506. const onFactoryTeardown = (): void => { abort.abort(this.ownership.signal.reason) }
  507. callerSignal?.addEventListener('abort', onCallerAbort, { once: true })
  508. this.ownership.signal.addEventListener('abort', onFactoryTeardown, { once: true })
  509. let machine: ReactLoopAgent | undefined
  510. let detachSession: (() => void) | undefined
  511. let detachAgent: (() => void) | undefined
  512. let disposing: Promise<void> | undefined
  513. const machineReady = Promise.withResolvers<void>()
  514. // Reverse teardown, memoized so every racing owner awaits one quiescence:
  515. // stop the machine, leave the registries, unwind the scope, release
  516. // bookkeeping.
  517. const dispose = (ownerTriggered = false): Promise<void> => (disposing ??= (async () => {
  518. abort.abort(new Error(`agent "${id}" lifecycle disposed`))
  519. callerSignal?.removeEventListener('abort', onCallerAbort)
  520. this.ownership.signal.removeEventListener('abort', onFactoryTeardown)
  521. try {
  522. // Disposal IS a disposed-cause cancel followed by quiescence. New work
  523. // sent after this point is the sender's bug — the registries are about
  524. // to drop the agent, so nothing should still hold it.
  525. if (machine === undefined) await machineReady.promise
  526. if (machine !== undefined) {
  527. machine.cancel({ kind: 'disposed' })
  528. await machine.whenIdle()
  529. await machine.scope.dispose()
  530. }
  531. } finally {
  532. try {
  533. detachAgent?.()
  534. detachSession?.()
  535. } finally {
  536. untrack()
  537. if (!ownerTriggered) await unfollowOwner()
  538. }
  539. }
  540. })())
  541. const untrack = this.ownership.track(dispose)
  542. let unfollowOwner: () => Promise<void> | void
  543. try {
  544. unfollowOwner = ownerCtx.effect(() => () => {
  545. // Owner disposal owns the same quiescence boundary. Its teardown skips
  546. // unregistering this already-running owner effect from inside itself.
  547. if (disposing !== undefined) return
  548. abort.abort(new Error(`agent "${id}" setup aborted: owner disposed during setup`))
  549. return dispose(true)
  550. }, `agentLoop.lifecycle(${id})`)
  551. /* v8 ignore start -- ctx.effect throws only on an inactive fiber, which assertActive() above already rejected */
  552. } catch (error: unknown) {
  553. untrack()
  554. callerSignal?.removeEventListener('abort', onCallerAbort)
  555. this.ownership.signal.removeEventListener('abort', onFactoryTeardown)
  556. throw error
  557. }
  558. /* v8 ignore stop */
  559. const assertLive = (): void => {
  560. if (!abort.signal.aborted) return
  561. // Every fused abort source carries an Error reason: onCallerAbort and
  562. // raceAbort wrap non-Error caller reasons, and the factory/lifecycle
  563. // owners abort with constructed Errors.
  564. /* v8 ignore next -- unreachable String() arm, see above */
  565. throw abort.signal.reason instanceof Error ? abort.signal.reason : new Error(String(abort.signal.reason))
  566. }
  567. try {
  568. const agent = machine = new ReactLoopAgent(loopCtx, id, options, session)
  569. machineReady.resolve()
  570. assertLive()
  571. return {
  572. agent,
  573. signal: abort.signal,
  574. publish: (source) => {
  575. assertLive()
  576. detachSession = agent.ctx.sessions.enter(session)
  577. detachAgent = loopCtx.agents.enter(agent, ownerCtx.agent)
  578. agent.ctx.sessions.announce(session)
  579. assertLive()
  580. loopCtx.agents.announce(agent)
  581. assertLive()
  582. // A synchronous announce/session-start listener may have started
  583. // teardown; the machine is already live (delivery works from the
  584. // session-start extension point), so only the liveness recheck is owed.
  585. emitAgentEvent(loopCtx, agent, 'agent/session-start', { source })
  586. assertLive()
  587. return { agent, dispose }
  588. },
  589. dispose,
  590. }
  591. } catch (error: unknown) {
  592. machineReady.resolve()
  593. void dispose()
  594. throw error
  595. }
  596. }
  597. /**
  598. * Create an agent and session under one caller-supplied identity, owned by
  599. * the accessing fiber. Constructor-driven config calls mint a fresh combined
  600. * id before entering this boundary.
  601. * @param id - shared agent/session identity.
  602. * @param options - concrete loop options.
  603. * @param meta - optional fresh-session workspace metadata.
  604. * @returns the published running agent.
  605. */
  606. create(id: SessionId, options: AgentOptions = {}, meta: Pick<SessionHeader, 'cwd'> = {}): Agent {
  607. using preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(id, { meta }))
  608. const prepared = this.prepare(this.ctx, id, options, preparation.session)
  609. try {
  610. return prepared.publish('startup').agent
  611. } catch (error: unknown) {
  612. void prepared.dispose()
  613. throw error
  614. }
  615. }
  616. /**
  617. * Create an owned agent on a caller-supplied session id.
  618. * @param ownerCtx - caller context that structurally owns the lifecycle.
  619. * @param options - identities, session seed/metadata, loop options, setup, and cancellation.
  620. * @returns the published handle.
  621. */
  622. async createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> {
  623. const preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(options.sessionId, {
  624. ...options.seed === undefined ? {} : { seed: options.seed },
  625. ...options.meta === undefined ? {} : { meta: options.meta },
  626. }))
  627. const published = this.setupAndPublish(
  628. ownerCtx,
  629. options.sessionId,
  630. preparation,
  631. options.agentOptions ?? {},
  632. options.setup,
  633. options.signal,
  634. 'startup',
  635. )
  636. this.ownership.trackWrapper(published)
  637. return published
  638. }
  639. /** Prepare one Agent around an acquired Session, run setup, and publish it. */
  640. private async setupAndPublish(
  641. ownerCtx: Context,
  642. id: SessionId,
  643. preparation: SessionPreparation,
  644. agentOptions: AgentOptions,
  645. setup: AgentSetup | undefined,
  646. signal: AbortSignal | undefined,
  647. source: SessionStartSource,
  648. ): Promise<AgentHandle> {
  649. using ownedPreparation = preparation
  650. const session = ownedPreparation.session
  651. const prepared = this.prepare(ownerCtx, id, agentOptions, session, signal)
  652. try {
  653. const setupCommit = await raceAbort(setup?.(prepared.agent.ctx), prepared.signal, id)
  654. setupCommit?.commit()
  655. return prepared.publish(source)
  656. } catch (error: unknown) {
  657. await prepared.dispose()
  658. throw error
  659. }
  660. }
  661. /**
  662. * Resume an owned agent from the configured persistence service.
  663. * @param ownerCtx - caller context that owns load, setup, and the live lifecycle.
  664. * @param options - persisted identity, loop options, setup, and cancellation.
  665. * @returns the published handle.
  666. */
  667. async resume(ownerCtx: Context, options: ResumeAgentOptions): Promise<AgentHandle> {
  668. const persistence = this.runtime.ctx.get('sessionPersistence')
  669. if (persistence === undefined) {
  670. throw new Error('cannot resume: session persistence is not configured (load a dsh-session-persistence backend)')
  671. }
  672. return this.resumeWith(ownerCtx, persistence, options)
  673. }
  674. /** Resume through an explicit persistence handle used by the deferred config path. */
  675. private resumeWith(
  676. ownerCtx: Context,
  677. persistence: SessionPersistence,
  678. options: ResumeAgentOptions,
  679. ): Promise<AgentHandle> {
  680. const id = options.resumeSessionId
  681. const published = (async () => {
  682. // The load may outlive its owner: race it against caller cancellation,
  683. // owner-fiber unload, and factory teardown so a never-settling backend
  684. // cannot pin the identity.
  685. const ownerAbort = new AbortController()
  686. const unfollowOwner = ownerCtx.effect(() => () => {
  687. ownerAbort.abort(new Error(`agent "${id}" setup aborted: owner disposed during setup`))
  688. }, `agentLoop.resume-load(${id})`)
  689. const fused = AbortSignal.any([
  690. ...options.signal === undefined ? [] : [options.signal],
  691. ownerAbort.signal,
  692. this.ownership.signal,
  693. ])
  694. let preparation: SessionPreparation | undefined
  695. try {
  696. try {
  697. preparation = await raceAbortCall(
  698. () => persistence.prepare(id, fused),
  699. fused,
  700. id,
  701. (abandoned) => { abandoned[Symbol.dispose]() },
  702. )
  703. } finally {
  704. await unfollowOwner()
  705. }
  706. ownerCtx.fiber.assertActive()
  707. if (!this.ownership.isActive()) throw new Error('agent loop is not active')
  708. return await this.setupAndPublish(
  709. ownerCtx,
  710. id,
  711. preparation,
  712. options.agentOptions ?? {},
  713. options.setup,
  714. options.signal,
  715. 'resume',
  716. )
  717. } finally {
  718. preparation?.[Symbol.dispose]()
  719. }
  720. })()
  721. this.ownership.trackWrapper(published)
  722. return published
  723. }
  724. }
  725. export default AgentLoop