scope-lifecycle.spec.ts 44 KB

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