subagent-inprocess.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342
  1. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  2. import { describe, expect, it } from 'vitest'
  3. import { Context } from 'cordis'
  4. import { type Agent, type AgentOptions } from '@deepseek-ai/dsh-agent'
  5. import { SessionId } from '@deepseek-ai/dsh-session'
  6. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  7. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  8. import InvariantService from '@deepseek-ai/dsh-invariants'
  9. import * as SessionInvariant from '@deepseek-ai/dsh-session/invariant'
  10. import * as AgentInvariant from '@deepseek-ai/dsh-agent/invariant'
  11. import * as AgentLoopInvariant from '@deepseek-ai/dsh-agent-loop/invariant'
  12. import SubagentService, { snapshotSubagentDescriptor } from '@deepseek-ai/dsh-subagent'
  13. import { maxTokensResponse, MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
  14. import { startInProcessRun } from '../src/index.ts'
  15. type Script = ConstructorParameters<typeof MockAdapter>[0]
  16. async function mountInvariants(ctx: Context): Promise<void> {
  17. await ctx.plugin(InvariantService)
  18. await ctx.plugin(SessionInvariant)
  19. await ctx.plugin(AgentInvariant)
  20. await ctx.plugin(AgentLoopInvariant)
  21. }
  22. async function setup(script: Script, parentOptions: Partial<AgentOptions> = {}) {
  23. const ctx = new Context()
  24. await mountAgentLoopTestDependencies(ctx)
  25. await mountInvariants(ctx)
  26. await ctx.plugin(AgentLoop, { agents: [] })
  27. await ctx.plugin(SubagentService)
  28. const adapter = new MockAdapter(script)
  29. ctx.llm.registerAdapter(['mock'], adapter)
  30. const parent = ctx.agentLoop.create(SessionId('parent'), { provider: 'mock', model: 'mock', ...parentOptions })
  31. return { ctx, parent, adapter }
  32. }
  33. function request(parent: Agent, signal = new AbortController().signal) {
  34. return {
  35. label: 'child task',
  36. prompt: [{ type: 'text' as const, text: 'child task' }],
  37. parent,
  38. signal,
  39. descriptor: snapshotSubagentDescriptor({
  40. mode: 'one-shot',
  41. provider: 'test',
  42. label: 'child task',
  43. }),
  44. }
  45. }
  46. function text(blocks: readonly { type: string; text?: string }[]): string {
  47. return blocks.filter(block => block.type === 'text').map(block => block.text).join('')
  48. }
  49. describe('startInProcessRun', () => {
  50. it('returns only after publication, drives a fresh child, and disposes it', async () => {
  51. const { ctx, parent } = await setup([textResponse('driver answer')])
  52. const run = await startInProcessRun(request(parent), {})
  53. expect(ctx.agents.get(run.id)).toBeDefined()
  54. const result = await run.result
  55. expect(result.stopReason).toBe('completed')
  56. expect(text(result.output)).toBe('driver answer')
  57. expect(ctx.agents.get(run.id)!.options.subagentDepth).toBe(1)
  58. await run.dispose()
  59. await run.dispose()
  60. expect(ctx.agents.get(run.id)).toBeUndefined()
  61. })
  62. it('uses explicit child model selectors when the parent has none and preserves its cwd', async () => {
  63. const { ctx } = await setup([textResponse('driver answer')])
  64. const parent = ctx.agentLoop.create(SessionId('bare-parent'), {}, { cwd: '/workspace' })
  65. const run = await startInProcessRun({
  66. ...request(parent),
  67. agentOptions: { provider: 'mock', model: 'mock' },
  68. }, {})
  69. const child = ctx.agents.get(run.id)!
  70. expect(child.options).toMatchObject({ provider: 'mock', model: 'mock' })
  71. expect(child.session.header.cwd).toBe('/workspace')
  72. await expect(run.result).resolves.toMatchObject({ stopReason: 'completed' })
  73. await run.dispose()
  74. })
  75. it('does not add a final durability checkpoint to a foreground run', async () => {
  76. const { ctx, parent } = await setup([textResponse('driver answer')])
  77. let flushes = 0
  78. ctx.on('session/flush', (session) => {
  79. if (session.header.parentSession === undefined) return
  80. flushes++
  81. throw new Error('disk full')
  82. })
  83. const run = await startInProcessRun(request(parent), {})
  84. await expect(run.result).resolves.toMatchObject({ stopReason: 'completed' })
  85. expect(flushes).toBe(1)
  86. await run.dispose()
  87. })
  88. it('keeps published run and handle disposal failures on separate channels', async () => {
  89. const { ctx, parent } = await setup([])
  90. const runError = new Error('published run failed')
  91. const disposalError = new Error('published handle disposal failed')
  92. const beforeAgents = ctx.agents.list().length
  93. const beforeSessions = ctx.sessions.list().length
  94. const parentWithFailedDisposal = {
  95. options: parent.options,
  96. session: parent.session,
  97. ctx: {
  98. get: () => undefined,
  99. agents: {
  100. create: async (options: Parameters<typeof ctx.agents.create>[0]) => {
  101. const handle = await ctx.agents.create(options)
  102. handle.agent.followup = () => { throw runError }
  103. return {
  104. ...handle,
  105. dispose: async () => {
  106. await handle.dispose()
  107. throw disposalError
  108. },
  109. }
  110. },
  111. },
  112. },
  113. } as unknown as Agent
  114. const run = await startInProcessRun(request(parentWithFailedDisposal), {})
  115. expect(ctx.agents.get(run.id)).toBeDefined()
  116. await expect(run.result).rejects.toBe(runError)
  117. await expect(run.dispose()).rejects.toBe(disposalError)
  118. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  119. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  120. })
  121. it('reports the message-turn outcome when a later non-message turn completes during flush', async () => {
  122. const { ctx, parent } = await setup([maxTokensResponse('partial answer')])
  123. let injected = false
  124. ctx.on('session/flush', (session) => {
  125. if (injected || session.header.parentSession === undefined) return
  126. const lastEnd = session.events.findLast(event => event.type === 'turn/end')
  127. if (lastEnd?.type !== 'turn/end' || lastEnd.data.reason.kind !== 'max-tokens') return
  128. injected = true
  129. const turn = lastEnd.data.turn + 1
  130. session.append('turn/start', {
  131. turn,
  132. trigger: { kind: 'injection', source: { kind: 'plugin', plugin: 'late-metadata' } },
  133. })
  134. session.append('user/message', createUserMessage({
  135. content: [{ type: 'text', text: 'late metadata' }],
  136. source: { kind: 'plugin', plugin: 'late-metadata' },
  137. }), { surfaceOp: 'append' })
  138. session.append('turn/end', { turn, reason: { kind: 'completed' } })
  139. })
  140. const run = await startInProcessRun(request(parent), {})
  141. const result = await run.result
  142. const child = ctx.agents.get(run.id)!
  143. expect(injected).toBe(true)
  144. expect(child.session.events.findLast(event => event.type === 'turn/end'))
  145. .toMatchObject({ data: { reason: { kind: 'completed' } } })
  146. expect(result.stopReason).toBe('max-tokens')
  147. await run.dispose()
  148. })
  149. it('seeds a forked child but reads only the child-owned output', async () => {
  150. const { ctx, parent } = await setup([textResponse('parent answer'), textResponse('child answer')])
  151. parent.followup(createUserMessage({ content: [{ type: 'text', text: 'parent question' }], source: { kind: 'user' } }))
  152. await parent.whenIdle()
  153. const seed = parent.session.events.slice()
  154. const run = await startInProcessRun(request(parent), { seed })
  155. const result = await run.result
  156. expect(text(result.output)).toBe('child answer')
  157. const child = ctx.agents.get(run.id)!
  158. expect(child.session.header.seedLength).toBe(seed.length)
  159. expect(child.session.events.slice(0, seed.length)).toEqual(seed)
  160. await run.dispose()
  161. })
  162. it('persists the child origin and depth in its session header', async () => {
  163. const { ctx, parent } = await setup([textResponse('child answer')])
  164. const run = await startInProcessRun(request(parent), {})
  165. await run.result
  166. // The recursion budget is durable session data, not only runtime options —
  167. // a depth that lived only in AgentOptions would reset to 0 on resume.
  168. expect(ctx.agents.get(run.id)!.session.header).toMatchObject({
  169. origin: 'subagent',
  170. delegationDepth: 1,
  171. })
  172. await run.dispose()
  173. })
  174. it('inherits the parent output-token cap and accepts an explicit child override', async () => {
  175. const { ctx, parent, adapter } = await setup(
  176. [textResponse('inherited'), textResponse('overridden')],
  177. { maxTokens: 111 },
  178. )
  179. const inherited = await startInProcessRun(request(parent), {})
  180. await inherited.result
  181. expect(adapter.requests[0]?.maxTokens).toBe(111)
  182. expect(ctx.agents.get(inherited.id)?.options.maxTokens).toBe(111)
  183. await inherited.dispose()
  184. const overridden = await startInProcessRun({
  185. ...request(parent),
  186. agentOptions: { maxTokens: 222 },
  187. }, {})
  188. await overridden.result
  189. expect(adapter.requests[1]?.maxTokens).toBe(222)
  190. expect(ctx.agents.get(overridden.id)?.options.maxTokens).toBe(222)
  191. await overridden.dispose()
  192. })
  193. it('counts a RESUMED child by its persisted header depth, not the absent runtime depth', async () => {
  194. // Resume rebuilds runtime options, so the durable header must keep this
  195. // depth-1 child from delegating as though it were top-level.
  196. const { ctx } = await setup([textResponse('unused')])
  197. const resumed = (await ctx.agents.create({
  198. sessionId: SessionId('resumed-child'),
  199. meta: { parentSession: SessionId('root'), delegationDepth: 1 },
  200. agentOptions: { provider: 'mock', model: 'mock' },
  201. signal: new AbortController().signal,
  202. })).agent
  203. await expect(startInProcessRun({ ...request(resumed), maxDepth: 1 }, {}))
  204. .rejects.toMatchObject({ name: 'SubagentDepthError', attemptedDepth: 2, maxDepth: 1 })
  205. })
  206. it('lets runtime options deepen but never lower the persisted depth', async () => {
  207. const { ctx } = await setup([textResponse('unused')])
  208. const parent = (await ctx.agents.create({
  209. sessionId: SessionId('deep-parent'),
  210. meta: { delegationDepth: 2 },
  211. agentOptions: { provider: 'mock', model: 'mock', subagentDepth: 1 },
  212. signal: new AbortController().signal,
  213. })).agent
  214. // Persisted 2 vs runtime 1: the child is depth 3, so maxDepth 2 rejects.
  215. await expect(startInProcessRun({ ...request(parent), maxDepth: 2 }, {}))
  216. .rejects.toMatchObject({ name: 'SubagentDepthError', attemptedDepth: 3, maxDepth: 2 })
  217. })
  218. it('rejects invalid and exceeded depth before publication', async () => {
  219. const { parent } = await setup([])
  220. await expect(startInProcessRun({ ...request(parent), maxDepth: -1 }, {}))
  221. .rejects.toThrow('non-negative safe integer')
  222. await expect(startInProcessRun({ ...request(parent), maxDepth: 0 }, {}))
  223. .rejects.toMatchObject({ name: 'SubagentDepthError' })
  224. for (const value of [Number.NaN, 1.5, -1, -0, Number.MAX_SAFE_INTEGER + 1]) {
  225. const malformed = { options: { subagentDepth: value }, session: { header: {} } } as unknown as Agent
  226. await expect(startInProcessRun(request(malformed), {}))
  227. .rejects.toThrow('agent subagentDepth must be a non-negative safe integer')
  228. }
  229. const maxParent = { options: { subagentDepth: Number.MAX_SAFE_INTEGER }, session: { header: {} } } as unknown as Agent
  230. await expect(startInProcessRun(request(maxParent), {})).rejects.toBeInstanceOf(RangeError)
  231. })
  232. it('rejects an already-aborted request without publishing a child', async () => {
  233. const { ctx, parent } = await setup([])
  234. const beforeAgents = ctx.agents.list().length
  235. const beforeSessions = ctx.sessions.list().length
  236. const controller = new AbortController()
  237. controller.abort('too late')
  238. await expect(startInProcessRun(request(parent, controller.signal), {}))
  239. .rejects.toThrow('aborted before child publication')
  240. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  241. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  242. })
  243. it('stamps only the resolved depth when neither parent nor request declares a model route', async () => {
  244. // The one-shot analogue of the deleted resume coverage ("resumes without
  245. // inventing undeclared agent model options"): a bare parent with no request
  246. // agentOptions yields a child whose options carry ONLY the stamped depth —
  247. // no provider/model is fabricated, so the child's turn errors for want of a
  248. // route rather than silently adopting one.
  249. const { ctx } = await setup([])
  250. const parent = ctx.agentLoop.create(SessionId('routeless-parent'), {})
  251. const run = await startInProcessRun(request(parent), {})
  252. const child = ctx.agents.get(run.id)!
  253. expect(child.options).toEqual({ subagentDepth: 1 })
  254. await expect(run.result).resolves.toMatchObject({ stopReason: 'error' })
  255. await run.dispose()
  256. })
  257. it('uses the request signal after publication and dispose as cancellation paths', async () => {
  258. const { parent, adapter } = await setup(['hang', 'hang'])
  259. const controller = new AbortController()
  260. const signalled = await startInProcessRun(request(parent, controller.signal), {})
  261. await new Promise(resolve => setTimeout(resolve, 30))
  262. controller.abort('stop child')
  263. await expect(signalled.result).resolves.toMatchObject({ stopReason: 'aborted' })
  264. expect(adapter.requests[0]?.signal?.reason).toEqual({ kind: 'parent' })
  265. const child = parent.ctx.agents.get(signalled.id)
  266. const turnEnd = child?.session.events.findLast(event => event.type === 'turn/end')
  267. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  268. await signalled.dispose()
  269. const disposed = await startInProcessRun(request(parent), {})
  270. await new Promise(resolve => setTimeout(resolve, 30))
  271. await disposed.dispose()
  272. await expect(disposed.result).resolves.toMatchObject({ stopReason: 'aborted' })
  273. })
  274. it('cleans a failed unpublished setup before rejecting', async () => {
  275. const { ctx, parent } = await setup([])
  276. const beforeAgents = ctx.agents.list().length
  277. const beforeSessions = ctx.sessions.list().length
  278. await expect(startInProcessRun({
  279. ...request(parent),
  280. toolFilter: { deny: ['unknown-tool'] },
  281. }, {})).rejects.toThrow('unknown global tool')
  282. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  283. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  284. })
  285. it('treats abort after factory publication as a cancelled run with an id', async () => {
  286. const { ctx, parent } = await setup([])
  287. const controller = new AbortController()
  288. const beforeAgents = ctx.agents.list().length
  289. const beforeSessions = ctx.sessions.list().length
  290. const parentWithAbortAtHandoff = {
  291. options: parent.options,
  292. session: parent.session,
  293. ctx: {
  294. // The driver's synchronous inheritance capture probes both policy
  295. // services opportunistically; this stub composes neither.
  296. get: () => undefined,
  297. agents: {
  298. create: async (options: Parameters<typeof ctx.agents.create>[0]) => {
  299. const handle = await ctx.agents.create(options)
  300. // `create()` has detached its creation-only listener, but the
  301. // published run has not installed its live listener yet.
  302. controller.abort('handoff race')
  303. return handle
  304. },
  305. },
  306. },
  307. } as unknown as Agent
  308. const run = await startInProcessRun(request(parentWithAbortAtHandoff, controller.signal), {})
  309. expect(ctx.agents.get(run.id)).toBeDefined()
  310. await expect(run.result).resolves.toEqual({ output: [], stopReason: 'aborted' })
  311. await run.dispose()
  312. expect(ctx.agents.list()).toHaveLength(beforeAgents)
  313. expect(ctx.sessions.list()).toHaveLength(beforeSessions)
  314. })
  315. })