headless.spec.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449
  1. /** Direct one-shot Agent driving, durable aggregation, flushing, and exit mapping. */
  2. import { afterEach, describe, expect, it } from 'vitest'
  3. import { Context } from '@deepseek-ai/cordis'
  4. import AgentRegistry from '@deepseek-ai/dsh-agent'
  5. import type { Agent, AgentHandle, AssistantStreamFrame, CreateAgentOptions } from '@deepseek-ai/dsh-agent'
  6. import AgentDefaultModelConfig from '@deepseek-ai/dsh-agent-default-model'
  7. import { LlmAttemptId, createAssistantMessage, type StreamChunk } from '@deepseek-ai/dsh-llm'
  8. import SessionStore from '@deepseek-ai/dsh-session'
  9. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  10. import type { Session, UserMessage } from '@deepseek-ai/dsh-session'
  11. import { createInboxStub } from '@deepseek-ai/dsh-agent-loop-testkit'
  12. import { apply, Config, internals } from '../src/index.ts'
  13. const originalInternals = { ...internals }
  14. afterEach(() => { Object.assign(internals, originalInternals) })
  15. interface Script {
  16. before?(session: Session): void
  17. afterPrompt(session: Session, message: UserMessage, agent: Agent): Promise<void> | void
  18. }
  19. const frameStates = new WeakMap<Agent, { attemptId: ReturnType<typeof LlmAttemptId>; revision: number; index: number }>()
  20. function startFrames(agent: Agent, turn = 1, step = 1): void {
  21. const state = { attemptId: LlmAttemptId(`${agent.id}:test`), revision: 1, index: 0 }
  22. frameStates.set(agent, state)
  23. agent.ctx.emit('agent/assistant-stream', {
  24. agent,
  25. frame: {
  26. type: 'start', attemptId: state.attemptId, revision: state.revision, turn, step,
  27. },
  28. })
  29. }
  30. function emitChunk(agent: Agent, chunk: StreamChunk): void {
  31. const state = frameStates.get(agent)
  32. if (state === undefined) throw new Error('test Assistant frames have not started')
  33. const frame: AssistantStreamFrame = {
  34. type: 'chunk', attemptId: state.attemptId, revision: ++state.revision,
  35. index: state.index++, time: Date.now(), chunk,
  36. }
  37. agent.ctx.emit('agent/assistant-stream', { agent, frame })
  38. }
  39. function appendTurn(
  40. session: Session,
  41. turn: number,
  42. message: UserMessage,
  43. text: string | undefined,
  44. completed: boolean,
  45. ): void {
  46. session.append('turn/start', { turn })
  47. session.append('step/start', { turn, step: 1 })
  48. session.append('user/message', message, { surfaceOp: 'append' })
  49. if (text !== undefined) {
  50. session.append('assistant/message', {
  51. stream: [],
  52. turn,
  53. step: 1,
  54. message: createAssistantMessage({
  55. content: [{ type: 'text', text }],
  56. source: { provider: 'test-provider', model: 'test-model' },
  57. }),
  58. }, { surfaceOp: 'append' })
  59. }
  60. session.append('step/end', { turn, step: 1 })
  61. session.append('turn/end', {
  62. turn,
  63. reason: completed
  64. ? { kind: 'completed' }
  65. : { kind: 'aborted', reason: { kind: 'user' } },
  66. })
  67. }
  68. /** Mount the real registries around a small scripted Agent factory. */
  69. async function bench(script: Script): Promise<{
  70. ctx: Context
  71. output(): { out: string; err: string; order: string[] }
  72. run(): Promise<{ code: number; out: string; err: string; order: string[] }>
  73. }> {
  74. const ctx = new Context()
  75. let out = ''
  76. let err = ''
  77. const order: string[] = []
  78. await ctx.plugin(SessionStore)
  79. await ctx.plugin(SessionProjectionRegistry)
  80. await ctx.plugin(AgentRegistry)
  81. await ctx.plugin(AgentDefaultModelConfig, { provider: 'test-provider', model: 'test-model' })
  82. ctx.agents.setFactory({
  83. async createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> {
  84. const session = ctx.sessions.create(options.sessionId, {
  85. ...options.meta === undefined ? {} : { meta: options.meta },
  86. })
  87. const inbox = createInboxStub()
  88. let idle = Promise.resolve()
  89. const agent: Agent = {
  90. id: session.id,
  91. options: options.agentOptions ?? {},
  92. session,
  93. inbox,
  94. status: 'idle',
  95. ctx: ownerCtx,
  96. cancel: () => {},
  97. runMaintenance: () => Promise.reject(new Error('not used')),
  98. send: () => {},
  99. followup: (message: UserMessage) => {
  100. agent.inbox.append('next-turn', message)
  101. idle = Promise.resolve().then(() => script.afterPrompt(session, message, agent))
  102. },
  103. steer: () => {},
  104. inject: () => {},
  105. whenIdle: () => idle,
  106. }
  107. await options.setup?.(ownerCtx, agent)
  108. script.before?.(session)
  109. ctx.agents.register(agent)
  110. return { agent, dispose: () => Promise.resolve() }
  111. },
  112. resume: () => Promise.reject(new Error('not used')),
  113. })
  114. return {
  115. ctx,
  116. output: () => ({ out, err, order: [...order] }),
  117. run: async () => {
  118. ctx.on('session/flush', () => { order.push('flush') })
  119. internals.stdout = { write: (chunk: string) => { out += chunk; return true } }
  120. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  121. const exited = new Promise<number>((resolve) => {
  122. ctx.provide('appExit', (code: number) => { order.push('exit'); resolve(code) })
  123. })
  124. apply(ctx, { task: 'do the thing' })
  125. return { code: await exited, out, err, order }
  126. },
  127. }
  128. }
  129. describe('headless runner', () => {
  130. it('aggregates the final text across the complete idle-to-idle interval and flushes before exit', async () => {
  131. const test = await bench({
  132. before(session) {
  133. const setupMessage = {
  134. role: 'user', content: [{ type: 'text', text: 'setup' }], source: { kind: 'user' }, id: 'setup',
  135. } as UserMessage
  136. appendTurn(session, 0, setupMessage, 'pre-task noise', true)
  137. },
  138. async afterPrompt(session, message) {
  139. await Promise.resolve()
  140. appendTurn(session, 1, message, '', true)
  141. appendTurn(session, 2, message, 'final answer', true)
  142. },
  143. })
  144. const result = await test.run()
  145. expect(result).toEqual({
  146. code: 0,
  147. out: 'final answer\n',
  148. err: '',
  149. order: ['flush', 'exit'],
  150. })
  151. await test.ctx.fiber.dispose()
  152. })
  153. it('ignores durable inbox events before the first owned turn', async () => {
  154. const test = await bench({
  155. afterPrompt(session, message) {
  156. session.append('agent/inbox/spliced', {
  157. target: 'next-turn',
  158. start: 0,
  159. inserted: [message],
  160. })
  161. appendTurn(session, 1, message, 'answer after inbox activity', true)
  162. },
  163. })
  164. expect(await test.run()).toMatchObject({
  165. code: 0,
  166. out: 'answer after inbox activity\n',
  167. err: '',
  168. })
  169. await test.ctx.fiber.dispose()
  170. })
  171. it('waits for asynchronously appended events instead of racing Agent idleness', async () => {
  172. const test = await bench({
  173. afterPrompt: async (session, message) => {
  174. await new Promise(resolve => setTimeout(resolve, 5))
  175. appendTurn(session, 1, message, 'race-free answer', true)
  176. },
  177. })
  178. expect(await test.run()).toMatchObject({ code: 0, out: 'race-free answer\n', err: '' })
  179. await test.ctx.fiber.dispose()
  180. })
  181. it('streams reasoning before the Agent becomes idle and terminates its stderr line', async () => {
  182. const reasoningAppended = Promise.withResolvers<undefined>()
  183. const release = Promise.withResolvers<undefined>()
  184. const test = await bench({
  185. async afterPrompt(session, message, agent) {
  186. session.append('turn/start', { turn: 1 })
  187. session.append('step/start', { turn: 1, step: 1 })
  188. session.append('user/message', message, { surfaceOp: 'append' })
  189. startFrames(agent)
  190. emitChunk(agent, { type: 'block-start', index: 0, blockType: 'reasoning' })
  191. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: '' })
  192. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'checking the workspace' })
  193. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: ' safely\n' })
  194. emitChunk(agent, { type: 'block-end', index: 0, block: { type: 'reasoning', text: 'checking the workspace safely\n' } })
  195. emitChunk(agent, { type: 'usage', usage: { inputTokens: 1, outputTokens: 2, reasoningTokens: 2 } })
  196. emitChunk(agent, { type: 'block-start', index: 1, blockType: 'reasoning' })
  197. emitChunk(agent, { type: 'reasoning-delta', index: 1, text: 'second pass\n' })
  198. reasoningAppended.resolve(undefined)
  199. await release.promise
  200. emitChunk(agent, { type: 'block-start', index: 2, blockType: 'text' })
  201. emitChunk(agent, { type: 'text-delta', index: 2, text: 'done' })
  202. emitChunk(agent, { type: 'block-end', index: 2, block: { type: 'text', text: 'done' } })
  203. session.append('assistant/message', {
  204. stream: [],
  205. turn: 1,
  206. step: 1,
  207. message: createAssistantMessage({
  208. content: [{ type: 'text', text: 'done' }],
  209. source: { provider: 'test-provider', model: 'test-model' },
  210. }),
  211. }, { surfaceOp: 'append' })
  212. session.append('step/end', { turn: 1, step: 1 })
  213. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  214. },
  215. })
  216. const running = test.run()
  217. await reasoningAppended.promise
  218. const other = test.ctx.sessions.create()
  219. other.append('turn/start', { turn: 1 })
  220. other.append('step/start', { turn: 1, step: 1 })
  221. test.ctx.emit('agent/assistant-stream', {
  222. agent: { session: other } as Agent,
  223. frame: {
  224. type: 'chunk', attemptId: LlmAttemptId('other'), revision: 1,
  225. index: 0, time: Date.now(), chunk: { type: 'reasoning-delta', index: 0, text: 'other session' },
  226. },
  227. })
  228. const streamed = test.output()
  229. release.resolve(undefined)
  230. const result = await running
  231. expect(streamed).toEqual({
  232. out: '',
  233. err: 'dsh: reasoning:\nchecking the workspace safely\nsecond pass\n',
  234. order: [],
  235. })
  236. expect(result).toEqual({
  237. code: 0,
  238. out: 'done\n',
  239. err: 'dsh: reasoning:\nchecking the workspace safely\nsecond pass\n',
  240. order: ['flush', 'exit'],
  241. })
  242. await test.ctx.fiber.dispose()
  243. })
  244. it('closes an unterminated reasoning line as soon as the attempt ends', async () => {
  245. const reasoningAppended = Promise.withResolvers<undefined>()
  246. const releaseEnd = Promise.withResolvers<undefined>()
  247. const ended = Promise.withResolvers<undefined>()
  248. const finish = Promise.withResolvers<undefined>()
  249. const test = await bench({
  250. async afterPrompt(session, message, agent) {
  251. session.append('turn/start', { turn: 1 })
  252. session.append('step/start', { turn: 1, step: 1 })
  253. session.append('user/message', message, { surfaceOp: 'append' })
  254. startFrames(agent)
  255. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'unfinished reasoning' })
  256. reasoningAppended.resolve(undefined)
  257. await releaseEnd.promise
  258. const state = frameStates.get(agent)
  259. if (state === undefined) throw new Error('test Assistant frames have not started')
  260. agent.ctx.emit('agent/assistant-stream', {
  261. agent,
  262. frame: {
  263. type: 'end', attemptId: state.attemptId, revision: ++state.revision,
  264. index: state.index, outcome: { kind: 'abandoned' },
  265. },
  266. })
  267. ended.resolve(undefined)
  268. await finish.promise
  269. session.append('step/end', { turn: 1, step: 1 })
  270. session.append('turn/end', {
  271. turn: 1, reason: { kind: 'aborted', reason: { kind: 'user' } },
  272. })
  273. },
  274. })
  275. const running = test.run()
  276. await reasoningAppended.promise
  277. expect(test.output().err).toBe('dsh: reasoning:\nunfinished reasoning')
  278. releaseEnd.resolve(undefined)
  279. await ended.promise
  280. expect(test.output().err).toBe('dsh: reasoning:\nunfinished reasoning\n')
  281. finish.resolve(undefined)
  282. await expect(running).resolves.toMatchObject({ code: 1 })
  283. await test.ctx.fiber.dispose()
  284. })
  285. it('exits 1 when the final turn does not complete', async () => {
  286. const test = await bench({
  287. afterPrompt(session, message) { appendTurn(session, 1, message, undefined, false) },
  288. })
  289. expect(await test.run()).toMatchObject({ code: 1, out: '\n', err: '' })
  290. await test.ctx.fiber.dispose()
  291. })
  292. it('prints the durable model failure when the final turn ends in error', async () => {
  293. const test = await bench({
  294. afterPrompt(session, message) {
  295. session.append('turn/start', { turn: 1 })
  296. session.append('step/start', { turn: 1, step: 1 })
  297. session.append('user/message', message, { surfaceOp: 'append' })
  298. session.append('step/end', { turn: 1, step: 1 })
  299. session.append('turn/end', {
  300. turn: 1,
  301. reason: { kind: 'error', error: { code: 'SERVER', message: 'provider unavailable' } },
  302. })
  303. },
  304. })
  305. expect(await test.run()).toMatchObject({
  306. code: 1,
  307. out: '\n',
  308. err: 'dsh: SERVER: provider unavailable\n',
  309. })
  310. await test.ctx.fiber.dispose()
  311. })
  312. it('separates an unterminated reasoning prefix from the terminal model failure', async () => {
  313. const test = await bench({
  314. afterPrompt(session, message, agent) {
  315. session.append('turn/start', { turn: 1 })
  316. session.append('step/start', { turn: 1, step: 1 })
  317. session.append('user/message', message, { surfaceOp: 'append' })
  318. startFrames(agent)
  319. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'trying recovery' })
  320. session.append('step/end', { turn: 1, step: 1 })
  321. session.append('turn/end', {
  322. turn: 1,
  323. reason: { kind: 'error', error: { code: 'SERVER', message: 'provider unavailable' } },
  324. })
  325. },
  326. })
  327. expect(await test.run()).toMatchObject({
  328. code: 1,
  329. out: '\n',
  330. err: 'dsh: reasoning:\ntrying recovery\ndsh: SERVER: provider unavailable\n',
  331. })
  332. await test.ctx.fiber.dispose()
  333. })
  334. it('exits 1 when the owned interval contains no turn', async () => {
  335. const test = await bench({ afterPrompt: () => {} })
  336. expect(await test.run()).toMatchObject({ code: 1, out: '\n', err: '' })
  337. await test.ctx.fiber.dispose()
  338. })
  339. it('fails when an event below the captured Session length cannot be read', async () => {
  340. let capturedLength = 0
  341. const test = await bench({
  342. afterPrompt(session, message) {
  343. appendTurn(session, 1, message, 'unreachable', true)
  344. capturedLength = session.seq
  345. Object.defineProperty(session, 'eventAt', { value: () => undefined })
  346. },
  347. })
  348. const result = await test.run()
  349. expect(capturedLength).toBeGreaterThan(0)
  350. expect(result).toMatchObject({
  351. code: 1,
  352. out: '',
  353. err: `dsh: headless summary cannot read seq 0 below captured length ${String(capturedLength)}\n`,
  354. })
  355. await test.ctx.fiber.dispose()
  356. })
  357. it('reports a direct Agent creation failure', async () => {
  358. const ctx = new Context()
  359. let err = ''
  360. internals.stdout = { write: () => true }
  361. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  362. const exited = new Promise<number>((resolve) => {
  363. ctx.provide('appExit', resolve)
  364. })
  365. ctx.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  366. ctx.provide('sessions', { flush: () => Promise.resolve(true) } as never)
  367. ctx.provide('agents', { create: () => Promise.reject(new Error('factory exploded')) } as never)
  368. apply(ctx, { task: 't' })
  369. expect(await exited).toBe(1)
  370. expect(err).toBe('dsh: factory exploded\n')
  371. await ctx.fiber.dispose()
  372. })
  373. it('stringifies a non-Error Agent creation failure', async () => {
  374. const ctx = new Context()
  375. let err = ''
  376. internals.stdout = { write: () => true }
  377. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  378. const exited = new Promise<number>((resolve) => {
  379. ctx.provide('appExit', resolve)
  380. })
  381. ctx.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  382. ctx.provide('sessions', { flush: () => Promise.resolve(true) } as never)
  383. const rejected = {
  384. then(_resolve: (value: never) => void, reject: (reason: unknown) => void): void {
  385. reject('factory exploded')
  386. },
  387. }
  388. ctx.provide('agents', { create: () => rejected } as never)
  389. apply(ctx, { task: 't' })
  390. expect(await exited).toBe(1)
  391. expect(err).toBe('dsh: factory exploded\n')
  392. await ctx.fiber.dispose()
  393. })
  394. it('abandons a run when the tree is disposed during Loader settlement', async () => {
  395. const ctx = new Context()
  396. let exited = false
  397. internals.stdout = { write: () => true }
  398. internals.stderr = { write: () => true }
  399. ctx.provide('appExit', () => { exited = true })
  400. const services = ctx.plugin((child: Context) => {
  401. child.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  402. child.provide('sessions', {} as never)
  403. child.provide('agents', {} as never)
  404. })
  405. await services
  406. let release: () => void
  407. const settlement = new Promise<void>((resolve) => { release = resolve })
  408. ctx.provide('loader', { await: () => settlement } as never)
  409. apply(ctx, { task: 't' })
  410. await services.dispose()
  411. release!()
  412. await new Promise(resolve => setTimeout(resolve, 10))
  413. expect(exited).toBe(false)
  414. await ctx.fiber.dispose()
  415. })
  416. it('fails loud without the launcher-provided exit request', () => {
  417. const ctx = new Context()
  418. expect(() => { apply(ctx, { task: 't' }) }).toThrow('must provide ctx.appExit')
  419. })
  420. it('validates config: the task is required', () => {
  421. expect(() => new Config({} as never)).toThrow()
  422. expect(new Config({ task: 'x' })).toEqual({ task: 'x' })
  423. })
  424. })