agent.spec.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272
  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 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 AgentLoop, { LoopAgent } from '@deepseek-ai/dsh-agent-loop'
  10. import { MockAdapter, textResponse } from './mock-adapter.ts'
  11. async function harness(adapter: MockAdapter) {
  12. const ctx = new Context()
  13. await ctx.plugin(LlmService)
  14. await ctx.plugin(SessionStore)
  15. await ctx.plugin(SystemPrompt)
  16. await ctx.plugin(ToolRegistry)
  17. await ctx.plugin(AgentRegistry)
  18. await ctx.plugin(AgentLoop, { agents: [] })
  19. ctx.llm.registerAdapter(['mock'], adapter)
  20. return ctx
  21. }
  22. function waitForIdle(ctx: Context, agent: LoopAgent): Promise<void> {
  23. return new Promise((resolve) => {
  24. const dispose = ctx.on('agent/status', (subject, status) => {
  25. if (subject === agent && status === 'idle') {
  26. dispose()
  27. resolve()
  28. }
  29. })
  30. })
  31. }
  32. function send(agent: LoopAgent, text: string) {
  33. agent.send([{ type: 'text', text }])
  34. }
  35. describe('LoopAgent', () => {
  36. it('send() throws after disposal', async () => {
  37. const adapter = new MockAdapter(['hang'])
  38. const ctx = await harness(adapter)
  39. let agent!: LoopAgent
  40. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  41. agent = inner.agentLoop.create('scoped', { model: 'mock' })
  42. }, { inject: ['agentLoop'] }))
  43. send(agent, 'go')
  44. await new Promise(r => setTimeout(r, 30))
  45. await fiber.dispose()
  46. await agent.done
  47. expect(() => { agent.send([{ type: 'text', text: 'too late' }]) }).toThrow('disposed')
  48. })
  49. it('steer() throws after disposal', async () => {
  50. const adapter = new MockAdapter(['hang'])
  51. const ctx = await harness(adapter)
  52. let agent!: LoopAgent
  53. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  54. agent = inner.agentLoop.create('scoped', { model: 'mock' })
  55. }, { inject: ['agentLoop'] }))
  56. send(agent, 'go')
  57. await new Promise(r => setTimeout(r, 30))
  58. await fiber.dispose()
  59. await agent.done
  60. expect(() => { agent.steer([{ type: 'text', text: 'too late' }]) }).toThrow('disposed')
  61. })
  62. it('inject() throws after disposal', async () => {
  63. const adapter = new MockAdapter(['hang'])
  64. const ctx = await harness(adapter)
  65. let agent!: LoopAgent
  66. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  67. agent = inner.agentLoop.create('scoped', { model: 'mock' })
  68. }, { inject: ['agentLoop'] }))
  69. send(agent, 'go')
  70. await new Promise(r => setTimeout(r, 30))
  71. await fiber.dispose()
  72. await agent.done
  73. expect(() => { agent.inject([{ type: 'text', text: 'too late' }]) }).toThrow('disposed')
  74. })
  75. it('inject() decides enclosure from the LOG (open turn), not agent status', async () => {
  76. const adapter = new MockAdapter([textResponse('ok')])
  77. const ctx = await harness(adapter)
  78. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  79. // Simulate an OPEN turn in the log while the agent is idle (status is not a
  80. // reliable open-turn signal). inject must append into that open turn, NOT
  81. // wrap a new one.
  82. agent.session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  83. agent.inject([{ type: 'text', text: 'mid' }], { source: { kind: 'plugin', plugin: 'p' } })
  84. expect(agent.session.events.filter(e => e.type === 'turn/start')).toHaveLength(1)
  85. expect(agent.session.events.at(-1)!.type).toBe('context/message')
  86. // Close the turn; now inject must wrap its own one-shot injection turn.
  87. agent.session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  88. agent.inject([{ type: 'text', text: 'after' }], { source: { kind: 'plugin', plugin: 'p' } })
  89. const starts = agent.session.events.filter(e => e.type === 'turn/start')
  90. expect(starts).toHaveLength(2)
  91. const last = starts[1]!
  92. expect(last.type === 'turn/start' && last.data.trigger.kind).toBe('injection')
  93. expect(agent.session.events.at(-1)!.type).toBe('turn/end') // turn-enclosed
  94. })
  95. it('idle inject() contains a failing flush (logs, does not throw into the caller)', async () => {
  96. const adapter = new MockAdapter([textResponse('ok')])
  97. const ctx = await harness(adapter)
  98. // A persistence-like listener whose flush rejects.
  99. ctx.on('session/flush', () => { throw new Error('disk gone') })
  100. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  101. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  102. // inject() is synchronous and fires a fire-and-forget flush; a rejecting
  103. // flush must be contained (logged), never thrown into the caller.
  104. expect(() => { agent.inject([{ type: 'text', text: 'notice' }], { source: { kind: 'plugin', plugin: 'p' } }) }).not.toThrow()
  105. await new Promise(r => setTimeout(r, 20)) // let the contained flush settle
  106. expect(warn).toHaveBeenCalledWith(expect.stringContaining('flush after idle injection failed'))
  107. warn.mockRestore()
  108. })
  109. it('idle inject() closes its one-shot turn AND still checkpoints even if the append throws', async () => {
  110. const adapter = new MockAdapter([textResponse('ok')])
  111. const ctx = await harness(adapter)
  112. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  113. let flushes = 0
  114. ctx.on('session/flush', () => { flushes += 1 })
  115. // Non-serializable injected content makes Session.append throw AFTER
  116. // turn/start was recorded. The turn/end must still be appended (finally),
  117. // AND the durability checkpoint must still fire — the balanced turn is in
  118. // memory and a crash before the next turn/dispose would otherwise lose it.
  119. expect(() => {
  120. agent.inject([{ type: 'text', text: 'x', bad: 1n } as never], { source: { kind: 'plugin', plugin: 'p' } })
  121. }).toThrow(/non-JSON-serializable/)
  122. const types = agent.session.events.map(e => e.type)
  123. expect(types).toEqual(['turn/start', 'turn/end']) // balanced, no open turn
  124. await new Promise(r => setTimeout(r, 10)) // let the fire-and-forget flush run
  125. expect(flushes).toBe(1) // checkpoint fired despite the throw
  126. })
  127. it('idle inject() still checkpoints when a listener throws on the synthetic turn/end', async () => {
  128. const adapter = new MockAdapter([textResponse('ok')])
  129. const ctx = await harness(adapter)
  130. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  131. let flushes = 0
  132. ctx.on('session/flush', () => { flushes += 1 })
  133. // A session/event listener that throws on the synthetic turn/end. Append
  134. // pushes before notifying, so turn/end is in the log (turn balanced) but the
  135. // throw must NOT skip the durability checkpoint — the flush decision is made
  136. // from the log, not a flag set after the (throwing) append.
  137. let threw = false
  138. ctx.on('session/event', (_s, event) => {
  139. if (!threw && event.type === 'turn/end') { threw = true; throw new Error('boom turn/end') }
  140. })
  141. expect(() => { agent.inject([{ type: 'text', text: 'notice' }], { source: { kind: 'plugin', plugin: 'p' } }) }).not.toThrow()
  142. const types = agent.session.events.map(e => e.type)
  143. expect(types).toEqual(['turn/start', 'context/message', 'turn/end']) // balanced
  144. await new Promise(r => setTimeout(r, 10))
  145. expect(flushes).toBe(1) // checkpoint fired despite the throwing turn/end listener
  146. })
  147. it('idle inject() reports a failing flush via agent/error (step 0) AND the logger', async () => {
  148. const adapter = new MockAdapter([textResponse('ok')])
  149. const ctx = await harness(adapter)
  150. // A non-Error rejection exercises the String() normalization branch.
  151. ctx.on('session/flush', () => { throw 'disk gone' })
  152. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
  153. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  154. const errors: { turn: number; step: number; message: string }[] = []
  155. ctx.on('agent/error', (_a, turn, step, error) => void errors.push({ turn, step, message: error.message }))
  156. agent.inject([{ type: 'text', text: 'notice' }], { source: { kind: 'plugin', plugin: 'p' } })
  157. await new Promise(r => setTimeout(r, 20)) // let the contained flush settle
  158. // Reported via agent/error (step 0 — the idle-injection convention) so
  159. // plugins monitoring agent/error see idle-injection persistence failures,
  160. // mirroring the loop's post-turn/end flush path. A non-Error throw is
  161. // normalized to an Error.
  162. expect(errors).toEqual([{ turn: 1, step: 0, message: 'disk gone' }])
  163. expect(warn).toHaveBeenCalledWith(expect.stringContaining('flush after idle injection failed'))
  164. warn.mockRestore()
  165. })
  166. it('idle inject() with a non-serializable source opens no turn (nothing to close)', async () => {
  167. const adapter = new MockAdapter([textResponse('ok')])
  168. const ctx = await harness(adapter)
  169. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  170. // A non-serializable source makes the turn/start append throw BEFORE the
  171. // event is pushed (Session.append validates before push), so NO turn opens.
  172. // The finally's isTurnOpen() guard sees no open turn and appends nothing —
  173. // the log stays empty, not left with a dangling turn/start.
  174. expect(() => {
  175. agent.inject([{ type: 'text', text: 'x' }], { source: { kind: 'plugin', plugin: 'p', bad: 1n } as never })
  176. }).toThrow(/non-JSON-serializable/)
  177. expect(agent.session.events).toHaveLength(0)
  178. })
  179. it('steer() when idle falls through to send() and starts a turn', async () => {
  180. const adapter = new MockAdapter([textResponse('ok')])
  181. const ctx = await harness(adapter)
  182. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  183. // steer while idle delegates to send
  184. agent.steer([{ type: 'text', text: 'steer idle' }], { source: { kind: 'plugin', plugin: 'test' } })
  185. await waitForIdle(ctx, agent)
  186. // The message was recorded as a user-level message (send path)
  187. expect(agent.session.events.some(e => e.type === 'user/message')).toBe(true)
  188. expect(adapter.requests).toHaveLength(1)
  189. })
  190. it('disposer is idempotent (double-stop)', async () => {
  191. // Create a bare LoopAgent and call start() directly to get the disposer.
  192. // Then call it twice — the second call hits the early-return branch.
  193. const ctx = new Context()
  194. await ctx.plugin(SessionStore)
  195. const session = ctx.sessions.create('test')
  196. const agent = new LoopAgent(ctx, AgentId('bare'), { model: 'mock' }, session)
  197. // Start the loop to get the disposer; the agent waits for messages
  198. // (idle, never-resolving cancel), so it will stay idle.
  199. const dispose = agent.start()
  200. // First dispose
  201. dispose()
  202. expect(agent.status).toBe('disposed')
  203. // Second dispose — idempotent, no throw
  204. expect(() => { dispose() }).not.toThrow()
  205. expect(agent.status).toBe('disposed')
  206. })
  207. it('setting the same status does not emit agent/status again', async () => {
  208. const adapter = new MockAdapter([textResponse('ok')])
  209. const ctx = await harness(adapter)
  210. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  211. const statuses: string[] = []
  212. ctx.on('agent/status', (subject, status) => {
  213. if (subject === agent) statuses.push(status)
  214. })
  215. send(agent, 'hi')
  216. await waitForIdle(ctx, agent)
  217. // After the turn, agent is idle. Send again to trigger another attempt
  218. // to go idle — but it's already idle, so no emission.
  219. const idleTransitionCount = statuses.filter(s => s === 'idle').length
  220. expect(idleTransitionCount).toBe(1) // only the final transition from running
  221. })
  222. it('abort() resolves reason to "aborted" when no reason provided', async () => {
  223. const adapter = new MockAdapter(['hang'])
  224. const ctx = await harness(adapter)
  225. const agent = ctx.agentLoop.create('a1', { model: 'mock' })
  226. const reasons: { kind: string; reason?: string }[] = []
  227. ctx.on('agent/turn-end', (_agent, _turn, reason) => void reasons.push(reason))
  228. send(agent, 'go')
  229. await new Promise(r => setTimeout(r, 30))
  230. agent.abort() // no reason string
  231. await waitForIdle(ctx, agent)
  232. expect(reasons[0]).toMatchObject({ kind: 'aborted', reason: 'aborted' })
  233. })
  234. })