agent.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277
  1. import { describe, expect, expectTypeOf, it } from 'vitest'
  2. import { Context, Service, symbols } from 'cordis'
  3. import type { Events } from 'cordis'
  4. import { Session, SessionId } from '@deepseek-ai/dsh-session'
  5. import AgentRegistry, {
  6. agentEvents,
  7. } from '@deepseek-ai/dsh-agent'
  8. import type {
  9. Agent,
  10. AgentCancelCause,
  11. AgentFactory,
  12. CreateAgentOptions,
  13. ResumeAgentOptions,
  14. } from '@deepseek-ai/dsh-agent'
  15. function stubAgent(rawId: string, overrides: Partial<Agent> = {}): Agent {
  16. const id = SessionId(rawId)
  17. const agent: Agent = {
  18. id,
  19. options: {},
  20. session: Session.create(id),
  21. status: 'idle',
  22. acceptsNextStep: false,
  23. ctx: new Context(),
  24. send: () => {},
  25. updateInbox: () => 'not-found',
  26. followup: () => {},
  27. steer: () => ({ outcome: Promise.resolve({ status: 'rejected' as const }) }),
  28. inject: () => {},
  29. reserveTurnAdmission: () => undefined,
  30. cancel() {},
  31. whenIdle() { return Promise.resolve() },
  32. }
  33. return Object.assign(agent, overrides)
  34. }
  35. describe('AgentRegistry', () => {
  36. it('registers exact entries, emits lifecycle events, and unregisters on owner disposal', async () => {
  37. const ctx = new Context()
  38. await ctx.plugin(AgentRegistry)
  39. const lifecycle: string[] = []
  40. ctx.on('agent/created', agent => void lifecycle.push(`created:${agent.id}`))
  41. ctx.on('agent/disposed', agent => void lifecycle.push(`disposed:${agent.id}`))
  42. const agent = stubAgent('a1')
  43. const dispose = ctx.agents.register(agent)
  44. expect(ctx.agents.get(agent.id)).toBe(agent)
  45. expect(ctx.agents.list()).toEqual([agent])
  46. expect(ctx.agents.roots()).toEqual([agent])
  47. expect(() => ctx.agents.register(stubAgent('a1'))).toThrow(/already registered/)
  48. dispose()
  49. expect(ctx.agents.get(agent.id)).toBeUndefined()
  50. expect(lifecycle).toEqual(['created:a1', 'disposed:a1'])
  51. })
  52. it('rejects an agent whose registry and session identities differ', async () => {
  53. const ctx = new Context()
  54. await ctx.plugin(AgentRegistry)
  55. const agent = stubAgent('agent-id', { session: Session.create(SessionId('session-id')) })
  56. expect(() => ctx.agents.enter(agent, undefined))
  57. .toThrow('agent id "agent-id" does not match session id "session-id"')
  58. expect(ctx.agents.list()).toEqual([])
  59. })
  60. it('tracks runtime creator ownership separately from registry order', async () => {
  61. const ctx = new Context()
  62. await ctx.plugin(AgentRegistry)
  63. const root = stubAgent('root')
  64. const child = stubAgent('child')
  65. const detachRoot = ctx.agents.enter(root, undefined)
  66. ctx.agents.announce(root)
  67. const detachChild = ctx.agents.enter(child, root)
  68. ctx.agents.announce(child)
  69. expect(ctx.agents.list()).toEqual([root, child])
  70. expect(ctx.agents.roots()).toEqual([root])
  71. expect(ctx.agents.isOwnedBy(child.id, root)).toBe(true)
  72. expect(ctx.agents.isOwnedBy(root.id, root)).toBe(false)
  73. expect(ctx.agents.isOwnedBy(SessionId('missing'), root)).toBe(false)
  74. detachChild()
  75. expect(ctx.agents.isOwnedBy(child.id, root)).toBe(false)
  76. detachRoot()
  77. })
  78. it('rolls an entry back and pairs a partially delivered creation when a listener throws', async () => {
  79. const ctx = new Context()
  80. await ctx.plugin(AgentRegistry)
  81. const lifecycle: string[] = []
  82. ctx.on('agent/created', agent => void lifecycle.push(`created:${agent.id}`))
  83. ctx.on('agent/created', () => { throw new Error('creation veto') })
  84. ctx.on('agent/disposed', agent => void lifecycle.push(`disposed:${agent.id}`))
  85. expect(() => ctx.agents.register(stubAgent('vetoed'))).toThrow('creation veto')
  86. expect(ctx.agents.get(SessionId('vetoed'))).toBeUndefined()
  87. expect(lifecycle).toEqual(['created:vetoed', 'disposed:vetoed'])
  88. })
  89. it('contains asynchronous creation rejection and every disposal-listener failure', async () => {
  90. const ctx = new Context()
  91. await ctx.plugin(AgentRegistry)
  92. const warnings: string[] = []
  93. const heard: string[] = []
  94. ctx.logger.warn = ((message: unknown) => { warnings.push(String(message)) }) as typeof ctx.logger.warn
  95. ctx.on('agent/created', () => Promise.reject(new Error('created async')) as never)
  96. ctx.on('agent/disposed', () => { throw new Error('disposed sync') })
  97. ctx.on('agent/disposed', () => Promise.reject(new Error('disposed async')) as never)
  98. ctx.on('agent/disposed', agent => void heard.push(agent.id))
  99. const dispose = ctx.agents.register(stubAgent('contained'))
  100. await Promise.resolve()
  101. dispose()
  102. await Promise.resolve()
  103. expect(heard).toEqual(['contained'])
  104. expect(warnings).toEqual([
  105. 'agent "contained": agent/created listener rejected: Error: created async',
  106. 'agent "contained": agent/disposed listener threw: Error: disposed sync',
  107. 'agent "contained": agent/disposed listener rejected: Error: disposed async',
  108. ])
  109. })
  110. it('separates entry from announcement and stale/idempotent detach cannot remove a replacement', async () => {
  111. const ctx = new Context()
  112. await ctx.plugin(AgentRegistry)
  113. const lifecycle: string[] = []
  114. ctx.on('agent/created', agent => void lifecycle.push(`created:${agent.id}`))
  115. ctx.on('agent/disposed', agent => void lifecycle.push(`disposed:${agent.id}`))
  116. const first = stubAgent('split')
  117. const detachFirst = ctx.agents.enter(first, undefined)
  118. expect(lifecycle).toEqual([])
  119. ctx.agents.announce(first)
  120. expect(() => { ctx.agents.announce(first) }).toThrow(/already announced/)
  121. detachFirst()
  122. detachFirst()
  123. const replacement = stubAgent('split')
  124. const detachReplacement = ctx.agents.enter(replacement, undefined)
  125. detachFirst()
  126. expect(ctx.agents.get(replacement.id)).toBe(replacement)
  127. expect(() => { ctx.agents.announce(first) }).toThrow(/not live/)
  128. detachReplacement()
  129. expect(lifecycle).toEqual(['created:split', 'disposed:split'])
  130. })
  131. it('defers detach requested by a creation listener until that dispatch unwinds', async () => {
  132. const ctx = new Context()
  133. await ctx.plugin(AgentRegistry)
  134. const order: string[] = []
  135. const agent = stubAgent('reentrant')
  136. ctx.on('agent/created', () => {
  137. order.push(`first:${ctx.agents.get(agent.id) === agent}`)
  138. detach()
  139. order.push(`after-detach:${ctx.agents.get(agent.id) === agent}`)
  140. })
  141. ctx.on('agent/created', () => void order.push(`second:${ctx.agents.get(agent.id) === agent}`))
  142. ctx.on('agent/disposed', () => void order.push('disposed'))
  143. const detach = ctx.agents.enter(agent, undefined)
  144. ctx.agents.announce(agent)
  145. expect(order).toEqual(['first:true', 'after-detach:true', 'second:true', 'disposed'])
  146. expect(ctx.agents.get(agent.id)).toBeUndefined()
  147. })
  148. })
  149. describe('agentEvents()', () => {
  150. it('contains each synchronous throw and returned-promise rejection', async () => {
  151. const ctx = new Context()
  152. const warnings: string[] = []
  153. const heard: string[] = []
  154. ctx.logger.warn = ((message: unknown) => { warnings.push(String(message)) }) as typeof ctx.logger.warn
  155. const agent = stubAgent('event')
  156. ctx.on('agent/status', () => { throw new Error('sync listener') })
  157. ctx.on('agent/status', () => Promise.reject(new Error('async listener')) as never)
  158. ctx.on('agent/status', (_agent, status) => void heard.push(status))
  159. agentEvents(ctx, agent).emit('agent/status', 'running')
  160. await Promise.resolve()
  161. expect(heard).toEqual(['running'])
  162. expect(warnings).toEqual([
  163. 'agent event "agent/status" listener threw: Error: sync listener',
  164. 'agent event "agent/status" listener rejected: Error: async listener',
  165. ])
  166. })
  167. })
  168. describe('explicit cancellation contract', () => {
  169. it('exposes the closed typed cancellation cause at the Agent seam', () => {
  170. expectTypeOf<Parameters<Agent['cancel']>[0]>().toEqualTypeOf<AgentCancelCause>()
  171. expectTypeOf<Parameters<Events['agent/cancel-requested']>[1]>().toEqualTypeOf<AgentCancelCause>()
  172. })
  173. })
  174. describe('AgentRegistry factory seam', () => {
  175. function stubFactory() {
  176. const calls: {
  177. create: Array<{ ownerCtx: Context; options: CreateAgentOptions }>
  178. resume: Array<{ ownerCtx: Context; options: ResumeAgentOptions }>
  179. } = { create: [], resume: [] }
  180. const factory: AgentFactory = {
  181. async createAgent(ownerCtx, options) {
  182. calls.create.push({ ownerCtx, options })
  183. return { agent: stubAgent(options.sessionId), dispose: () => Promise.resolve() }
  184. },
  185. async resume(ownerCtx, options) {
  186. calls.resume.push({ ownerCtx, options })
  187. return { agent: stubAgent(options.resumeSessionId), dispose: () => Promise.resolve() }
  188. },
  189. }
  190. return { factory, calls }
  191. }
  192. it('requires a factory and delegates through the calling context', async () => {
  193. const ctx = new Context()
  194. await ctx.plugin(AgentRegistry)
  195. await expect(ctx.agents.create({ sessionId: SessionId('s') })).rejects.toThrow(/no agent factory/)
  196. const { factory, calls } = stubFactory()
  197. ctx.agents.setFactory(factory)
  198. let callerFiber: Context['fiber'] | undefined
  199. await ctx.plugin(Object.assign(async (inner: Context) => {
  200. callerFiber = inner.fiber
  201. await inner.agents.create({ sessionId: SessionId('create-s') })
  202. await inner.agents.resume({ resumeSessionId: SessionId('resume-s') })
  203. }, { inject: ['agents'] }))
  204. expect(calls.create[0]?.ownerCtx.fiber).toBe(callerFiber)
  205. expect(calls.resume[0]?.ownerCtx.fiber).toBe(callerFiber)
  206. })
  207. it('rejects a second factory and clears the slot with its owner (HMR)', async () => {
  208. const ctx = new Context()
  209. await ctx.plugin(AgentRegistry)
  210. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  211. inner.agents.setFactory(stubFactory().factory)
  212. expect(() => inner.agents.setFactory(stubFactory().factory)).toThrow(/already registered/)
  213. }, { inject: ['agents'] }))
  214. await expect(ctx.agents.create({ sessionId: SessionId('before-s') })).resolves.toBeDefined()
  215. await owner.dispose()
  216. await expect(ctx.agents.create({ sessionId: SessionId('after-s') })).rejects.toThrow(/no agent factory/)
  217. })
  218. it('canonicalizes an already traced Service before tracing it for the caller', async () => {
  219. const ctx = new Context()
  220. await ctx.plugin(AgentRegistry)
  221. const states = new WeakMap<object, string[]>()
  222. class TracedFactory extends Service implements AgentFactory {
  223. constructor(inner: Context) {
  224. super(inner, 'tracedFactory')
  225. states.set(this, [])
  226. }
  227. private calls(): string[] {
  228. const original = (this as unknown as { [symbols.original]?: TracedFactory })[symbols.original] ?? this
  229. const calls = states.get(original)
  230. if (calls === undefined) throw new Error('factory receiver was not canonicalized')
  231. return calls
  232. }
  233. async createAgent(_ownerCtx: Context, options: CreateAgentOptions) {
  234. this.calls().push('create')
  235. return { agent: stubAgent(options.sessionId), dispose: () => Promise.resolve() }
  236. }
  237. async resume(_ownerCtx: Context, options: ResumeAgentOptions) {
  238. this.calls().push('resume')
  239. return { agent: stubAgent(options.resumeSessionId), dispose: () => Promise.resolve() }
  240. }
  241. }
  242. await ctx.plugin(TracedFactory)
  243. const traced = (ctx as Context & { tracedFactory: TracedFactory }).tracedFactory
  244. ctx.agents.setFactory(traced)
  245. await ctx.agents.create({ sessionId: SessionId('create-s') })
  246. await ctx.agents.resume({ resumeSessionId: SessionId('resume-s') })
  247. const raw = (traced as unknown as { [symbols.original]?: TracedFactory })[symbols.original]
  248. expect(states.get(raw!)).toEqual(['create', 'resume'])
  249. })
  250. })