headless.spec.ts 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218
  1. /**
  2. * One-shot runner behavior over a scripted in-process API: idle-to-idle
  3. * aggregation (last text of the whole interval), exit-code mapping by the
  4. * final turn-end reason, stream-error and RPC-error paths, and the
  5. * launcher-owned `ctx.headlessIo` requirement.
  6. */
  7. import { describe, expect, it } from 'vitest'
  8. import { Context } from 'cordis'
  9. import type { Agent } from '@deepseek-ai/dsh-agent'
  10. import { apply, Config, type HeadlessIo } from '../src/index.ts'
  11. interface ScriptedEvent { type: string; seq?: number; time?: number; sessionId?: string; data: Record<string, unknown> }
  12. let nextSeq = 0
  13. /** Stamp the envelope fields the wire schema requires. */
  14. function stamped(event: ScriptedEvent): ScriptedEvent {
  15. nextSeq += 1
  16. return { seq: nextSeq, time: nextSeq, ...event }
  17. }
  18. interface RpcShapedRequest { rpcId: string }
  19. /** Build a fake apiProxy (echoing rpcIds like the real gateway) whose mux stream replays `events` for the created session. */
  20. function scriptedApi(events: ScriptedEvent[], options: { promptFails?: boolean } = {}): unknown {
  21. return {
  22. sessions: {
  23. create: (request: RpcShapedRequest) =>
  24. Promise.resolve({ rpcId: request.rpcId, result: { ok: true, value: { sessionId: 'S1' } } }),
  25. prompt: (request: RpcShapedRequest) => Promise.resolve(options.promptFails === true
  26. // A code from the closed wire union: the carrier schema rejects invented codes.
  27. ? { rpcId: request.rpcId, result: { ok: false, error: { code: 'agent-busy', message: 'agent is busy', details: { reason: 'test' } } } }
  28. : { rpcId: request.rpcId, result: { ok: true, value: { accepted: true } } }),
  29. },
  30. events: {
  31. mux: async function* () {
  32. for (const event of events) {
  33. if (event.type === 'stream/error') {
  34. yield { rpcId: 'e', payload: { type: 'stream/error', error: { code: 'cancelled', message: 'stream broke', details: {} } } }
  35. continue
  36. }
  37. const { sessionId = 'S1', ...rest } = event
  38. yield { rpcId: 'e', payload: { type: 'session/event', sessionId, event: stamped(rest) } }
  39. }
  40. },
  41. },
  42. }
  43. }
  44. /**
  45. * Mount the runner against a scripted API, emit the idle transition after the
  46. * scripted frames drain, and wait for its exit request.
  47. */
  48. async function run(events: ScriptedEvent[], options: { promptFails?: boolean } = {}): Promise<{ code: number; out: string; err: string }> {
  49. const ctx = new Context()
  50. let out = ''
  51. let err = ''
  52. const exited = new Promise<number>((resolve) => {
  53. const io: HeadlessIo = {
  54. stdout: { write: (chunk: string) => { out += chunk; return true } },
  55. stderr: { write: (chunk: string) => { err += chunk; return true } },
  56. exit: resolve,
  57. }
  58. ctx.provide('headlessIo', io)
  59. })
  60. ctx.provide('apiProxy', scriptedApi(events, options) as never)
  61. ctx.provide('httpServer', { port: 12345 } as never)
  62. apply(ctx, { task: 'do the thing' })
  63. // Quiescence is out of band: give the scripted stream a beat to drain, then
  64. // flip the agent idle exactly as the loop would. Foreign agents and
  65. // non-idle transitions must not settle the run.
  66. await new Promise(resolve => setTimeout(resolve, 10))
  67. ctx.emit('agent/status', { agent: { id: 'OTHER' } as Agent, status: 'idle' })
  68. ctx.emit('agent/status', { agent: { id: 'S1' } as Agent, status: 'running' })
  69. ctx.emit('agent/status', { agent: { id: 'S1' } as Agent, status: 'idle' })
  70. const code = await exited
  71. await ctx.fiber.dispose()
  72. return { code, out, err }
  73. }
  74. const startupTurn: ScriptedEvent = { type: 'turn/start', data: { turn: 0, trigger: { kind: 'startup' } } }
  75. const messageTurn: ScriptedEvent = { type: 'turn/start', data: { turn: 1, trigger: { kind: 'message' } } }
  76. const text = (turn: number, value: string): ScriptedEvent => ({
  77. type: 'assistant/message',
  78. data: { turn, message: { content: [{ type: 'text', text: value }] } },
  79. })
  80. const end = (turn: number, reason: string): ScriptedEvent => ({ type: 'turn/end', data: { turn, reason: { kind: reason } } })
  81. describe('headless runner', () => {
  82. it('aggregates to quiescence: last text wins across turns, final turn-end reason maps to exit 0', async () => {
  83. const { code, out, err } = await run([
  84. // Frames before the first turn/start are outside the task interval.
  85. { type: 'assistant/message', data: { turn: 0, message: { content: [{ type: 'text', text: 'pre-task noise' }] } } },
  86. startupTurn,
  87. // Off-session, non-text, and text-empty frames never affect the aggregate.
  88. { type: 'assistant/message', sessionId: 'OTHER', data: { turn: 1, message: { content: [{ type: 'text', text: 'other session' }] } } },
  89. { type: 'assistant/message', data: { turn: 1, message: { content: [{ type: 'tool_call', text: 'ignored' }] } } },
  90. text(0, 'draft'),
  91. end(0, 'completed'),
  92. messageTurn,
  93. text(1, 'final answer'),
  94. end(1, 'completed'),
  95. ])
  96. expect(code).toBe(0)
  97. expect(out).toBe('final answer\n')
  98. expect(err).toContain('observing at http://127.0.0.1:12345')
  99. })
  100. it('exits 1 when the final turn ends for any other reason', async () => {
  101. const { code } = await run([messageTurn, end(1, 'aborted')])
  102. expect(code).toBe(1)
  103. })
  104. it('exits 1 when no turn ever starts (idle without work)', async () => {
  105. const { code, out } = await run([])
  106. expect(code).toBe(1)
  107. expect(out).toBe('\n')
  108. })
  109. it('keeps the error outcome after a stream error ends the frame consumer early', async () => {
  110. const { code } = await run([messageTurn, { type: 'stream/error', data: {} }, end(1, 'completed')])
  111. // The consumer stopped at the stream error; the completed turn-end after
  112. // it is never observed, so the reason stays 'error'.
  113. expect(code).toBe(1)
  114. })
  115. it('prints an RPC business error and exits 1 without waiting for idle', async () => {
  116. const ctx = new Context()
  117. let err = ''
  118. const exited = new Promise<number>((resolve) => {
  119. ctx.provide('headlessIo', {
  120. stdout: { write: () => true },
  121. stderr: { write: (chunk: string) => { err += chunk; return true } },
  122. exit: resolve,
  123. } satisfies HeadlessIo)
  124. })
  125. ctx.provide('apiProxy', scriptedApi([messageTurn, end(1, 'completed')], { promptFails: true }) as never)
  126. ctx.provide('httpServer', { port: 1 } as never)
  127. apply(ctx, { task: 't' })
  128. expect(await exited).toBe(1)
  129. expect(err).toContain('agent-busy')
  130. await ctx.fiber.dispose()
  131. })
  132. it('reports the stream-failed diagnostic when the event channel dies, still settling at idle', async () => {
  133. const ctx = new Context()
  134. let err = ''
  135. const exited = new Promise<number>((resolve) => {
  136. ctx.provide('headlessIo', {
  137. stdout: { write: () => true },
  138. stderr: { write: (chunk: string) => { err += chunk; return true } },
  139. exit: resolve,
  140. } satisfies HeadlessIo)
  141. })
  142. ctx.provide('apiProxy', {
  143. sessions: {
  144. create: (request: RpcShapedRequest) =>
  145. Promise.resolve({ rpcId: request.rpcId, result: { ok: true, value: { sessionId: 'S1' } } }),
  146. prompt: (request: RpcShapedRequest) =>
  147. Promise.resolve({ rpcId: request.rpcId, result: { ok: true, value: { accepted: true } } }),
  148. },
  149. events: {
  150. // Synchronous throw: the SSE response never forms, so the client-side
  151. // iterable rejects — the runner's own catch path, not a carrier frame.
  152. mux: () => { throw new Error('channel exploded') },
  153. },
  154. } as never)
  155. ctx.provide('httpServer', { port: 1 } as never)
  156. apply(ctx, { task: 't' })
  157. await new Promise(resolve => setTimeout(resolve, 10))
  158. ctx.emit('agent/status', { agent: { id: 'S1' } as Agent, status: 'idle' })
  159. expect(await exited).toBe(1)
  160. expect(err).toContain('event stream failed')
  161. await ctx.fiber.dispose()
  162. })
  163. it('waits for Loader settlement and abandons the run when the tree died during it', async () => {
  164. const ctx = new Context()
  165. let err = ''
  166. let exited = false
  167. ctx.provide('headlessIo', {
  168. stdout: { write: () => true },
  169. stderr: { write: (chunk: string) => { err += chunk; return true } },
  170. exit: () => { exited = true },
  171. } satisfies HeadlessIo)
  172. ctx.provide('apiProxy', scriptedApi([]) as never)
  173. // The webserver is provided by a child fiber whose disposal (early
  174. // SIGTERM during the boot window) removes the service; settlement
  175. // resolves only afterwards, and the runner must abandon rather than
  176. // crash on the torn-down port read.
  177. const webserverFiber = ctx.plugin((childCtx: Context) => {
  178. childCtx.provide('httpServer', { port: 1 } as never)
  179. })
  180. await webserverFiber
  181. let release: () => void
  182. const settlement = new Promise<void>((resolve) => { release = resolve })
  183. ctx.provide('loader', { await: () => settlement } as never)
  184. apply(ctx, { task: 't' })
  185. await webserverFiber.dispose()
  186. release!()
  187. await new Promise(resolve => setTimeout(resolve, 10))
  188. expect(err).toBe('')
  189. expect(exited).toBe(false)
  190. await ctx.fiber.dispose()
  191. })
  192. it('fails loud without the launcher-owned headlessIo seam', () => {
  193. const ctx = new Context()
  194. ctx.provide('apiProxy', scriptedApi([]) as never)
  195. ctx.provide('httpServer', { port: 1 } as never)
  196. expect(() => { apply(ctx, { task: 't' }) }).toThrow('must provide ctx.headlessIo')
  197. })
  198. it('validates config: the task is required', () => {
  199. expect(() => new Config({ } as never)).toThrow()
  200. expect(new Config({ task: 'x' })).toEqual({ task: 'x' })
  201. })
  202. })