agent.spec.ts 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { AgentId } from '@deepseek-ai/dsh-agent'
  4. import LlmService from '@deepseek-ai/dsh-llm'
  5. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  6. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  7. import ToolRegistry from '@deepseek-ai/dsh-tools'
  8. import AgentRegistry from '@deepseek-ai/dsh-agent'
  9. import AgentExecutionProvider from '@deepseek-ai/dsh-agent-execution'
  10. import AgentLoop, { ReactLoopAgent } from '@deepseek-ai/dsh-agent-loop'
  11. import { bindReactLoopAgentContext, prepareReactLoopAgent } from '../src/agent.ts'
  12. import { MockAdapter, textResponse } from './mock-adapter.ts'
  13. async function harness(adapter: MockAdapter) {
  14. const ctx = new Context()
  15. await ctx.plugin(LlmService)
  16. await ctx.plugin(SessionStore)
  17. await ctx.plugin(SystemPrompt)
  18. await ctx.plugin(ToolRegistry)
  19. await ctx.plugin(AgentRegistry)
  20. await ctx.plugin(AgentExecutionProvider)
  21. await ctx.plugin(AgentLoop, { agents: [] })
  22. ctx.llm.registerAdapter(['mock'], adapter)
  23. return ctx
  24. }
  25. function waitForIdle(ctx: Context, agent: ReactLoopAgent): Promise<void> {
  26. return new Promise((resolve) => {
  27. const dispose = ctx.on('agent/status', (subject, status) => {
  28. if (subject === agent && status === 'idle') {
  29. dispose()
  30. resolve()
  31. }
  32. })
  33. })
  34. }
  35. function waitForStatus(ctx: Context, agent: ReactLoopAgent, expected: ReactLoopAgent['status']): Promise<void> {
  36. return new Promise((resolve) => {
  37. const dispose = ctx.on('agent/status', (subject, status) => {
  38. if (subject === agent && status === expected) {
  39. dispose()
  40. resolve()
  41. }
  42. })
  43. })
  44. }
  45. function send(agent: ReactLoopAgent, text: string) {
  46. agent.send([{ type: 'text', text }])
  47. }
  48. describe('ReactLoopAgent', () => {
  49. it('rejects access before context binding and a second driver for one session', async () => {
  50. const ctx = new Context()
  51. await ctx.plugin(AgentExecutionProvider)
  52. await ctx.plugin(SessionStore)
  53. const session = ctx.sessions.create(SessionId('exclusive-driver'))
  54. const prepared = prepareReactLoopAgent(ctx, AgentId('first-driver'), { provider: 'mock', model: 'mock' }, session)
  55. expect(() => prepared.agent.ctx).toThrow('context is not bound')
  56. expect(() => prepareReactLoopAgent(ctx, AgentId('second-driver'), { provider: 'mock', model: 'mock' }, session))
  57. .toThrow('already has a concrete agent driver')
  58. await prepared.dispose()
  59. await ctx.fiber.dispose()
  60. })
  61. it('borrows caller options and binds its scoped context exactly once', async () => {
  62. const ctx = await harness(new MockAdapter([textResponse('unused')]))
  63. const options = { provider: 'mock', model: 'mock' }
  64. const agent = ctx.agentLoop.create(AgentId('owned-bindings'), options)
  65. expect(agent.options).toBe(options)
  66. expect(agent.id).toBe('owned-bindings')
  67. expect(agent.session.id).toMatch(/^owned-bindings-session-/)
  68. expect(() => { bindReactLoopAgentContext(agent, new Context()) }).toThrow(/context is already bound/)
  69. await ctx.fiber.dispose()
  70. })
  71. it('send() throws after disposal', async () => {
  72. const adapter = new MockAdapter(['hang'])
  73. const ctx = await harness(adapter)
  74. let agent!: ReactLoopAgent
  75. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  76. agent = inner.agentLoop.create(AgentId('scoped'), { provider: 'mock', model: 'mock' })
  77. }, { inject: ['agentLoop'] }))
  78. send(agent, 'go')
  79. await new Promise(r => setTimeout(r, 30))
  80. await fiber.dispose()
  81. await agent.done
  82. expect(() => { agent.send([{ type: 'text', text: 'too late' }]) }).toThrow('disposed')
  83. })
  84. it('steer() throws after disposal', async () => {
  85. const adapter = new MockAdapter(['hang'])
  86. const ctx = await harness(adapter)
  87. let agent!: ReactLoopAgent
  88. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  89. agent = inner.agentLoop.create(AgentId('scoped'), { provider: 'mock', model: 'mock' })
  90. }, { inject: ['agentLoop'] }))
  91. send(agent, 'go')
  92. await new Promise(r => setTimeout(r, 30))
  93. await fiber.dispose()
  94. await agent.done
  95. expect(() => { agent.steer([{ type: 'text', text: 'too late' }]) }).toThrow('disposed')
  96. })
  97. it('inject() throws after disposal', async () => {
  98. const adapter = new MockAdapter(['hang'])
  99. const ctx = await harness(adapter)
  100. let agent!: ReactLoopAgent
  101. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  102. agent = inner.agentLoop.create(AgentId('scoped'), { provider: 'mock', model: 'mock' })
  103. }, { inject: ['agentLoop'] }))
  104. send(agent, 'go')
  105. await new Promise(r => setTimeout(r, 30))
  106. await fiber.dispose()
  107. await agent.done
  108. expect(() => { agent.inject([{ type: 'text', text: 'too late' }]) }).toThrow('disposed')
  109. })
  110. it('inject() decides enclosure from the LOG (open turn), not agent status', async () => {
  111. const adapter = new MockAdapter([textResponse('ok')])
  112. const ctx = await harness(adapter)
  113. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  114. // Simulate an OPEN turn in the log while the agent is idle (status is not a
  115. // reliable open-turn signal). inject must append into that open turn, NOT
  116. // wrap a new one.
  117. agent.session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  118. agent.inject([{ type: 'text', text: 'mid' }], { source: { kind: 'plugin', plugin: 'p' } })
  119. expect(agent.session.events.filter(e => e.type === 'turn/start')).toHaveLength(1)
  120. expect(agent.session.events.at(-1)!.type).toBe('context/message')
  121. // Close the turn; now inject must wrap its own one-shot injection turn.
  122. agent.session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  123. agent.inject([{ type: 'text', text: 'after' }], { source: { kind: 'plugin', plugin: 'p' } })
  124. const starts = agent.session.events.filter(e => e.type === 'turn/start')
  125. expect(starts).toHaveLength(2)
  126. const last = starts[1]!
  127. expect(last.type === 'turn/start' && last.data.trigger.kind).toBe('injection')
  128. expect(agent.session.events.at(-1)!.type).toBe('turn/end') // turn-enclosed
  129. })
  130. it('idle inject() contains a failing flush (logs, does not throw into the caller)', async () => {
  131. const adapter = new MockAdapter([textResponse('ok')])
  132. const ctx = await harness(adapter)
  133. // A persistence-like listener whose flush rejects.
  134. ctx.on('session/flush', () => { throw new Error('disk gone') })
  135. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  136. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  137. // inject() is synchronous and fires a fire-and-forget flush; a rejecting
  138. // flush must be contained (logged), never thrown into the caller.
  139. expect(() => { agent.inject([{ type: 'text', text: 'notice' }], { source: { kind: 'plugin', plugin: 'p' } }) }).not.toThrow()
  140. await new Promise(r => setTimeout(r, 20)) // let the contained flush settle
  141. expect(warn).toHaveBeenCalledWith(expect.stringContaining('flush after idle injection failed'))
  142. warn.mockRestore()
  143. })
  144. it('idle inject() closes its one-shot turn AND still checkpoints even if the append throws', async () => {
  145. const adapter = new MockAdapter([textResponse('ok')])
  146. const ctx = await harness(adapter)
  147. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  148. let flushes = 0
  149. ctx.on('session/flush', () => { flushes += 1 })
  150. // Invalid injected content throws after turn/start. `finally` must still append turn/end and
  151. // flush the balanced in-memory turn so a crash cannot lose it before the next checkpoint.
  152. expect(() => {
  153. agent.inject([{ type: 'text', text: 'x', bad: 1n } as never], { source: { kind: 'plugin', plugin: 'p' } })
  154. }).toThrow(/non-JSON-serializable/)
  155. const types = agent.session.events.map(e => e.type)
  156. expect(types).toEqual(['turn/start', 'turn/end']) // balanced, no open turn
  157. await new Promise(r => setTimeout(r, 10)) // let the fire-and-forget flush run
  158. expect(flushes).toBe(1) // checkpoint fired despite the throw
  159. })
  160. it('idle inject() still checkpoints when a listener throws on the synthetic turn/end', async () => {
  161. const adapter = new MockAdapter([textResponse('ok')])
  162. const ctx = await harness(adapter)
  163. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  164. let flushes = 0
  165. ctx.on('session/flush', () => { flushes += 1 })
  166. // Session contains a throwing post-commit turn/end observer. The accepted
  167. // boundary still triggers the idle injection's durability checkpoint.
  168. let threw = false
  169. ctx.on('session/event', (_s, event) => {
  170. if (!threw && event.type === 'turn/end') { threw = true; throw new Error('boom turn/end') }
  171. })
  172. expect(() => { agent.inject([{ type: 'text', text: 'notice' }], { source: { kind: 'plugin', plugin: 'p' } }) }).not.toThrow()
  173. const types = agent.session.events.map(e => e.type)
  174. expect(types).toEqual(['turn/start', 'context/message', 'turn/end']) // balanced
  175. await new Promise(r => setTimeout(r, 10))
  176. expect(flushes).toBe(1) // checkpoint fired despite the throwing turn/end listener
  177. })
  178. it('idle inject() reports a failing flush via agent/error (step 0) AND the logger', async () => {
  179. const adapter = new MockAdapter([textResponse('ok')])
  180. const ctx = await harness(adapter)
  181. // A non-Error rejection exercises the String() normalization branch.
  182. ctx.on('session/flush', () => { throw 'disk gone' })
  183. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  184. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  185. const errors: { turn: number; step: number; message: string }[] = []
  186. ctx.on('agent/error', (_a, turn, step, error) => void errors.push({ turn, step, message: error.message }))
  187. agent.inject([{ type: 'text', text: 'notice' }], { source: { kind: 'plugin', plugin: 'p' } })
  188. await new Promise(r => setTimeout(r, 20)) // let the contained flush settle
  189. // Reported via agent/error (step 0 — the idle-injection convention) so
  190. // plugins monitoring agent/error see idle-injection persistence failures,
  191. // mirroring the loop's post-turn/end flush path. A non-Error throw is
  192. // normalized to an Error.
  193. expect(errors).toEqual([{ turn: 1, step: 0, message: 'disk gone' }])
  194. expect(warn).toHaveBeenCalledWith(expect.stringContaining('flush after idle injection failed'))
  195. warn.mockRestore()
  196. })
  197. it('idle inject() with a non-serializable source opens no turn (nothing to close)', async () => {
  198. const adapter = new MockAdapter([textResponse('ok')])
  199. const ctx = await harness(adapter)
  200. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  201. // A non-serializable source makes the turn/start append throw BEFORE the
  202. // event is pushed (Session.append validates before push), so NO turn opens.
  203. // The finally's isTurnOpen() guard sees no open turn and appends nothing —
  204. // the log stays empty, not left with a dangling turn/start.
  205. expect(() => {
  206. agent.inject([{ type: 'text', text: 'x' }], { source: { kind: 'plugin', plugin: 'p', bad: 1n } as never })
  207. }).toThrow(/non-JSON-serializable/)
  208. expect(agent.session.events).toHaveLength(0)
  209. })
  210. it('steer() when idle falls through to send() and starts a turn', async () => {
  211. const adapter = new MockAdapter([textResponse('ok')])
  212. const ctx = await harness(adapter)
  213. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  214. // steer while idle delegates to send
  215. agent.steer([{ type: 'text', text: 'steer idle' }], { source: { kind: 'plugin', plugin: 'test' } })
  216. await waitForIdle(ctx, agent)
  217. // The message was recorded as a user-level message (send path)
  218. expect(agent.session.events.some(e => e.type === 'user/message')).toBe(true)
  219. expect(adapter.requests).toHaveLength(1)
  220. })
  221. it('disposer is idempotent (double-stop)', async () => {
  222. // The internal start seam exposes one idle driver's disposer for repeated invocation.
  223. const ctx = new Context()
  224. await ctx.plugin(AgentExecutionProvider)
  225. await ctx.plugin(SessionStore)
  226. const session = ctx.sessions.create(SessionId('test'))
  227. const prepared = prepareReactLoopAgent(ctx, AgentId('bare'), { provider: 'mock', model: 'mock' }, session)
  228. const { agent } = prepared
  229. prepared.markPublished()
  230. const dispose = prepared.startDriver()
  231. const firstDisposal = dispose()
  232. expect(agent.status).toBe('disposed')
  233. await firstDisposal
  234. await expect(dispose()).resolves.toBeUndefined()
  235. expect(agent.status).toBe('disposed')
  236. })
  237. it('a pre-start disposal makes a later driver-start attempt inert', async () => {
  238. const ctx = new Context()
  239. await ctx.plugin(SessionStore)
  240. const session = ctx.sessions.create(SessionId('pre-start-dispose'))
  241. const prepared = prepareReactLoopAgent(ctx, AgentId('pre-start-dispose'), { provider: 'mock', model: 'mock' }, session)
  242. await prepared.dispose()
  243. expect(prepared.agent.status).toBe('disposed')
  244. const dispose = prepared.startDriver()
  245. await dispose()
  246. await expect(prepared.agent.done).resolves.toBeUndefined()
  247. expect(prepared.agent.session.events).toEqual([])
  248. await ctx.fiber.dispose()
  249. })
  250. it('setting the same status does not emit agent/status again', async () => {
  251. const adapter = new MockAdapter([textResponse('ok')])
  252. const ctx = await harness(adapter)
  253. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  254. const statuses: string[] = []
  255. ctx.on('agent/status', (subject, status) => {
  256. if (subject === agent) statuses.push(status)
  257. })
  258. send(agent, 'hi')
  259. await waitForIdle(ctx, agent)
  260. // After the turn, agent is idle. Send again to trigger another attempt
  261. // to go idle — but it's already idle, so no emission.
  262. const idleTransitionCount = statuses.filter(s => s === 'idle').length
  263. expect(idleTransitionCount).toBe(1) // only the final transition from running
  264. })
  265. it('whenIdle() resolves immediately when the agent is not running', async () => {
  266. const adapter = new MockAdapter([textResponse('ok')])
  267. const ctx = await harness(adapter)
  268. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  269. // Fresh agent is idle — whenIdle() takes the not-running fast path and
  270. // resolves without subscribing. await must not hang.
  271. await agent.whenIdle()
  272. expect(agent.status).not.toBe('running')
  273. })
  274. it('whenIdle() waits for queued work that has not flipped status yet', async () => {
  275. const adapter = new MockAdapter(['hang'])
  276. const ctx = await harness(adapter)
  277. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  278. send(agent, 'queued')
  279. let settled = false
  280. const idle = agent.whenIdle().then(() => { settled = true })
  281. await Promise.resolve()
  282. expect(settled).toBe(false)
  283. await waitForStatus(ctx, agent, 'running')
  284. agent.cancel({ kind: 'user' })
  285. await idle
  286. expect(settled).toBe(true)
  287. expect(agent.status).toBe('idle')
  288. })
  289. it('whenIdle() awaits the running→idle transition, ignoring other subjects/running events', async () => {
  290. const adapter = new MockAdapter([textResponse('ok'), textResponse('ok')])
  291. const ctx = await harness(adapter)
  292. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  293. const other = ctx.agentLoop.create(AgentId('a2'), { provider: 'mock', model: 'mock' })
  294. // Drive `agent` into `running`, then await whenIdle() — it subscribes to
  295. // agent/status and resolves on the first transition out of running.
  296. const running = new Promise<void>((resolve) => {
  297. const dispose = ctx.on('agent/status', (subject, status) => {
  298. if (subject === agent && status === 'running') { dispose(); resolve() }
  299. })
  300. })
  301. send(agent, 'go')
  302. await running
  303. expect(agent.status).toBe('running')
  304. // While `agent`'s whenIdle is pending, churn `other` through running→idle:
  305. // every status event it emits hits whenIdle's guard with `subject !== this`,
  306. // so the wait must ignore them and only resolve on `agent`'s own idle.
  307. send(other, 'go')
  308. await agent.whenIdle()
  309. expect(agent.status).toBe('idle')
  310. })
  311. it('whenIdle() subscribed while running resolves via done when the agent is then disposed', async () => {
  312. // Queue the internal waiter while running, then dispose the bare driver. Its disposed branch
  313. // must chain the loop's `done` promise rather than resolve before exit.
  314. const ctx = new Context()
  315. await ctx.plugin(AgentExecutionProvider)
  316. await ctx.plugin(LlmService)
  317. await ctx.plugin(SessionStore)
  318. await ctx.plugin(SystemPrompt)
  319. await ctx.plugin(ToolRegistry)
  320. await ctx.plugin(AgentRegistry)
  321. const adapter = new MockAdapter(['hang'])
  322. ctx.llm.registerAdapter(['mock'], adapter)
  323. const session = ctx.sessions.create(SessionId('bare'))
  324. const prepared = prepareReactLoopAgent(ctx, AgentId('bare'), { provider: 'mock', model: 'mock' }, session)
  325. const { agent } = prepared
  326. prepared.markPublished()
  327. const dispose = prepared.startDriver()
  328. agent.send([{ type: 'text', text: 'go' }])
  329. await new Promise(r => setTimeout(r, 30))
  330. expect(agent.status).toBe('running')
  331. const idle = agent.whenIdle() // queues an internal waiter (running)
  332. const disposal = dispose() // settles the waiter synchronously; whenIdle chains done
  333. await idle
  334. expect(agent.status).toBe('disposed')
  335. await disposal
  336. })
  337. it('whenIdle() subscribed while running survives a FIBER dispose (no hung promise)', async () => {
  338. // The waiter is agent-owned state, not an effect-scoped listener that owner disposal would
  339. // remove before the disposed transition. Fiber teardown must still settle it.
  340. const adapter = new MockAdapter(['hang'])
  341. const ctx = await harness(adapter)
  342. let agent!: ReactLoopAgent
  343. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  344. agent = inner.agentLoop.create(AgentId('scoped'), { provider: 'mock', model: 'mock' })
  345. }, { inject: ['agentLoop'] }))
  346. send(agent, 'go')
  347. await new Promise(r => setTimeout(r, 30))
  348. expect(agent.status).toBe('running')
  349. const idle = agent.whenIdle() // queued while running
  350. await fiber.dispose() // tears the fiber down (drops agent listeners)
  351. await idle // must resolve, not hang
  352. expect(agent.status).toBe('disposed')
  353. })
  354. it('whenIdle() on a disposed agent awaits the loop exit (done), not just the status flip', async () => {
  355. // Disposed status is emitted before the driver unwinds. `whenIdle()` must chain `done` so it
  356. // resolves only after true loop exit.
  357. const adapter = new MockAdapter(['hang'])
  358. const ctx = await harness(adapter)
  359. let agent!: ReactLoopAgent
  360. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  361. agent = inner.agentLoop.create(AgentId('scoped'), { provider: 'mock', model: 'mock' })
  362. }, { inject: ['agentLoop'] }))
  363. send(agent, 'go')
  364. await new Promise(r => setTimeout(r, 30))
  365. let doneResolved = false
  366. void agent.done.then(() => { doneResolved = true })
  367. await fiber.dispose() // sets status disposed, aborts, drains the loop
  368. expect(agent.status).toBe('disposed')
  369. // whenIdle() must not resolve before `done` has — chaining `done` is the
  370. // quiescence guarantee. By here dispose() awaited the loop, so done is
  371. // settled; whenIdle resolves and done is observed resolved.
  372. await agent.whenIdle()
  373. expect(doneResolved).toBe(true)
  374. })
  375. it('contains a throwing agent/status listener on the running transition', async () => {
  376. const adapter = new MockAdapter([textResponse('ok')])
  377. const ctx = await harness(adapter)
  378. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  379. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  380. ctx.on('agent/status', (_subject, status) => {
  381. if (status === 'running') throw new Error('bad running listener')
  382. })
  383. send(agent, 'go')
  384. await agent.whenIdle()
  385. expect(adapter.requests).toHaveLength(1)
  386. expect(agent.status).toBe('idle')
  387. expect(warn).toHaveBeenCalledWith(expect.stringContaining('agent event "agent/status" listener threw'))
  388. warn.mockRestore()
  389. })
  390. it('contains a throwing agent/status listener on the idle transition', async () => {
  391. const adapter = new MockAdapter([textResponse('ok')])
  392. const ctx = await harness(adapter)
  393. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  394. const agent = ctx.agentLoop.create(AgentId('a1'), { provider: 'mock', model: 'mock' })
  395. ctx.on('agent/status', (_subject, status) => {
  396. if (status === 'idle') throw new Error('bad idle listener')
  397. })
  398. send(agent, 'go')
  399. await agent.whenIdle()
  400. expect(adapter.requests).toHaveLength(1)
  401. expect(agent.status).toBe('idle')
  402. expect(warn).toHaveBeenCalledWith(expect.stringContaining('agent event "agent/status" listener threw'))
  403. warn.mockRestore()
  404. })
  405. })