agent.spec.ts 14 KB

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