subagent-inprocess.spec.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. import { describe, expect, it } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { type Agent } from '@deepseek-ai/dsh-agent'
  4. import { SessionId } from '@deepseek-ai/dsh-session'
  5. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  6. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  7. import InvariantService from '@deepseek-ai/dsh-invariants'
  8. import * as SessionInvariant from '@deepseek-ai/dsh-session/invariant'
  9. import * as AgentInvariant from '@deepseek-ai/dsh-agent/invariant'
  10. import * as AgentLoopInvariant from '@deepseek-ai/dsh-agent-loop/invariant'
  11. import SubagentService from '@deepseek-ai/dsh-subagent'
  12. import { maxTokensResponse, MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
  13. import { startInProcessRun } from '../src/index.ts'
  14. type Script = ConstructorParameters<typeof MockAdapter>[0]
  15. async function mountInvariants(ctx: Context): Promise<void> {
  16. await ctx.plugin(InvariantService)
  17. await ctx.plugin(SessionInvariant)
  18. await ctx.plugin(AgentInvariant)
  19. await ctx.plugin(AgentLoopInvariant)
  20. }
  21. async function setup(script: Script) {
  22. const ctx = new Context()
  23. await mountAgentLoopTestDependencies(ctx)
  24. await mountInvariants(ctx)
  25. await ctx.plugin(AgentLoop, { agents: [] })
  26. await ctx.plugin(SubagentService)
  27. const adapter = new MockAdapter(script)
  28. ctx.llm.registerAdapter(['mock'], adapter)
  29. const parent = ctx.agentLoop.create(SessionId('parent'), { provider: 'mock', model: 'mock' })
  30. return { ctx, parent, adapter }
  31. }
  32. function request(parent: Agent, signal = new AbortController().signal) {
  33. return { prompt: [{ type: 'text' as const, text: 'child task' }], parent, signal }
  34. }
  35. function text(blocks: readonly { type: string; text?: string }[]): string {
  36. return blocks.filter(block => block.type === 'text').map(block => block.text).join('')
  37. }
  38. describe('startInProcessRun', () => {
  39. it('returns only after publication, drives a fresh child, and disposes it', async () => {
  40. const { ctx, parent } = await setup([textResponse('driver answer')])
  41. const run = await startInProcessRun(request(parent), {})
  42. expect(ctx.agents.get(run.id)).toBeDefined()
  43. const result = await run.result
  44. expect(result.stopReason).toBe('completed')
  45. expect(text(result.output)).toBe('driver answer')
  46. expect(ctx.agents.get(run.id)!.options.subagentDepth).toBe(1)
  47. await run.dispose()
  48. await run.dispose()
  49. expect(ctx.agents.get(run.id)).toBeUndefined()
  50. })
  51. it('reports the message-turn outcome when a later non-message turn completes during flush', async () => {
  52. const { ctx, parent } = await setup([maxTokensResponse('partial answer')])
  53. let injected = false
  54. ctx.on('session/flush', (session) => {
  55. if (injected || session.header.parentSession === undefined) return
  56. const lastEnd = session.events.findLast(event => event.type === 'turn/end')
  57. if (lastEnd?.type !== 'turn/end' || lastEnd.data.reason.kind !== 'max-tokens') return
  58. injected = true
  59. const turn = lastEnd.data.turn + 1
  60. session.append('turn/start', {
  61. turn,
  62. trigger: { kind: 'injection', source: { kind: 'plugin', plugin: 'late-metadata' } },
  63. })
  64. session.append('user/message', {
  65. content: [{ type: 'text', text: 'late metadata' }],
  66. source: { kind: 'plugin', plugin: 'late-metadata' },
  67. }, { surfaceOp: 'append' })
  68. session.append('turn/end', { turn, reason: { kind: 'completed' } })
  69. })
  70. const run = await startInProcessRun(request(parent), {})
  71. const result = await run.result
  72. const child = ctx.agents.get(run.id)!
  73. expect(injected).toBe(true)
  74. expect(child.session.events.findLast(event => event.type === 'turn/end'))
  75. .toMatchObject({ data: { reason: { kind: 'completed' } } })
  76. expect(result.stopReason).toBe('max-tokens')
  77. await run.dispose()
  78. })
  79. it('seeds a forked child but reads only the child-owned output', async () => {
  80. const { ctx, parent } = await setup([textResponse('parent answer'), textResponse('child answer')])
  81. parent.followup({ content: [{ type: 'text', text: 'parent question' }], source: { kind: 'user' } })
  82. await parent.whenIdle()
  83. const seed = parent.session.events.slice()
  84. const run = await startInProcessRun(request(parent), { seed })
  85. const result = await run.result
  86. expect(text(result.output)).toBe('child answer')
  87. const child = ctx.agents.get(run.id)!
  88. expect(child.session.header.seedLength).toBe(seed.length)
  89. expect(child.session.events.slice(0, seed.length)).toEqual(seed)
  90. await run.dispose()
  91. })
  92. it('persists the child depth in its session header', async () => {
  93. const { ctx, parent } = await setup([textResponse('child answer')])
  94. const run = await startInProcessRun(request(parent), {})
  95. await run.result
  96. // The recursion budget is durable session data, not only runtime options —
  97. // a depth that lived only in AgentOptions would reset to 0 on resume.
  98. expect(ctx.agents.get(run.id)!.session.header.delegationDepth).toBe(1)
  99. await run.dispose()
  100. })
  101. it('counts a RESUMED child by its persisted header depth, not the absent runtime depth', async () => {
  102. // Resume rebuilds runtime options, so the durable header must keep this
  103. // depth-1 child from delegating as though it were top-level.
  104. const { ctx } = await setup([textResponse('unused')])
  105. const resumed = (await ctx.agents.create({
  106. sessionId: SessionId('resumed-child'),
  107. meta: { parentSession: SessionId('root'), delegationDepth: 1 },
  108. agentOptions: { provider: 'mock', model: 'mock' },
  109. signal: new AbortController().signal,
  110. })).agent
  111. await expect(startInProcessRun({ ...request(resumed), maxDepth: 1 }, {}))
  112. .rejects.toMatchObject({ name: 'SubagentDepthError', attemptedDepth: 2, maxDepth: 1 })
  113. })
  114. it('lets runtime options deepen but never lower the persisted depth', async () => {
  115. const { ctx } = await setup([textResponse('unused')])
  116. const parent = (await ctx.agents.create({
  117. sessionId: SessionId('deep-parent'),
  118. meta: { delegationDepth: 2 },
  119. agentOptions: { provider: 'mock', model: 'mock', subagentDepth: 1 },
  120. signal: new AbortController().signal,
  121. })).agent
  122. // Persisted 2 vs runtime 1: the child is depth 3, so maxDepth 2 rejects.
  123. await expect(startInProcessRun({ ...request(parent), maxDepth: 2 }, {}))
  124. .rejects.toMatchObject({ name: 'SubagentDepthError', attemptedDepth: 3, maxDepth: 2 })
  125. })
  126. it('rejects invalid and exceeded depth before publication', async () => {
  127. const { parent } = await setup([])
  128. await expect(startInProcessRun({ ...request(parent), maxDepth: -1 }, {}))
  129. .rejects.toThrow('non-negative safe integer')
  130. await expect(startInProcessRun({ ...request(parent), maxDepth: 0 }, {}))
  131. .rejects.toMatchObject({ name: 'SubagentDepthError' })
  132. for (const value of [Number.NaN, 1.5, -1, -0, Number.MAX_SAFE_INTEGER + 1]) {
  133. const malformed = { options: { subagentDepth: value }, session: { header: {} } } as unknown as Agent
  134. await expect(startInProcessRun(request(malformed), {}))
  135. .rejects.toThrow('agent subagentDepth must be a non-negative safe integer')
  136. }
  137. const maxParent = { options: { subagentDepth: Number.MAX_SAFE_INTEGER }, session: { header: {} } } as unknown as Agent
  138. await expect(startInProcessRun(request(maxParent), {})).rejects.toBeInstanceOf(RangeError)
  139. })
  140. it('rejects an already-aborted request without publishing a child', async () => {
  141. const { ctx, parent } = await setup([])
  142. const beforeAgents = ctx.agents.list().length
  143. const beforeSessions = ctx.sessions.list().length
  144. const controller = new AbortController()
  145. controller.abort('too late')
  146. await expect(startInProcessRun(request(parent, controller.signal), {}))
  147. .rejects.toThrow('aborted before child publication')
  148. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  149. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  150. })
  151. it('uses the request signal after publication and dispose as cancellation paths', async () => {
  152. const { parent, adapter } = await setup(['hang', 'hang'])
  153. const controller = new AbortController()
  154. const signalled = await startInProcessRun(request(parent, controller.signal), {})
  155. await new Promise(resolve => setTimeout(resolve, 30))
  156. controller.abort('stop child')
  157. await expect(signalled.result).resolves.toMatchObject({ stopReason: 'aborted' })
  158. expect(adapter.requests[0]?.signal?.reason).toEqual({ kind: 'parent' })
  159. const child = parent.ctx.agents.get(signalled.id)
  160. const turnEnd = child?.session.events.findLast(event => event.type === 'turn/end')
  161. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  162. await signalled.dispose()
  163. const disposed = await startInProcessRun(request(parent), {})
  164. await new Promise(resolve => setTimeout(resolve, 30))
  165. await disposed.dispose()
  166. await expect(disposed.result).resolves.toMatchObject({ stopReason: 'aborted' })
  167. })
  168. it('cleans a failed unpublished setup before rejecting', async () => {
  169. const { ctx, parent } = await setup([])
  170. const beforeAgents = ctx.agents.list().length
  171. const beforeSessions = ctx.sessions.list().length
  172. await expect(startInProcessRun({
  173. ...request(parent),
  174. toolFilter: { deny: ['unknown-tool'] },
  175. }, {})).rejects.toThrow('unknown global tool')
  176. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  177. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  178. })
  179. it('closes the abort handoff after the factory detaches its creation listener', async () => {
  180. const { ctx, parent } = await setup([])
  181. const controller = new AbortController()
  182. const beforeAgents = ctx.agents.list().length
  183. const beforeSessions = ctx.sessions.list().length
  184. const parentWithAbortAtHandoff = {
  185. options: parent.options,
  186. session: parent.session,
  187. ctx: {
  188. agents: {
  189. create: async (options: Parameters<typeof ctx.agents.create>[0]) => {
  190. const handle = await ctx.agents.create(options)
  191. // `create()` has detached its creation-only listener, but the
  192. // provider continuation has not installed its live-run listener.
  193. controller.abort('handoff race')
  194. return handle
  195. },
  196. },
  197. },
  198. } as unknown as Agent
  199. await expect(startInProcessRun(request(parentWithAbortAtHandoff, controller.signal), {}))
  200. .rejects.toThrow('aborted before child publication')
  201. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  202. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  203. })
  204. })