scope-lifecycle.spec.ts 50 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210
  1. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  2. import { describe, expect, it } from 'vitest'
  3. import { Context, symbols, type EffectMeta, type Fiber } from '@deepseek-ai/cordis'
  4. import LlmRuntime from '@deepseek-ai/dsh-llm'
  5. import SessionStore, { SessionId, type SessionEvent } from '@deepseek-ai/dsh-session'
  6. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  7. import ToolRuntime, { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
  8. import AgentRegistry, { agentEvents, assembleContextFor } from '@deepseek-ai/dsh-agent'
  9. import type { Agent } from '@deepseek-ai/dsh-agent'
  10. import { scopeOf } from '@deepseek-ai/dsh-scope'
  11. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  12. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  13. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  14. import { MockAdapter, textResponse } from './mock-adapter.ts'
  15. async function harnessWithLoop(adapter: MockAdapter = new MockAdapter([textResponse('ok')])): Promise<{ ctx: Context; loopFiber: Fiber }> {
  16. const ctx = new Context()
  17. await ctx.plugin(LlmRuntime)
  18. await ctx.plugin(SessionStore)
  19. await ctx.plugin(SessionProjectionRegistry)
  20. await ctx.plugin(SystemPrompt, { personaPrefix: 'You are the deployment.' })
  21. await ctx.plugin(ToolRuntime)
  22. await ctx.plugin(AgentRegistry)
  23. const loopFiber = await ctx.plugin(AgentLoop, { agents: [] })
  24. ctx.llm.registerAdapter(['mock'], adapter)
  25. return { ctx, loopFiber }
  26. }
  27. async function harness(adapter: MockAdapter = new MockAdapter([textResponse('ok')])): Promise<Context> {
  28. return (await harnessWithLoop(adapter)).ctx
  29. }
  30. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  31. return new Promise((resolve) => {
  32. const dispose = ctx.on('agent/status', ({ agent: subject, status }) => {
  33. if (subject === agent && status === 'idle') {
  34. dispose()
  35. resolve()
  36. }
  37. })
  38. })
  39. }
  40. const text = (t: string): ContentBlock[] => [{ type: 'text', text: t }]
  41. /** Throw an arbitrary callback value to exercise the public unknown-error boundary. */
  42. function throwUnknown(value: unknown): never {
  43. throw value
  44. }
  45. /** Invoke the exact lifecycle effect to exercise same-stack reentrant teardown. */
  46. function disposeCurrentLifecycle(ownerCtx: Context): void {
  47. const lifecycle = [...ownerCtx.fiber._disposables]
  48. .find((dispose) => {
  49. const effect = (dispose as typeof dispose & { [symbols.effect]?: EffectMeta })[symbols.effect]
  50. return effect?.label.startsWith('agentLoop.lifecycle(') === true
  51. })
  52. if (lifecycle === undefined) throw new Error('agent lifecycle effect not found')
  53. void lifecycle()
  54. }
  55. describe('agent scope lifecycle', () => {
  56. it('awaits serial creation before caller readiness', async () => {
  57. const ctx = await harness()
  58. const entered = Promise.withResolvers<undefined>()
  59. const release = Promise.withResolvers<undefined>()
  60. const order: string[] = []
  61. ctx.on('agent/created', async () => {
  62. order.push('created:start')
  63. entered.resolve(undefined)
  64. await release.promise
  65. order.push('created:end')
  66. })
  67. ctx.on('agent/created', () => { order.push('created:second') })
  68. const creating = ctx.agentLoop.create(SessionId('serial-start')).then((agent) => {
  69. order.push('ready')
  70. return agent
  71. })
  72. try {
  73. await entered.promise
  74. expect(order).toEqual(['created:start'])
  75. release.resolve(undefined)
  76. await creating
  77. expect(order).toEqual(['created:start', 'created:end', 'created:second', 'ready'])
  78. } finally {
  79. release.resolve(undefined)
  80. await ctx.fiber.dispose()
  81. }
  82. })
  83. it('retains both entries and scoped effects until an asynchronous creation listener settles on owner unload', async () => {
  84. const ctx = await harness()
  85. const entered = Promise.withResolvers<undefined>()
  86. const release = Promise.withResolvers<undefined>()
  87. let ownerCtx!: Context
  88. const owner = await ctx.plugin(Object.assign((inner: Context) => { ownerCtx = inner }, { inject: ['agents'] }))
  89. const order: string[] = []
  90. const id = SessionId('serial-owner-dispose')
  91. ctx.on('agent/created', async ({ agent }) => {
  92. entered.resolve(undefined)
  93. await release.promise
  94. expect(ctx.agents.get(id)).toBe(agent)
  95. expect(ctx.sessions.get(id)).toBe(agent.session)
  96. expect(order).toEqual([])
  97. order.push('created:end')
  98. })
  99. ctx.on('agent/disposed', () => { order.push('agent-disposed') })
  100. ctx.on('session/disposed', () => { order.push('session-disposed') })
  101. const creating = ownerCtx.agents.create({
  102. sessionId: id,
  103. setup(agentCtx) { agentCtx.effect(() => () => { order.push('scope-disposed') }) },
  104. })
  105. const rejected = expect(creating).rejects.toThrow('owner disposed during setup')
  106. try {
  107. await entered.promise
  108. const disposing = owner.dispose()
  109. expect(ctx.agents.get(id)).toBeDefined()
  110. release.resolve(undefined)
  111. await rejected
  112. await disposing
  113. expect(order).toEqual(['created:end', 'scope-disposed', 'agent-disposed', 'session-disposed'])
  114. expect(ctx.agents.get(id)).toBeUndefined()
  115. expect(ctx.sessions.get(id)).toBeUndefined()
  116. } finally {
  117. release.resolve(undefined)
  118. await ctx.fiber.dispose()
  119. }
  120. })
  121. it('rolls back both registries and the scope when asynchronous creation fails', async () => {
  122. const adapter = new MockAdapter([textResponse('unexpected')])
  123. const ctx = await harness(adapter)
  124. const reason = new Error('async creation veto')
  125. const order: string[] = []
  126. ctx.on('agent/created', async ({ agent }) => {
  127. agent.followup(createUserMessage({ content: text('queued during creation'), source: { kind: 'user' } }))
  128. await Promise.resolve()
  129. throw reason
  130. })
  131. ctx.on('agent/created', () => { order.push('later-listener') })
  132. ctx.on('agent/disposed', () => { order.push('agent-disposed') })
  133. ctx.on('session/disposed', () => { order.push('session-disposed') })
  134. await expect(ctx.agents.create({
  135. sessionId: SessionId('async-veto'),
  136. setup(agentCtx) { agentCtx.effect(() => () => { order.push('scope-disposed') }) },
  137. })).rejects.toBe(reason)
  138. expect(order).toEqual(['scope-disposed', 'agent-disposed', 'session-disposed'])
  139. expect(ctx.agents.list()).toEqual([])
  140. expect(ctx.sessions.list()).toEqual([])
  141. expect(adapter.requests).toHaveLength(0)
  142. await ctx.fiber.dispose()
  143. })
  144. it('rejects an already-aborted creation signal before publishing either object', async () => {
  145. const ctx = await harness()
  146. const reason = new Error('cancelled before creation')
  147. const controller = new AbortController()
  148. controller.abort(reason)
  149. await expect(ctx.agents.create({
  150. sessionId: SessionId('pre-aborted-s'),
  151. signal: controller.signal,
  152. })).rejects.toBe(reason)
  153. expect(ctx.agents.get(SessionId('pre-aborted-s'))).toBeUndefined()
  154. expect(ctx.sessions.get(SessionId('pre-aborted-s'))).toBeUndefined()
  155. const valueController = new AbortController()
  156. valueController.abort('plain cancellation reason')
  157. await expect(ctx.agents.create({
  158. sessionId: SessionId('pre-aborted-value-s'),
  159. signal: valueController.signal,
  160. })).rejects.toMatchObject({
  161. message: 'agent "pre-aborted-value-s" creation aborted',
  162. cause: 'plain cancellation reason',
  163. })
  164. expect(ctx.agents.get(SessionId('pre-aborted-value-s'))).toBeUndefined()
  165. expect(ctx.sessions.get(SessionId('pre-aborted-value-s'))).toBeUndefined()
  166. await ctx.fiber.dispose()
  167. })
  168. it('joins cleanup when an abort lands reentrantly during scope preparation', async () => {
  169. const ctx = await harness()
  170. const reason = new Error('cancelled while preparing')
  171. const controller = new AbortController()
  172. let aborted = false
  173. ctx.on('internal/plugin', (fiber) => {
  174. if (aborted || fiber.name !== 'scope') return
  175. aborted = true
  176. controller.abort(reason)
  177. })
  178. await expect(ctx.agents.create({
  179. sessionId: SessionId('prepare-abort-s'),
  180. signal: controller.signal,
  181. })).rejects.toBe(reason)
  182. expect(ctx.agents.get(SessionId('prepare-abort-s'))).toBeUndefined()
  183. expect(ctx.sessions.get(SessionId('prepare-abort-s'))).toBeUndefined()
  184. await ctx.fiber.dispose()
  185. })
  186. it('normalizes non-Error create failures for rollback while rethrowing the original value', async () => {
  187. const ctx = await harness()
  188. let thrown: unknown
  189. ctx.on('session/created', () => {
  190. if (thrown === undefined) return
  191. const value = thrown
  192. thrown = undefined
  193. throwUnknown(value)
  194. })
  195. const createFailure = { source: 'create' }
  196. thrown = createFailure
  197. let createCaught: unknown
  198. try {
  199. await ctx.agentLoop.create(SessionId('unknown-create'))
  200. } catch (error: unknown) {
  201. createCaught = error
  202. }
  203. expect(createCaught).toBe(createFailure)
  204. const ownedFailure = { source: 'createAgent' }
  205. thrown = ownedFailure
  206. await expect(ctx.agents.create({
  207. sessionId: SessionId('unknown-owned-create-s'),
  208. })).rejects.toBe(ownedFailure)
  209. expect(ctx.agents.get(SessionId('unknown-create'))).toBeUndefined()
  210. expect(ctx.agents.get(SessionId('unknown-owned-create-s'))).toBeUndefined()
  211. await ctx.fiber.dispose()
  212. })
  213. it('tags agent.ctx with the Agent scope key', async () => {
  214. const ctx = await harness()
  215. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  216. expect(scopeOf(agent.ctx)).toBe(agent)
  217. await agent.whenIdle()
  218. })
  219. it('records an explicitly owned Agent as a non-root runtime child', async () => {
  220. const ctx = await harness()
  221. const root = await ctx.agents.create({
  222. sessionId: SessionId('runtime-root'),
  223. agentOptions: { model: 'mock' },
  224. })
  225. const child = await root.agent.ctx.agents.create({
  226. sessionId: SessionId('runtime-child'),
  227. agentOptions: { model: 'mock' },
  228. parentAgent: root.agent,
  229. })
  230. expect(ctx.agents.list()).toEqual([root.agent, child.agent])
  231. expect(ctx.agents.roots()).toEqual([root.agent])
  232. await child.dispose()
  233. await root.dispose()
  234. })
  235. it('scoped registrations live in the agent world and die with the agent', async () => {
  236. const ctx = await harness()
  237. const handle = await ctx.agents.create({ sessionId: SessionId('s1'), agentOptions: { provider: 'mock', model: 'mock' } })
  238. const { agent } = handle
  239. agent.ctx.systemPrompt.section({ name: 'deployment:persona-prefix', order: 0, text: 'You run tests.' })
  240. agent.ctx.tools.register(defineContentToolFixture({
  241. name: 'mine', description: 'scoped', parameters: {},
  242. execute: () => Promise.resolve(text('ran')),
  243. }))
  244. const scopedAssembly = await ctx.systemPrompt.assemble(assembleContextFor(agent))
  245. expect(scopedAssembly.sections.find(s => s.name === 'deployment:persona-prefix')?.text).toBe('You run tests.')
  246. expect(scopedAssembly.tools.map(t => t.name)).toContain('mine')
  247. // Other assemblies are untouched.
  248. const globalAssembly = await ctx.systemPrompt.assemble()
  249. expect(globalAssembly.sections.find(s => s.name === 'deployment:persona-prefix')?.text).toBe('You are the deployment.')
  250. expect(globalAssembly.tools.map(t => t.name)).not.toContain('mine')
  251. await handle.dispose()
  252. // The scoped world unwound with the agent: nothing leaked into the registries.
  253. expect(ctx.tools.get('mine', agent)).toBeUndefined()
  254. const after = await ctx.systemPrompt.assemble(assembleContextFor(agent))
  255. expect(after.sections.find(s => s.name === 'deployment:persona-prefix')?.text).toBe('You are the deployment.')
  256. })
  257. it('keeps the inbox projection until the last owning agent fiber unloads', async () => {
  258. const ctx = await harness()
  259. let first!: Awaited<ReturnType<typeof ctx.agents.create>>
  260. let second!: Awaited<ReturnType<typeof ctx.agents.create>>
  261. const firstOwner = await ctx.plugin(Object.assign(async (inner: Context) => {
  262. first = await inner.agents.create({
  263. sessionId: SessionId('projection-owner-first'),
  264. agentOptions: { provider: 'mock', model: 'mock' },
  265. })
  266. }, { inject: ['agents'] }))
  267. const secondOwner = await ctx.plugin(Object.assign(async (inner: Context) => {
  268. second = await inner.agents.create({
  269. sessionId: SessionId('projection-owner-second'),
  270. agentOptions: { provider: 'mock', model: 'mock' },
  271. })
  272. }, { inject: ['agents'] }))
  273. expect(ctx.sessionProjections.stateOf(first.agent.session, 'inbox')).toBeDefined()
  274. await firstOwner.dispose()
  275. expect(ctx.sessionProjections.stateOf(second.agent.session, 'inbox')).toBeDefined()
  276. await secondOwner.dispose()
  277. expect(ctx.sessionProjections.stateOf(second.agent.session, 'inbox')).toBeUndefined()
  278. await Promise.all([first.dispose(), second.dispose()])
  279. await ctx.fiber.dispose()
  280. })
  281. it('agent.ctx listeners hear only their own agent (scoped dispatch end to end)', async () => {
  282. const ctx = await harness(new MockAdapter([textResponse('one'), textResponse('two')]))
  283. const a = await ctx.agentLoop.create(SessionId('a'), { provider: 'mock', model: 'mock' })
  284. const b = await ctx.agentLoop.create(SessionId('b'), { provider: 'mock', model: 'mock' })
  285. const heard: string[] = []
  286. a.ctx.on('agent/status', ({ agent: subject, status }) => void heard.push(`a-sees:${subject.id}:${status}`))
  287. a.ctx.on('session/event', (_s, event) => {
  288. if (event.type === 'user/message') heard.push('a-sees:user-message')
  289. })
  290. b.followup(createUserMessage({ content: text('for b'), source: { kind: 'user' } }))
  291. await waitForIdle(ctx, b)
  292. expect(heard).toEqual([]) // nothing of b's leaked into a's scope
  293. a.followup(createUserMessage({ content: text('for a'), source: { kind: 'user' } }))
  294. await waitForIdle(ctx, a)
  295. expect(heard).toContain('a-sees:a:running')
  296. expect(heard).toContain('a-sees:user-message')
  297. })
  298. it('runs setup in the guaranteed slot: scoped world complete before agent/created and the first assembly', async () => {
  299. const ctx = await harness()
  300. const order: string[] = []
  301. ctx.on('agent/created', ({ agent }) => {
  302. order.push('agent/created')
  303. // The scoped section is already registered by the time agent/created fires.
  304. void ctx.systemPrompt.assemble(assembleContextFor(agent)).then((assembly) => {
  305. order.push(`persona:${assembly.sections.find(s => s.name === 'deployment:persona-prefix')?.text}`)
  306. })
  307. })
  308. const handle = await ctx.agents.create({
  309. sessionId: SessionId('child-s'),
  310. agentOptions: { provider: 'mock', model: 'mock' },
  311. setup: async (agentCtx) => {
  312. order.push('setup')
  313. await Promise.resolve()
  314. agentCtx.systemPrompt.section({ name: 'deployment:persona-prefix', order: 0, text: 'You are the child.' })
  315. },
  316. })
  317. expect(order).toEqual(['setup', 'agent/created', 'persona:You are the child.'])
  318. await handle.dispose()
  319. })
  320. it('keeps both objects unpublished until async setup completes, then announces in order', async () => {
  321. const ctx = await harness()
  322. const gate = Promise.withResolvers<undefined>()
  323. const setupStarted = Promise.withResolvers<undefined>()
  324. const order: string[] = []
  325. ctx.on('session/created', (session) => {
  326. expect(ctx.sessions.get(session.id)).toBe(session)
  327. expect(ctx.agents.get(session.id)?.session).toBe(session)
  328. order.push('session/created')
  329. })
  330. ctx.on('agent/created', () => void order.push('agent/created'))
  331. const acceptedOptions = { provider: 'mock', model: 'mock' }
  332. const creating = ctx.agents.create({
  333. sessionId: SessionId('atomic'),
  334. agentOptions: acceptedOptions,
  335. setup: async (agentCtx, agent) => {
  336. expect(agent.id).toBe(SessionId('atomic'))
  337. agentCtx.on('session/created', () => void order.push('setup-listener:session/created'))
  338. agentCtx.on('agent/created', () => void order.push('setup-listener:agent/created'))
  339. order.push('setup:start')
  340. setupStarted.resolve(undefined)
  341. await gate.promise
  342. order.push('setup:end')
  343. return {
  344. commit: () => {
  345. expect(ctx.agents.get(SessionId('atomic'))).toBeUndefined()
  346. expect(ctx.sessions.get(SessionId('atomic'))).toBeUndefined()
  347. order.push('setup:commit')
  348. },
  349. }
  350. },
  351. })
  352. await setupStarted.promise
  353. expect(ctx.agents.get(SessionId('atomic'))).toBeUndefined()
  354. expect(ctx.sessions.get(SessionId('atomic'))).toBeUndefined()
  355. expect(order).toEqual(['setup:start'])
  356. gate.resolve(undefined)
  357. const handle = await creating
  358. expect(handle.agent.options).toBe(acceptedOptions)
  359. expect(order).toEqual([
  360. 'setup:start',
  361. 'setup:end',
  362. 'setup:commit',
  363. 'session/created',
  364. 'setup-listener:session/created',
  365. 'agent/created',
  366. 'setup-listener:agent/created',
  367. ])
  368. await handle.dispose()
  369. })
  370. it('lets the final enter arbitrate unsupported concurrent same-id creation and rolls the loser back', async () => {
  371. const ctx = await harness()
  372. const gate = Promise.withResolvers<undefined>()
  373. const bothStarted = Promise.withResolvers<undefined>()
  374. let started = 0
  375. const setup = async (): Promise<void> => {
  376. started += 1
  377. if (started === 2) bothStarted.resolve(undefined)
  378. await gate.promise
  379. }
  380. const sessionId = SessionId('concurrent-final-enter')
  381. const first = ctx.agents.create({
  382. sessionId,
  383. agentOptions: { provider: 'mock', model: 'mock' },
  384. setup,
  385. })
  386. const second = ctx.agents.create({
  387. sessionId,
  388. agentOptions: { provider: 'mock', model: 'mock' },
  389. setup,
  390. })
  391. await bothStarted.promise
  392. expect(ctx.agents.list()).toEqual([])
  393. expect(ctx.sessions.list()).toEqual([])
  394. gate.resolve(undefined)
  395. const outcomes = await Promise.allSettled([first, second])
  396. const fulfilled = outcomes.filter((outcome): outcome is PromiseFulfilledResult<Awaited<typeof first>> => outcome.status === 'fulfilled')
  397. const rejected = outcomes.filter((outcome): outcome is PromiseRejectedResult => outcome.status === 'rejected')
  398. expect(fulfilled).toHaveLength(1)
  399. expect(rejected).toHaveLength(1)
  400. expect(String(rejected[0]!.reason)).toMatch(/already exists/)
  401. expect(ctx.agents.list()).toEqual([fulfilled[0]!.value.agent])
  402. expect(ctx.sessions.list()).toEqual([fulfilled[0]!.value.agent.session])
  403. await fulfilled[0]!.value.dispose()
  404. expect(ctx.agents.list()).toEqual([])
  405. expect(ctx.sessions.list()).toEqual([])
  406. })
  407. it('uses signal only for creation: aborts pending setup but not a returned live handle', async () => {
  408. const ctx = await harness()
  409. const pendingController = new AbortController()
  410. const setupStarted = Promise.withResolvers<undefined>()
  411. const pending = ctx.agents.create({
  412. sessionId: SessionId('signal-pending-s'),
  413. agentOptions: { provider: 'mock', model: 'mock' },
  414. signal: pendingController.signal,
  415. setup: async () => {
  416. setupStarted.resolve(undefined)
  417. await new Promise<never>(() => {})
  418. },
  419. })
  420. await setupStarted.promise
  421. pendingController.abort(new Error('cancel pending creation'))
  422. await expect(pending).rejects.toThrow('cancel pending creation')
  423. expect(ctx.agents.get(SessionId('signal-pending-s'))).toBeUndefined()
  424. expect(ctx.sessions.get(SessionId('signal-pending-s'))).toBeUndefined()
  425. const liveController = new AbortController()
  426. const live = await ctx.agents.create({
  427. sessionId: SessionId('signal-live-s'),
  428. agentOptions: { provider: 'mock', model: 'mock' },
  429. signal: liveController.signal,
  430. })
  431. liveController.abort(new Error('too late'))
  432. await Promise.resolve()
  433. expect(ctx.agents.get(live.agent.id)).toBe(live.agent)
  434. expect(live.agent.status).toBe('idle')
  435. await live.dispose()
  436. })
  437. it('owner unload aborts a pending setup and publishes nothing', async () => {
  438. const ctx = await harness()
  439. const gate = Promise.withResolvers<undefined>()
  440. const setupStarted = Promise.withResolvers<undefined>()
  441. const published: string[] = []
  442. ctx.on('session/created', () => void published.push('session/created'))
  443. ctx.on('agent/created', () => void published.push('agent/created'))
  444. let creating!: ReturnType<typeof ctx.agents.create>
  445. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  446. creating = inner.agents.create({
  447. sessionId: SessionId('owner-race-s'),
  448. agentOptions: { provider: 'mock', model: 'mock' },
  449. setup: async () => {
  450. setupStarted.resolve(undefined)
  451. await gate.promise
  452. },
  453. })
  454. }, { inject: ['agents'] }))
  455. await setupStarted.promise
  456. await owner.dispose()
  457. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  458. expect(published).toEqual([])
  459. expect(ctx.agents.get(SessionId('owner-race-s'))).toBeUndefined()
  460. expect(ctx.sessions.get(SessionId('owner-race-s'))).toBeUndefined()
  461. // Let the losing callback settle; Promise.race already observes it.
  462. gate.resolve(undefined)
  463. await Promise.resolve()
  464. // The other ordering in the same race: setup resolves first (its reaction
  465. // is queued), then owner disposal flips active before that continuation can
  466. // publish. The post-race active check must still reject.
  467. const gate2 = Promise.withResolvers<undefined>()
  468. const setupStarted2 = Promise.withResolvers<undefined>()
  469. let creating2!: ReturnType<typeof ctx.agents.create>
  470. const owner2 = await ctx.plugin(Object.assign((inner: Context) => {
  471. creating2 = inner.agents.create({
  472. sessionId: SessionId('owner-race-s-2'),
  473. agentOptions: { provider: 'mock', model: 'mock' },
  474. setup: async () => {
  475. setupStarted2.resolve(undefined)
  476. await gate2.promise
  477. },
  478. })
  479. }, { inject: ['agents'] }))
  480. await setupStarted2.promise
  481. gate2.resolve(undefined)
  482. const unload2 = owner2.dispose()
  483. await expect(creating2).rejects.toThrow(/owner disposed during setup/)
  484. await unload2
  485. expect(ctx.agents.get(SessionId('owner-race-s-2'))).toBeUndefined()
  486. expect(ctx.sessions.get(SessionId('owner-race-s-2'))).toBeUndefined()
  487. })
  488. it('an AgentLoop unload aborts pending setup, awaits cleanup, and releases both ids', async () => {
  489. const { ctx, loopFiber } = await harnessWithLoop()
  490. const gate = Promise.withResolvers<undefined>()
  491. const setupStarted = Promise.withResolvers<undefined>()
  492. const published: string[] = []
  493. ctx.on('session/created', () => void published.push('session/created'))
  494. ctx.on('agent/created', () => void published.push('agent/created'))
  495. const creating = ctx.agents.create({
  496. sessionId: SessionId('factory-setup-race-s'),
  497. agentOptions: { provider: 'mock', model: 'mock' },
  498. setup: async () => {
  499. setupStarted.resolve(undefined)
  500. await gate.promise
  501. },
  502. })
  503. await setupStarted.promise
  504. await loopFiber.dispose()
  505. await expect(creating).rejects.toThrow(/agent loop is not active/)
  506. expect(published).toEqual([])
  507. expect(ctx.agents.get(SessionId('factory-setup-race-s'))).toBeUndefined()
  508. expect(ctx.sessions.get(SessionId('factory-setup-race-s'))).toBeUndefined()
  509. gate.resolve(undefined)
  510. await ctx.fiber.dispose()
  511. })
  512. it('factory unload during scope minting skips setup and awaits provisional cleanup', async () => {
  513. const { ctx, loopFiber } = await harnessWithLoop()
  514. let unloaded = false
  515. let setupCalls = 0
  516. ctx.on('internal/plugin', (fiber) => {
  517. if (unloaded || fiber.name !== 'scope') return
  518. unloaded = true
  519. void loopFiber.dispose()
  520. })
  521. const creating = ctx.agents.create({
  522. sessionId: SessionId('factory-scope-race-s'),
  523. agentOptions: { provider: 'mock', model: 'mock' },
  524. setup: () => { setupCalls += 1 },
  525. })
  526. await expect(creating).rejects.toThrow(/agent loop is not active/)
  527. await loopFiber.dispose()
  528. expect(setupCalls).toBe(1)
  529. expect(ctx.agents.get(SessionId('factory-scope-race-s'))).toBeUndefined()
  530. expect(ctx.sessions.get(SessionId('factory-scope-race-s'))).toBeUndefined()
  531. await ctx.fiber.dispose()
  532. })
  533. it('caller unload during scope minting owns and drains the half-built child', async () => {
  534. const ctx = await harness()
  535. const gate = Promise.withResolvers<undefined>()
  536. const cleanupStarted = Promise.withResolvers<undefined>()
  537. let ownerFiber!: Fiber
  538. let ownerDisposal!: Promise<void>
  539. let scopeFiber: Fiber | undefined
  540. let creating!: ReturnType<typeof ctx.agents.create>
  541. ctx.on('internal/plugin', (fiber) => {
  542. if (fiber.name !== 'scope' || scopeFiber !== undefined) return
  543. scopeFiber = fiber
  544. fiber.ctx.effect(() => async () => {
  545. cleanupStarted.resolve(undefined)
  546. await gate.promise
  547. })
  548. ownerDisposal = ownerFiber.dispose()
  549. })
  550. const owner = ctx.plugin(Object.assign((inner: Context) => {
  551. ownerFiber = inner.fiber
  552. creating = inner.agents.create({
  553. sessionId: SessionId('caller-scope-race-s'),
  554. agentOptions: { provider: 'mock', model: 'mock' },
  555. })
  556. }, { inject: ['agents'] }))
  557. await cleanupStarted.promise
  558. let ownerSettled = false
  559. void ownerDisposal.then(() => { ownerSettled = true })
  560. await Promise.resolve()
  561. expect(ownerSettled).toBe(false)
  562. gate.resolve(undefined)
  563. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  564. await ownerDisposal
  565. await owner
  566. expect(scopeFiber?.uid).toBeNull()
  567. expect(ctx.agents.get(SessionId('caller-scope-race-s'))).toBeUndefined()
  568. expect(ctx.sessions.get(SessionId('caller-scope-race-s'))).toBeUndefined()
  569. await owner.dispose()
  570. await ctx.fiber.dispose()
  571. })
  572. it('create rechecks provider liveness before its first publication edge', async () => {
  573. const { ctx, loopFiber } = await harnessWithLoop()
  574. const sessionsBefore = ctx.sessions.list().length
  575. let unloaded = false
  576. let unloading!: Promise<void>
  577. ctx.on('internal/plugin', (fiber) => {
  578. if (unloaded || fiber.name !== 'scope') return
  579. unloaded = true
  580. unloading = loopFiber.dispose()
  581. })
  582. // The mid-setup unload races publication: create either rejects or its
  583. // published agent is torn straight back down — both leave no state.
  584. await ctx.agentLoop.create(SessionId('config-scope-race'), { provider: 'mock', model: 'mock' })
  585. .then(() => undefined, () => undefined)
  586. await unloading
  587. expect(ctx.agents.get(SessionId('config-scope-race')) === undefined).toBe(true)
  588. expect(ctx.sessions.list().length).toBe(sessionsBefore)
  589. await ctx.fiber.dispose()
  590. })
  591. it('create leaves no lifecycle state when session preparation fails', async () => {
  592. const ctx = await harness()
  593. const id = SessionId('config-prepare-failure')
  594. await expect(ctx.agentLoop.create(id, { provider: 'mock', model: 'mock' }, { cwd: 'relative' }))
  595. .rejects.toThrow(/absolute path/)
  596. const replacement = await ctx.agentLoop.create(id, { provider: 'mock', model: 'mock' }, { cwd: '/recovered' })
  597. expect(ctx.agents.get(id)).toBe(replacement)
  598. await replacement.whenIdle()
  599. await ctx.fiber.dispose()
  600. })
  601. it('factory unload awaits provisional cleanup when scope preparation throws', async () => {
  602. const { ctx, loopFiber } = await harnessWithLoop()
  603. let triggered = false
  604. ctx.on('internal/plugin', (fiber) => {
  605. if (triggered || fiber.name !== 'scope') return
  606. triggered = true
  607. void loopFiber.dispose()
  608. throw new Error('scope preparation failed')
  609. })
  610. await expect(ctx.agents.create({
  611. sessionId: SessionId('factory-scope-throw-s'),
  612. agentOptions: { provider: 'mock', model: 'mock' },
  613. })).rejects.toThrow('scope preparation failed')
  614. await loopFiber.dispose()
  615. expect(ctx.agents.get(SessionId('factory-scope-throw-s'))).toBeUndefined()
  616. expect(ctx.sessions.get(SessionId('factory-scope-throw-s'))).toBeUndefined()
  617. await ctx.fiber.dispose()
  618. })
  619. it('AgentLoop unload is a structural co-owner of every live programmatic agent', async () => {
  620. const { ctx, loopFiber } = await harnessWithLoop()
  621. const loop = ctx.agentLoop
  622. const sessionId = SessionId('factory-live')
  623. const handle = await ctx.agents.create({
  624. sessionId: SessionId('factory-live-s'),
  625. agentOptions: { provider: 'mock', model: 'mock' },
  626. })
  627. await loopFiber.dispose()
  628. expect(handle.agent.status).toBe('idle')
  629. expect(ctx.agents.get(sessionId)).toBeUndefined()
  630. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  631. expect(ctx.fiber.getEffects().filter(effect => effect.label === `agentLoop.lifecycle(${sessionId})`)).toEqual([])
  632. // The consumer handle shares the provider's completed quiescence boundary.
  633. await handle.dispose()
  634. await expect(loop.createAgent(ctx, {
  635. sessionId: SessionId('factory-inactive-s'),
  636. })).rejects.toThrow(/agent loop is not active|inactive context/)
  637. await ctx.fiber.dispose()
  638. })
  639. it('keeps AgentLoop dependencies available when the caller injects only agents', async () => {
  640. const ctx = await harness()
  641. let creating!: ReturnType<typeof ctx.agents.create>
  642. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  643. creating = inner.agents.create({
  644. sessionId: SessionId('dependency-origin-s'),
  645. agentOptions: { provider: 'mock', model: 'mock' },
  646. setup: (agentCtx) => {
  647. agentCtx.tools.register(defineContentToolFixture({
  648. name: 'dependency-origin-tool',
  649. description: 'proves AgentLoop dependency origin',
  650. parameters: {},
  651. execute: () => Promise.resolve(text('ok')),
  652. }))
  653. agentCtx.systemPrompt.section({
  654. name: 'dependency-origin-section',
  655. order: 1,
  656. text: 'factory dependency API',
  657. })
  658. },
  659. })
  660. }, { inject: ['agents'] }))
  661. const handle = await creating
  662. const assembly = await ctx.systemPrompt.assemble(assembleContextFor(handle.agent))
  663. expect(assembly.tools.map(tool => tool.name)).toContain('dependency-origin-tool')
  664. expect(assembly.sections.map(section => section.name)).toContain('dependency-origin-section')
  665. await handle.dispose()
  666. await owner.dispose()
  667. await ctx.fiber.dispose()
  668. })
  669. it('keeps both entries and the scope live through a reentrant session/created teardown', async () => {
  670. const ctx = await harness()
  671. let ownerCtx!: Context
  672. let creating!: ReturnType<typeof ctx.agents.create>
  673. const lifecycle: string[] = []
  674. ctx.on('session/created', (session) => {
  675. if (session.id !== SessionId('session-created-barrier-s')) return
  676. lifecycle.push('session-created:dispose')
  677. disposeCurrentLifecycle(ownerCtx)
  678. })
  679. ctx.on('session/created', (session) => {
  680. if (session.id !== SessionId('session-created-barrier-s')) return
  681. const agent = ctx.agents.get(SessionId('session-created-barrier-s'))!
  682. expect(ctx.sessions.get(session.id)).toBe(session)
  683. expect(agent.session).toBe(session)
  684. agent.ctx.effect(() => () => { lifecycle.push('scope-disposed') })
  685. lifecycle.push('session-created:observer')
  686. })
  687. ctx.on('agent/created', () => void lifecycle.push('agent-created'))
  688. ctx.on('agent/disposed', () => void lifecycle.push('agent-disposed'))
  689. ctx.on('session/disposed', (session) => {
  690. if (session.id === SessionId('session-created-barrier-s')) lifecycle.push('session-disposed')
  691. })
  692. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  693. ownerCtx = inner
  694. creating = inner.agents.create({
  695. sessionId: SessionId('session-created-barrier-s'),
  696. agentOptions: { provider: 'mock', model: 'mock' },
  697. })
  698. }, { inject: ['agents'] }))
  699. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  700. await owner.dispose()
  701. expect(lifecycle).toEqual([
  702. 'session-created:dispose',
  703. 'session-created:observer',
  704. 'scope-disposed',
  705. 'session-disposed',
  706. ])
  707. expect(ctx.agents.get(SessionId('session-created-barrier-s'))).toBeUndefined()
  708. expect(ctx.sessions.get(SessionId('session-created-barrier-s'))).toBeUndefined()
  709. await ctx.fiber.dispose()
  710. })
  711. it('keeps both entries and the scope live through a reentrant agent/created teardown', async () => {
  712. const ctx = await harness()
  713. let ownerCtx!: Context
  714. let creating!: ReturnType<typeof ctx.agents.create>
  715. const lifecycle: string[] = []
  716. ctx.on('session/created', (session) => {
  717. if (session.id === SessionId('agent-created-barrier-s')) lifecycle.push('session-created')
  718. })
  719. ctx.on('agent/created', ({ agent }) => {
  720. if (agent.id !== SessionId('agent-created-barrier-s')) return
  721. lifecycle.push('agent-created:dispose')
  722. disposeCurrentLifecycle(ownerCtx)
  723. })
  724. ctx.on('agent/created', ({ agent }) => {
  725. if (agent.id !== SessionId('agent-created-barrier-s')) return
  726. expect(ctx.agents.get(agent.id)).toBe(agent)
  727. expect(ctx.sessions.get(agent.session.id)).toBe(agent.session)
  728. agent.ctx.effect(() => () => { lifecycle.push('scope-disposed') })
  729. lifecycle.push('agent-created:observer')
  730. })
  731. ctx.on('agent/disposed', ({ agent }) => {
  732. if (agent.id === SessionId('agent-created-barrier-s')) lifecycle.push('agent-disposed')
  733. })
  734. ctx.on('session/disposed', (session) => {
  735. if (session.id === SessionId('agent-created-barrier-s')) lifecycle.push('session-disposed')
  736. })
  737. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  738. ownerCtx = inner
  739. creating = inner.agents.create({
  740. sessionId: SessionId('agent-created-barrier-s'),
  741. agentOptions: { provider: 'mock', model: 'mock' },
  742. })
  743. }, { inject: ['agents'] }))
  744. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  745. await owner.dispose()
  746. expect(lifecycle).toEqual([
  747. 'session-created',
  748. 'agent-created:dispose',
  749. 'agent-created:observer',
  750. 'scope-disposed',
  751. 'agent-disposed',
  752. 'session-disposed',
  753. ])
  754. expect(ctx.agents.get(SessionId('agent-created-barrier-s'))).toBeUndefined()
  755. expect(ctx.sessions.get(SessionId('agent-created-barrier-s'))).toBeUndefined()
  756. await ctx.fiber.dispose()
  757. })
  758. it('rechecks caller liveness after creation listeners before unlocking the driver', async () => {
  759. const ctx = await harness()
  760. const starts: string[] = []
  761. let ownerCtx!: Context
  762. let creating!: ReturnType<typeof ctx.agents.create>
  763. ctx.on('agent/status', ({ agent, status }) => { if (status === 'running') starts.push(agent.id) })
  764. ctx.on('agent/created', ({ agent }) => {
  765. if (agent.id === SessionId('listener-dispose-s')) disposeCurrentLifecycle(ownerCtx)
  766. })
  767. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  768. ownerCtx = inner
  769. creating = inner.agents.create({
  770. sessionId: SessionId('listener-dispose-s'),
  771. agentOptions: { provider: 'mock', model: 'mock' },
  772. })
  773. }, { inject: ['agents'] }))
  774. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  775. await owner.dispose()
  776. expect(starts).toEqual([])
  777. expect(ctx.agents.get(SessionId('listener-dispose-s')) === undefined).toBe(true)
  778. expect(ctx.sessions.get(SessionId('listener-dispose-s')) === undefined).toBe(true)
  779. await ctx.fiber.dispose()
  780. })
  781. it('rechecks caller liveness after agent/created before starting the driver', async () => {
  782. const ctx = await harness()
  783. let ownerCtx!: Context
  784. let creating!: ReturnType<typeof ctx.agents.create>
  785. let announced!: Agent
  786. const statuses: string[] = []
  787. let scopeDisposed = false
  788. let observerSawLive = false
  789. ctx.on('agent/status', ({ agent, status }) => {
  790. if (agent.id === SessionId('agent/created-dispose-s')) statuses.push(status)
  791. })
  792. ctx.on('agent/created', ({ agent }) => {
  793. if (agent.id !== SessionId('agent/created-dispose-s')) return
  794. announced = agent
  795. disposeCurrentLifecycle(ownerCtx)
  796. })
  797. ctx.on('agent/created', ({ agent }) => {
  798. if (agent.id !== SessionId('agent/created-dispose-s')) return
  799. expect(ctx.agents.get(agent.id)).toBe(agent)
  800. expect(ctx.sessions.get(agent.session.id)).toBe(agent.session)
  801. agent.ctx.effect(() => () => { scopeDisposed = true })
  802. observerSawLive = true
  803. })
  804. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  805. ownerCtx = inner
  806. creating = inner.agents.create({
  807. sessionId: SessionId('agent/created-dispose-s'),
  808. agentOptions: { provider: 'mock', model: 'mock' },
  809. })
  810. }, { inject: ['agents'] }))
  811. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  812. await owner.dispose()
  813. expect(announced.status).toBe('idle')
  814. expect(statuses).toEqual([])
  815. expect(observerSawLive).toBe(true)
  816. expect(scopeDisposed).toBe(true)
  817. expect(announced.session.snapshotEvents()).toEqual([])
  818. expect(ctx.agents.get(SessionId('agent/created-dispose-s'))).toBeUndefined()
  819. expect(ctx.sessions.get(SessionId('agent/created-dispose-s'))).toBeUndefined()
  820. await ctx.fiber.dispose()
  821. })
  822. it('a rejecting setup publishes nothing and unwinds the unpublished scope', async () => {
  823. const ctx = await harness()
  824. const published: string[] = []
  825. ctx.on('session/created', () => void published.push('session/created'))
  826. ctx.on('agent/created', () => void published.push('agent/created'))
  827. await expect(ctx.agents.create({
  828. sessionId: SessionId('bad-s'),
  829. agentOptions: { provider: 'mock', model: 'mock' },
  830. setup: async () => {
  831. await Promise.resolve()
  832. throw new Error('boom setup')
  833. },
  834. })).rejects.toThrow('boom setup')
  835. // Nothing leaked: no agent, no session, and the ids are reusable.
  836. expect(published).toEqual([])
  837. expect(ctx.agents.get(SessionId('bad-s'))).toBeUndefined()
  838. expect(ctx.sessions.get(SessionId('bad-s'))).toBeUndefined()
  839. const retry = await ctx.agents.create({ sessionId: SessionId('bad-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  840. await retry.dispose()
  841. })
  842. it('rejects an exotic durable seed before publishing either object', async () => {
  843. const ctx = await harness()
  844. const published: string[] = []
  845. ctx.on('session/created', () => { published.push('session') })
  846. ctx.on('agent/created', () => { published.push('agent') })
  847. class ExoticData { readonly value = 'not durable JSON' }
  848. const seed = [{
  849. seq: 0,
  850. type: 'test/exotic-seed',
  851. data: new ExoticData(),
  852. }] as unknown as SessionEvent[]
  853. await expect(ctx.agents.create({
  854. sessionId: SessionId('exotic-seed-session'),
  855. agentOptions: { provider: 'mock', model: 'mock' },
  856. seed,
  857. })).rejects.toThrow(/seed event at index 0 is not losslessly JSON-serializable/)
  858. expect(published).toEqual([])
  859. expect(ctx.agents.get(SessionId('exotic-seed-session'))).toBeUndefined()
  860. expect(ctx.sessions.get(SessionId('exotic-seed-session'))).toBeUndefined()
  861. const retry = await ctx.agents.create({
  862. sessionId: SessionId('exotic-seed-session'),
  863. agentOptions: { provider: 'mock', model: 'mock' },
  864. })
  865. await retry.dispose()
  866. })
  867. it('a throwing session/created listener disposes the scope (pre-nesting rollback window)', async () => {
  868. const ctx = await harness()
  869. let boom = true
  870. const disposed: string[] = []
  871. ctx.on('agent/disposed', ({ agent }) => void disposed.push(agent.id))
  872. ctx.on('session/created', () => {
  873. if (boom) { boom = false; throw new Error('boom created') }
  874. })
  875. await expect(ctx.agents.create({
  876. sessionId: SessionId('bad-s'), agentOptions: { provider: 'mock', model: 'mock' },
  877. })).rejects.toThrow('boom created')
  878. expect(ctx.agents.get(SessionId('bad-s'))).toBeUndefined()
  879. expect(ctx.sessions.get(SessionId('bad-s'))).toBeUndefined()
  880. expect(disposed).toEqual([]) // inserted but never announced: no impossible disposed edge
  881. // The rollback also disposed the scope fiber: re-creating works cleanly.
  882. const retry = await ctx.agents.create({ sessionId: SessionId('bad-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  883. expect(scopeOf(retry.agent.ctx)).toBe(retry.agent)
  884. await retry.dispose()
  885. })
  886. it('pairs session and agent announcements when agent creation aborts publication', async () => {
  887. const ctx = await harness()
  888. const lifecycle: string[] = []
  889. ctx.on('session/created', (session) => { lifecycle.push(`session-created:${session.id}`) })
  890. ctx.on('session/disposed', (session) => { lifecycle.push(`session-disposed:${session.id}`) })
  891. ctx.on('agent/created', ({ agent }) => {
  892. lifecycle.push(`agent-created:${agent.id}`)
  893. throw new Error('agent observer failed')
  894. })
  895. ctx.on('agent/disposed', ({ agent }) => { lifecycle.push(`agent-disposed:${agent.id}`) })
  896. await expect(ctx.agents.create({
  897. sessionId: SessionId('partial-session'),
  898. agentOptions: { provider: 'mock', model: 'mock' },
  899. })).rejects.toThrow('agent observer failed')
  900. expect(lifecycle).toEqual([
  901. 'session-created:partial-session',
  902. 'agent-created:partial-session',
  903. 'agent-disposed:partial-session',
  904. 'session-disposed:partial-session',
  905. ])
  906. expect(ctx.agents.get(SessionId('partial-session'))).toBeUndefined()
  907. expect(ctx.sessions.get(SessionId('partial-session'))).toBeUndefined()
  908. })
  909. it('the config create helper rolls back when publication throws', async () => {
  910. const ctx = await harness()
  911. const sessionsBefore = ctx.sessions.list().length
  912. let boom = true
  913. ctx.on('session/created', () => {
  914. if (boom) {
  915. boom = false
  916. throw new Error('config publish failed')
  917. }
  918. })
  919. await expect(ctx.agentLoop.create(SessionId('config-bad'), { provider: 'mock', model: 'mock' }))
  920. .rejects.toThrow('config publish failed')
  921. await expect.poll(() => ctx.agents.get(SessionId('config-bad')) === undefined).toBe(true)
  922. await expect.poll(() => ctx.sessions.list().length).toBe(sessionsBefore)
  923. })
  924. it('registrations through a disposed agent ctx throw INACTIVE_EFFECT', async () => {
  925. const ctx = await harness()
  926. const handle = await ctx.agents.create({ sessionId: SessionId('s1'), agentOptions: { provider: 'mock', model: 'mock' } })
  927. await handle.dispose()
  928. expect(() => handle.agent.ctx.on('agent/status', () => {})).toThrow(/inactive context/)
  929. })
  930. it('agentEvents fuses carrier and subject for custom drivers', async () => {
  931. const ctx = await harness()
  932. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  933. const other = await ctx.agentLoop.create(SessionId('a2'), { provider: 'mock', model: 'mock' })
  934. const heard: string[] = []
  935. agent.ctx.on('agent/error', ({ agent: subject, turn }) => void heard.push(`${subject.id}:${turn}`))
  936. agentEvents(ctx, other).emit('agent/error', { turn: 1, step: 0, error: new Error('not for a1') })
  937. agentEvents(ctx, agent).emit('agent/error', { turn: 2, step: 0, error: new Error('for a1') })
  938. expect(heard).toEqual(['a1:2'])
  939. })
  940. it('owner unload honors the documented teardown order: unregistration AFTER the drain, before detach', async () => {
  941. const ctx = await harness()
  942. let handle!: Awaited<ReturnType<typeof ctx.agents.create>>
  943. const owner = await ctx.plugin(Object.assign(async (inner: Context) => {
  944. handle = await inner.agents.create({ sessionId: SessionId('o1-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  945. }, { inject: ['agents'] }))
  946. const { agent } = handle
  947. const order: string[] = []
  948. ctx.on('session/event', (_s, event) => {
  949. if (event.type === 'turn/end') order.push('turn-end')
  950. })
  951. ctx.on('agent/disposed', () => {
  952. order.push(`disposed(listed=${ctx.agents.get(SessionId('o1-s')) !== undefined})`)
  953. order.push(`session-still-stored=${ctx.sessions.get(SessionId('o1-s')) !== undefined}`)
  954. })
  955. // Open a turn so disposal must drain real work before registry removal.
  956. // Waiting for turn/start avoids pre-step disposal dropping the queued prompt
  957. // before a turn opens.
  958. const turnOpen = new Promise<void>((resolve) => {
  959. const off = ctx.on('session/event', (_s, event) => {
  960. if (event.type === 'turn/start') { off(); resolve() }
  961. })
  962. })
  963. agent.followup(createUserMessage({ content: text('work'), source: { kind: 'user' } }))
  964. await turnOpen
  965. await owner.dispose()
  966. expect(order).toEqual([
  967. 'turn-end',
  968. 'disposed(listed=false)',
  969. 'session-still-stored=true',
  970. ])
  971. expect(ctx.sessions.get(SessionId('o1-s'))).toBeUndefined()
  972. })
  973. it('handle.dispose() during owner unload still awaits true quiescence (shared boundary)', async () => {
  974. const ctx = await harness()
  975. let handle!: Awaited<ReturnType<typeof ctx.agents.create>>
  976. const owner = await ctx.plugin(Object.assign(async (inner: Context) => {
  977. handle = await inner.agents.create({ sessionId: SessionId('h1-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  978. }, { inject: ['agents'] }))
  979. const teardownDone: string[] = []
  980. ctx.on('agent/disposed', () => void teardownDone.push('unregistered'))
  981. // Owner unload begins FIRST (invokes the raw cordis wrapper)…
  982. const unload = owner.dispose()
  983. // …and a concurrent handle.dispose() must not resolve before the chain
  984. // actually finished (the raw wrapper returns undefined on a repeat call).
  985. await handle.dispose()
  986. expect(teardownDone).toContain('unregistered')
  987. expect(ctx.agents.get(SessionId('h1-s'))).toBeUndefined()
  988. expect(ctx.sessions.get(SessionId('h1-s'))).toBeUndefined()
  989. await unload
  990. })
  991. it('successful handle disposal retires its caller ownership effect', async () => {
  992. const ctx = await harness()
  993. const sessionId = SessionId('retired-owner-effect')
  994. const handle = await ctx.agents.create({
  995. sessionId,
  996. agentOptions: { provider: 'mock', model: 'mock' },
  997. })
  998. expect(ctx.fiber.getEffects().map(effect => effect.label)).toContain(`agentLoop.lifecycle(${sessionId})`)
  999. await handle.dispose()
  1000. expect(ctx.fiber.getEffects().filter(effect => effect.label === `agentLoop.lifecycle(${sessionId})`)).toEqual([])
  1001. await ctx.fiber.dispose()
  1002. })
  1003. it('owner unload after handle-first teardown follows the same in-flight boundary', async () => {
  1004. const ctx = await harness()
  1005. const gate = Promise.withResolvers<undefined>()
  1006. const cleanupStarted = Promise.withResolvers<undefined>()
  1007. let handle!: Awaited<ReturnType<typeof ctx.agents.create>>
  1008. const owner = await ctx.plugin(Object.assign(async (inner: Context) => {
  1009. handle = await inner.agents.create({
  1010. sessionId: SessionId('manual-first-s'),
  1011. agentOptions: { provider: 'mock', model: 'mock' },
  1012. setup(agentCtx) {
  1013. agentCtx.effect(() => async () => {
  1014. cleanupStarted.resolve(undefined)
  1015. await gate.promise
  1016. })
  1017. },
  1018. })
  1019. }, { inject: ['agents'] }))
  1020. const disposing = handle.dispose()
  1021. await cleanupStarted.promise
  1022. let ownerSettled = false
  1023. const unloading = owner.dispose().then(() => { ownerSettled = true })
  1024. await Promise.resolve()
  1025. expect(ownerSettled).toBe(false)
  1026. gate.resolve(undefined)
  1027. await Promise.all([disposing, unloading])
  1028. expect(ctx.agents.get(SessionId('manual-first-s'))).toBeUndefined()
  1029. expect(ctx.sessions.get(SessionId('manual-first-s'))).toBeUndefined()
  1030. await ctx.fiber.dispose()
  1031. })
  1032. it('reopens ids after the prior private scope finishes quiescing', async () => {
  1033. const ctx = await harness()
  1034. const gate = Promise.withResolvers<undefined>()
  1035. const cleanupStarted = Promise.withResolvers<undefined>()
  1036. const sessionId = SessionId('quiescent-reuse')
  1037. const first = await ctx.agents.create({
  1038. sessionId,
  1039. agentOptions: { provider: 'mock', model: 'mock' },
  1040. setup(agentCtx) {
  1041. agentCtx.effect(() => async () => {
  1042. cleanupStarted.resolve(undefined)
  1043. await gate.promise
  1044. })
  1045. },
  1046. })
  1047. const disposing = first.dispose()
  1048. await cleanupStarted.promise
  1049. expect(ctx.agents.get(sessionId)).toBe(first.agent)
  1050. expect(ctx.sessions.get(sessionId)).toBe(first.agent.session)
  1051. gate.resolve(undefined)
  1052. await disposing
  1053. expect(ctx.agents.get(sessionId)).toBeUndefined()
  1054. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  1055. const replacement = await ctx.agents.create({ sessionId, agentOptions: { provider: 'mock', model: 'mock' } })
  1056. expect(ctx.agents.get(sessionId)).toBe(replacement.agent)
  1057. expect(ctx.sessions.get(sessionId)).toBe(replacement.agent.session)
  1058. await replacement.dispose()
  1059. await ctx.fiber.dispose()
  1060. })
  1061. it('drains a run re-entered by cancel\'s own idle transition before removing the scope', async () => {
  1062. // Automation shaped like goal-round-driver: the running→idle transition that
  1063. // disposal's cancel produces immediately queues a follow-up prompt. The
  1064. // teardown must drain that replacement run to true quiescence instead of
  1065. // awaiting only the first captured done and unwinding under a live run.
  1066. const adapter = new MockAdapter([textResponse('one'), textResponse('never awaited')])
  1067. const ctx = await harness(adapter)
  1068. const handle = await ctx.agents.create({
  1069. sessionId: SessionId('drain-reentered-run'),
  1070. agentOptions: { provider: 'mock', model: 'mock' },
  1071. })
  1072. const agent = handle.agent
  1073. let reentered = false
  1074. ctx.on('agent/status', ({ agent: subject, status }) => {
  1075. if (subject !== agent || status !== 'idle' || reentered) return
  1076. reentered = true
  1077. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'reentrant' }], source: { kind: 'user' } }))
  1078. })
  1079. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'go' }], source: { kind: 'user' } }))
  1080. await waitForIdle(ctx, agent)
  1081. expect(reentered).toBe(true)
  1082. // Idle again: the reentrant batch was already claimed and settled (its
  1083. // prompt was blocked by nothing, so it ran) — arm a SECOND reentry that
  1084. // fires from the disposal cancel's idle transition itself.
  1085. reentered = false
  1086. await handle.dispose()
  1087. // The reentrant run either never started or was drained: the registries
  1088. // are empty and nothing still drives the detached session.
  1089. expect(ctx.agents.get(agent.id)).toBeUndefined()
  1090. expect(ctx.sessions.get(agent.id)).toBeUndefined()
  1091. const eventsAfter = agent.session.snapshotEvents().length
  1092. await new Promise(resolve => setTimeout(resolve, 30))
  1093. expect(agent.session.snapshotEvents().length).toBe(eventsAfter)
  1094. await ctx.fiber.dispose()
  1095. })
  1096. })