headless.spec.ts 9.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252
  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. run(): Promise<{ code: number; out: string; err: string; order: string[] }>
  49. }> {
  50. const ctx = new Context()
  51. await ctx.plugin(SessionStore)
  52. await ctx.plugin(AgentRegistry)
  53. await ctx.plugin(AgentDefaultModelConfig, { provider: 'test-provider', model: 'test-model' })
  54. ctx.agents.setFactory({
  55. async createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> {
  56. const session = ctx.sessions.create(options.sessionId, {
  57. ...options.meta === undefined ? {} : { meta: options.meta },
  58. })
  59. let idle = Promise.resolve()
  60. const agent = {} as Agent
  61. const agentCtx = ownerCtx.extend({ agent })
  62. Object.assign(agent, {
  63. id: session.id,
  64. options: options.agentOptions ?? {},
  65. session,
  66. inbox: new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} }),
  67. status: 'idle',
  68. ctx: agentCtx,
  69. cancel: () => {},
  70. runMaintenance: () => Promise.reject(new Error('not used')),
  71. send: () => {},
  72. followup: (message: UserMessage) => {
  73. agent.inbox.append('next-turn', message)
  74. idle = Promise.resolve().then(() => script.afterPrompt(session, message))
  75. },
  76. steer: () => {},
  77. inject: () => {},
  78. whenIdle: () => idle,
  79. } satisfies Partial<Agent>)
  80. await options.setup?.(agentCtx)
  81. script.before?.(session)
  82. ctx.agents.register(agent)
  83. return { agent, dispose: () => Promise.resolve() }
  84. },
  85. resume: () => Promise.reject(new Error('not used')),
  86. })
  87. return {
  88. ctx,
  89. run: async () => {
  90. let out = ''
  91. let err = ''
  92. const order: string[] = []
  93. ctx.on('session/flush', () => { order.push('flush') })
  94. internals.stdout = { write: (chunk: string) => { out += chunk; return true } }
  95. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  96. const exited = new Promise<number>((resolve) => {
  97. ctx.provide('appExit', (code: number) => { order.push('exit'); resolve(code) })
  98. })
  99. apply(ctx, { task: 'do the thing' })
  100. return { code: await exited, out, err, order }
  101. },
  102. }
  103. }
  104. describe('headless runner', () => {
  105. it('aggregates the final text across the complete idle-to-idle interval and flushes before exit', async () => {
  106. const test = await bench({
  107. before(session) {
  108. const setupMessage = {
  109. role: 'user', content: [{ type: 'text', text: 'setup' }], source: { kind: 'user' }, id: 'setup',
  110. } as UserMessage
  111. appendTurn(session, 0, setupMessage, 'pre-task noise', true)
  112. },
  113. async afterPrompt(session, message) {
  114. await Promise.resolve()
  115. appendTurn(session, 1, message, '', true)
  116. appendTurn(session, 2, message, 'final answer', true)
  117. },
  118. })
  119. const result = await test.run()
  120. expect(result).toEqual({
  121. code: 0,
  122. out: 'final answer\n',
  123. err: '',
  124. order: ['flush', 'exit'],
  125. })
  126. await test.ctx.fiber.dispose()
  127. })
  128. it('waits for asynchronously appended events instead of racing Agent idleness', async () => {
  129. const test = await bench({
  130. afterPrompt: async (session, message) => {
  131. await new Promise(resolve => setTimeout(resolve, 5))
  132. appendTurn(session, 1, message, 'race-free answer', true)
  133. },
  134. })
  135. expect(await test.run()).toMatchObject({ code: 0, out: 'race-free answer\n', err: '' })
  136. await test.ctx.fiber.dispose()
  137. })
  138. it('exits 1 when the final turn does not complete', async () => {
  139. const test = await bench({
  140. afterPrompt(session, message) { appendTurn(session, 1, message, undefined, false) },
  141. })
  142. expect(await test.run()).toMatchObject({ code: 1, out: '\n', err: '' })
  143. await test.ctx.fiber.dispose()
  144. })
  145. it('prints the durable model failure when the final turn ends in error', async () => {
  146. const test = await bench({
  147. afterPrompt(session, message) {
  148. session.append('turn/start', { turn: 1 })
  149. session.append('step/start', { turn: 1, step: 1 })
  150. session.append('user/message', message, { surfaceOp: 'append' })
  151. session.append('step/end', { turn: 1, step: 1 })
  152. session.append('turn/end', {
  153. turn: 1,
  154. reason: { kind: 'error', error: { code: 'SERVER', message: 'provider unavailable' } },
  155. })
  156. },
  157. })
  158. expect(await test.run()).toMatchObject({
  159. code: 1,
  160. out: '\n',
  161. err: 'dsh: SERVER: provider unavailable\n',
  162. })
  163. await test.ctx.fiber.dispose()
  164. })
  165. it('exits 1 when the owned interval contains no turn', async () => {
  166. const test = await bench({ afterPrompt: () => {} })
  167. expect(await test.run()).toMatchObject({ code: 1, out: '\n', err: '' })
  168. await test.ctx.fiber.dispose()
  169. })
  170. it('reports a direct Agent creation failure', async () => {
  171. const ctx = new Context()
  172. let err = ''
  173. internals.stdout = { write: () => true }
  174. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  175. const exited = new Promise<number>((resolve) => {
  176. ctx.provide('appExit', resolve)
  177. })
  178. ctx.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  179. ctx.provide('sessions', { flush: () => Promise.resolve(true) } as never)
  180. ctx.provide('agents', { create: () => Promise.reject(new Error('factory exploded')) } as never)
  181. apply(ctx, { task: 't' })
  182. expect(await exited).toBe(1)
  183. expect(err).toBe('dsh: factory exploded\n')
  184. await ctx.fiber.dispose()
  185. })
  186. it('stringifies a non-Error Agent creation failure', async () => {
  187. const ctx = new Context()
  188. let err = ''
  189. internals.stdout = { write: () => true }
  190. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  191. const exited = new Promise<number>((resolve) => {
  192. ctx.provide('appExit', resolve)
  193. })
  194. ctx.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  195. ctx.provide('sessions', { flush: () => Promise.resolve(true) } as never)
  196. const rejected = {
  197. then(_resolve: (value: never) => void, reject: (reason: unknown) => void): void {
  198. reject('factory exploded')
  199. },
  200. }
  201. ctx.provide('agents', { create: () => rejected } as never)
  202. apply(ctx, { task: 't' })
  203. expect(await exited).toBe(1)
  204. expect(err).toBe('dsh: factory exploded\n')
  205. await ctx.fiber.dispose()
  206. })
  207. it('abandons a run when the tree is disposed during Loader settlement', async () => {
  208. const ctx = new Context()
  209. let exited = false
  210. internals.stdout = { write: () => true }
  211. internals.stderr = { write: () => true }
  212. ctx.provide('appExit', () => { exited = true })
  213. const services = ctx.plugin((child: Context) => {
  214. child.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  215. child.provide('sessions', {} as never)
  216. child.provide('agents', {} as never)
  217. })
  218. await services
  219. let release: () => void
  220. const settlement = new Promise<void>((resolve) => { release = resolve })
  221. ctx.provide('loader', { await: () => settlement } as never)
  222. apply(ctx, { task: 't' })
  223. await services.dispose()
  224. release!()
  225. await new Promise(resolve => setTimeout(resolve, 10))
  226. expect(exited).toBe(false)
  227. await ctx.fiber.dispose()
  228. })
  229. it('fails loud without the launcher-provided exit request', () => {
  230. const ctx = new Context()
  231. expect(() => { apply(ctx, { task: 't' }) }).toThrow('must provide ctx.appExit')
  232. })
  233. it('validates config: the task is required', () => {
  234. expect(() => new Config({} as never)).toThrow()
  235. expect(new Config({ task: 'x' })).toEqual({ task: 'x' })
  236. })
  237. })