headless.spec.ts 14 KB

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