scope-lifecycle.spec.ts 49 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184
  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 agent/created before starting the driver', async () => {
  759. const ctx = await harness()
  760. let ownerCtx!: Context
  761. let creating!: ReturnType<typeof ctx.agents.create>
  762. let announced!: Agent
  763. const statuses: string[] = []
  764. let scopeDisposed = false
  765. let observerSawLive = false
  766. ctx.on('agent/status', ({ agent, status }) => {
  767. if (agent.id === SessionId('agent/created-dispose-s')) statuses.push(status)
  768. })
  769. ctx.on('agent/created', ({ agent }) => {
  770. if (agent.id !== SessionId('agent/created-dispose-s')) return
  771. announced = agent
  772. disposeCurrentLifecycle(ownerCtx)
  773. })
  774. ctx.on('agent/created', ({ agent }) => {
  775. if (agent.id !== SessionId('agent/created-dispose-s')) return
  776. expect(ctx.agents.get(agent.id)).toBe(agent)
  777. expect(ctx.sessions.get(agent.session.id)).toBe(agent.session)
  778. agent.ctx.effect(() => () => { scopeDisposed = true })
  779. observerSawLive = true
  780. })
  781. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  782. ownerCtx = inner
  783. creating = inner.agents.create({
  784. sessionId: SessionId('agent/created-dispose-s'),
  785. agentOptions: { provider: 'mock', model: 'mock' },
  786. })
  787. }, { inject: ['agents'] }))
  788. await expect(creating).rejects.toThrow(/owner disposed during setup/)
  789. await owner.dispose()
  790. expect(announced.status).toBe('idle')
  791. expect(statuses).toEqual([])
  792. expect(observerSawLive).toBe(true)
  793. expect(scopeDisposed).toBe(true)
  794. expect(announced.session.snapshotEvents()).toEqual([])
  795. expect(ctx.agents.get(SessionId('agent/created-dispose-s'))).toBeUndefined()
  796. expect(ctx.sessions.get(SessionId('agent/created-dispose-s'))).toBeUndefined()
  797. await ctx.fiber.dispose()
  798. })
  799. it('a rejecting setup publishes nothing and unwinds the unpublished scope', async () => {
  800. const ctx = await harness()
  801. const published: string[] = []
  802. ctx.on('session/created', () => void published.push('session/created'))
  803. ctx.on('agent/created', () => void published.push('agent/created'))
  804. await expect(ctx.agents.create({
  805. sessionId: SessionId('bad-s'),
  806. agentOptions: { provider: 'mock', model: 'mock' },
  807. setup: async () => {
  808. await Promise.resolve()
  809. throw new Error('boom setup')
  810. },
  811. })).rejects.toThrow('boom setup')
  812. // Nothing leaked: no agent, no session, and the ids are reusable.
  813. expect(published).toEqual([])
  814. expect(ctx.agents.get(SessionId('bad-s'))).toBeUndefined()
  815. expect(ctx.sessions.get(SessionId('bad-s'))).toBeUndefined()
  816. const retry = await ctx.agents.create({ sessionId: SessionId('bad-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  817. await retry.dispose()
  818. })
  819. it('rejects an exotic durable seed before publishing either object', async () => {
  820. const ctx = await harness()
  821. const published: string[] = []
  822. ctx.on('session/created', () => { published.push('session') })
  823. ctx.on('agent/created', () => { published.push('agent') })
  824. class ExoticData { readonly value = 'not durable JSON' }
  825. const seed = [{
  826. seq: 0,
  827. type: 'test/exotic-seed',
  828. data: new ExoticData(),
  829. }] as unknown as SessionEvent[]
  830. await expect(ctx.agents.create({
  831. sessionId: SessionId('exotic-seed-session'),
  832. agentOptions: { provider: 'mock', model: 'mock' },
  833. seed,
  834. })).rejects.toThrow(/seed event at index 0 is not losslessly JSON-serializable/)
  835. expect(published).toEqual([])
  836. expect(ctx.agents.get(SessionId('exotic-seed-session'))).toBeUndefined()
  837. expect(ctx.sessions.get(SessionId('exotic-seed-session'))).toBeUndefined()
  838. const retry = await ctx.agents.create({
  839. sessionId: SessionId('exotic-seed-session'),
  840. agentOptions: { provider: 'mock', model: 'mock' },
  841. })
  842. await retry.dispose()
  843. })
  844. it('a throwing session/created listener disposes the scope (pre-nesting rollback window)', async () => {
  845. const ctx = await harness()
  846. let boom = true
  847. const disposed: string[] = []
  848. ctx.on('agent/disposed', ({ agent }) => void disposed.push(agent.id))
  849. ctx.on('session/created', () => {
  850. if (boom) { boom = false; throw new Error('boom created') }
  851. })
  852. await expect(ctx.agents.create({
  853. sessionId: SessionId('bad-s'), agentOptions: { provider: 'mock', model: 'mock' },
  854. })).rejects.toThrow('boom created')
  855. expect(ctx.agents.get(SessionId('bad-s'))).toBeUndefined()
  856. expect(ctx.sessions.get(SessionId('bad-s'))).toBeUndefined()
  857. expect(disposed).toEqual([]) // inserted but never announced: no impossible disposed edge
  858. // The rollback also disposed the scope fiber: re-creating works cleanly.
  859. const retry = await ctx.agents.create({ sessionId: SessionId('bad-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  860. expect(scopeOf(retry.agent.ctx)).toBe(retry.agent)
  861. await retry.dispose()
  862. })
  863. it('pairs session and agent announcements when agent creation aborts publication', async () => {
  864. const ctx = await harness()
  865. const lifecycle: string[] = []
  866. ctx.on('session/created', (session) => { lifecycle.push(`session-created:${session.id}`) })
  867. ctx.on('session/disposed', (session) => { lifecycle.push(`session-disposed:${session.id}`) })
  868. ctx.on('agent/created', ({ agent }) => {
  869. lifecycle.push(`agent-created:${agent.id}`)
  870. throw new Error('agent observer failed')
  871. })
  872. ctx.on('agent/disposed', ({ agent }) => { lifecycle.push(`agent-disposed:${agent.id}`) })
  873. await expect(ctx.agents.create({
  874. sessionId: SessionId('partial-session'),
  875. agentOptions: { provider: 'mock', model: 'mock' },
  876. })).rejects.toThrow('agent observer failed')
  877. expect(lifecycle).toEqual([
  878. 'session-created:partial-session',
  879. 'agent-created:partial-session',
  880. 'agent-disposed:partial-session',
  881. 'session-disposed:partial-session',
  882. ])
  883. expect(ctx.agents.get(SessionId('partial-session'))).toBeUndefined()
  884. expect(ctx.sessions.get(SessionId('partial-session'))).toBeUndefined()
  885. })
  886. it('the config create helper rolls back when publication throws', async () => {
  887. const ctx = await harness()
  888. const sessionsBefore = ctx.sessions.list().length
  889. let boom = true
  890. ctx.on('session/created', () => {
  891. if (boom) {
  892. boom = false
  893. throw new Error('config publish failed')
  894. }
  895. })
  896. await expect(ctx.agentLoop.create(SessionId('config-bad'), { provider: 'mock', model: 'mock' }))
  897. .rejects.toThrow('config publish failed')
  898. await expect.poll(() => ctx.agents.get(SessionId('config-bad')) === undefined).toBe(true)
  899. await expect.poll(() => ctx.sessions.list().length).toBe(sessionsBefore)
  900. })
  901. it('registrations through a disposed agent ctx throw INACTIVE_EFFECT', async () => {
  902. const ctx = await harness()
  903. const handle = await ctx.agents.create({ sessionId: SessionId('s1'), agentOptions: { provider: 'mock', model: 'mock' } })
  904. await handle.dispose()
  905. expect(() => handle.agent.ctx.on('agent/status', () => {})).toThrow(/inactive context/)
  906. })
  907. it('agentEvents fuses carrier and subject for custom drivers', async () => {
  908. const ctx = await harness()
  909. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  910. const other = await ctx.agentLoop.create(SessionId('a2'), { provider: 'mock', model: 'mock' })
  911. const heard: string[] = []
  912. agent.ctx.on('agent/error', ({ agent: subject, turn }) => void heard.push(`${subject.id}:${turn}`))
  913. agentEvents(ctx, other).emit('agent/error', { turn: 1, step: 0, error: new Error('not for a1') })
  914. agentEvents(ctx, agent).emit('agent/error', { turn: 2, step: 0, error: new Error('for a1') })
  915. expect(heard).toEqual(['a1:2'])
  916. })
  917. it('owner unload honors the documented teardown order: unregistration AFTER the drain, before detach', async () => {
  918. const ctx = await harness()
  919. let handle!: Awaited<ReturnType<typeof ctx.agents.create>>
  920. const owner = await ctx.plugin(Object.assign(async (inner: Context) => {
  921. handle = await inner.agents.create({ sessionId: SessionId('o1-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  922. }, { inject: ['agents'] }))
  923. const { agent } = handle
  924. const order: string[] = []
  925. ctx.on('session/event', (_s, event) => {
  926. if (event.type === 'turn/end') order.push('turn-end')
  927. })
  928. ctx.on('agent/disposed', () => {
  929. order.push(`disposed(listed=${ctx.agents.get(SessionId('o1-s')) !== undefined})`)
  930. order.push(`session-still-stored=${ctx.sessions.get(SessionId('o1-s')) !== undefined}`)
  931. })
  932. // Open a turn so disposal must drain real work before registry removal.
  933. // Waiting for turn/start avoids pre-step disposal dropping the queued prompt
  934. // before a turn opens.
  935. const turnOpen = new Promise<void>((resolve) => {
  936. const off = ctx.on('session/event', (_s, event) => {
  937. if (event.type === 'turn/start') { off(); resolve() }
  938. })
  939. })
  940. agent.followup(createUserMessage({ content: text('work'), source: { kind: 'user' } }))
  941. await turnOpen
  942. await owner.dispose()
  943. expect(order).toEqual([
  944. 'turn-end',
  945. 'disposed(listed=false)',
  946. 'session-still-stored=true',
  947. ])
  948. expect(ctx.sessions.get(SessionId('o1-s'))).toBeUndefined()
  949. })
  950. it('handle.dispose() during owner unload still awaits true quiescence (shared boundary)', async () => {
  951. const ctx = await harness()
  952. let handle!: Awaited<ReturnType<typeof ctx.agents.create>>
  953. const owner = await ctx.plugin(Object.assign(async (inner: Context) => {
  954. handle = await inner.agents.create({ sessionId: SessionId('h1-s'), agentOptions: { provider: 'mock', model: 'mock' } })
  955. }, { inject: ['agents'] }))
  956. const teardownDone: string[] = []
  957. ctx.on('agent/disposed', () => void teardownDone.push('unregistered'))
  958. // Owner unload begins FIRST (invokes the raw cordis wrapper)…
  959. const unload = owner.dispose()
  960. // …and a concurrent handle.dispose() must not resolve before the chain
  961. // actually finished (the raw wrapper returns undefined on a repeat call).
  962. await handle.dispose()
  963. expect(teardownDone).toContain('unregistered')
  964. expect(ctx.agents.get(SessionId('h1-s'))).toBeUndefined()
  965. expect(ctx.sessions.get(SessionId('h1-s'))).toBeUndefined()
  966. await unload
  967. })
  968. it('successful handle disposal retires its caller ownership effect', async () => {
  969. const ctx = await harness()
  970. const sessionId = SessionId('retired-owner-effect')
  971. const handle = await ctx.agents.create({
  972. sessionId,
  973. agentOptions: { provider: 'mock', model: 'mock' },
  974. })
  975. expect(ctx.fiber.getEffects().map(effect => effect.label)).toContain(`agentLoop.lifecycle(${sessionId})`)
  976. await handle.dispose()
  977. expect(ctx.fiber.getEffects().filter(effect => effect.label === `agentLoop.lifecycle(${sessionId})`)).toEqual([])
  978. await ctx.fiber.dispose()
  979. })
  980. it('owner unload after handle-first teardown follows the same in-flight boundary', async () => {
  981. const ctx = await harness()
  982. const gate = Promise.withResolvers<undefined>()
  983. const cleanupStarted = Promise.withResolvers<undefined>()
  984. let handle!: Awaited<ReturnType<typeof ctx.agents.create>>
  985. const owner = await ctx.plugin(Object.assign(async (inner: Context) => {
  986. handle = await inner.agents.create({
  987. sessionId: SessionId('manual-first-s'),
  988. agentOptions: { provider: 'mock', model: 'mock' },
  989. setup(agentCtx) {
  990. agentCtx.effect(() => async () => {
  991. cleanupStarted.resolve(undefined)
  992. await gate.promise
  993. })
  994. },
  995. })
  996. }, { inject: ['agents'] }))
  997. const disposing = handle.dispose()
  998. await cleanupStarted.promise
  999. let ownerSettled = false
  1000. const unloading = owner.dispose().then(() => { ownerSettled = true })
  1001. await Promise.resolve()
  1002. expect(ownerSettled).toBe(false)
  1003. gate.resolve(undefined)
  1004. await Promise.all([disposing, unloading])
  1005. expect(ctx.agents.get(SessionId('manual-first-s'))).toBeUndefined()
  1006. expect(ctx.sessions.get(SessionId('manual-first-s'))).toBeUndefined()
  1007. await ctx.fiber.dispose()
  1008. })
  1009. it('reopens ids after the prior private scope finishes quiescing', async () => {
  1010. const ctx = await harness()
  1011. const gate = Promise.withResolvers<undefined>()
  1012. const cleanupStarted = Promise.withResolvers<undefined>()
  1013. const sessionId = SessionId('quiescent-reuse')
  1014. const first = await ctx.agents.create({
  1015. sessionId,
  1016. agentOptions: { provider: 'mock', model: 'mock' },
  1017. setup(agentCtx) {
  1018. agentCtx.effect(() => async () => {
  1019. cleanupStarted.resolve(undefined)
  1020. await gate.promise
  1021. })
  1022. },
  1023. })
  1024. const disposing = first.dispose()
  1025. await cleanupStarted.promise
  1026. expect(ctx.agents.get(sessionId)).toBe(first.agent)
  1027. expect(ctx.sessions.get(sessionId)).toBe(first.agent.session)
  1028. gate.resolve(undefined)
  1029. await disposing
  1030. expect(ctx.agents.get(sessionId)).toBeUndefined()
  1031. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  1032. const replacement = await ctx.agents.create({ sessionId, agentOptions: { provider: 'mock', model: 'mock' } })
  1033. expect(ctx.agents.get(sessionId)).toBe(replacement.agent)
  1034. expect(ctx.sessions.get(sessionId)).toBe(replacement.agent.session)
  1035. await replacement.dispose()
  1036. await ctx.fiber.dispose()
  1037. })
  1038. it('drains a run re-entered by cancel\'s own idle transition before removing the scope', async () => {
  1039. // Automation shaped like goal-round-driver: the running→idle transition that
  1040. // disposal's cancel produces immediately queues a follow-up prompt. The
  1041. // teardown must drain that replacement run to true quiescence instead of
  1042. // awaiting only the first captured done and unwinding under a live run.
  1043. const adapter = new MockAdapter([textResponse('one'), textResponse('never awaited')])
  1044. const ctx = await harness(adapter)
  1045. const handle = await ctx.agents.create({
  1046. sessionId: SessionId('drain-reentered-run'),
  1047. agentOptions: { provider: 'mock', model: 'mock' },
  1048. })
  1049. const agent = handle.agent
  1050. let reentered = false
  1051. ctx.on('agent/status', ({ agent: subject, status }) => {
  1052. if (subject !== agent || status !== 'idle' || reentered) return
  1053. reentered = true
  1054. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'reentrant' }], source: { kind: 'user' } }))
  1055. })
  1056. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'go' }], source: { kind: 'user' } }))
  1057. await waitForIdle(ctx, agent)
  1058. expect(reentered).toBe(true)
  1059. // Idle again: the reentrant batch was already claimed and settled (its
  1060. // prompt was blocked by nothing, so it ran) — arm a SECOND reentry that
  1061. // fires from the disposal cancel's idle transition itself.
  1062. reentered = false
  1063. await handle.dispose()
  1064. // The reentrant run either never started or was drained: the registries
  1065. // are empty and nothing still drives the detached session.
  1066. expect(ctx.agents.get(agent.id)).toBeUndefined()
  1067. expect(ctx.sessions.get(agent.id)).toBeUndefined()
  1068. const eventsAfter = agent.session.snapshotEvents().length
  1069. await new Promise(resolve => setTimeout(resolve, 30))
  1070. expect(agent.session.snapshotEvents().length).toBe(eventsAfter)
  1071. await ctx.fiber.dispose()
  1072. })
  1073. })