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