agent.spec.ts 14 KB

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