headless.spec.ts 41 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052
  1. /** Direct one-shot Agent driving, exact Session adoption, machine-readable projection, and exit mapping. */
  2. import { Readable } from 'node:stream'
  3. import { afterEach, describe, expect, it } from 'vitest'
  4. import { Context } from '@deepseek-ai/cordis'
  5. import { brandString } from '@deepseek-ai/dsh-brand'
  6. import AgentRegistry from '@deepseek-ai/dsh-agent'
  7. import type {
  8. Agent,
  9. AgentHandle,
  10. AssistantStreamFrame,
  11. CreateAgentOptions,
  12. ResumeAgentOptions,
  13. } from '@deepseek-ai/dsh-agent'
  14. import AgentDefaultModelConfig from '@deepseek-ai/dsh-agent-default-model'
  15. import { LlmAttemptId, ToolCallId, createAssistantMessage, createToolResultMessage, type StreamChunk } from '@deepseek-ai/dsh-llm'
  16. import SessionStore from '@deepseek-ai/dsh-session'
  17. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  18. import type { Session, SessionId, UserMessage } from '@deepseek-ai/dsh-session'
  19. import { SessionQueryError } from '@deepseek-ai/dsh-session-query'
  20. import { createInboxStub } from '@deepseek-ai/dsh-agent-loop-testkit'
  21. import { apply, Config } from '../src/index.ts'
  22. import { internals } from '../src/runner-internals.ts'
  23. const originalInternals = { ...internals }
  24. afterEach(() => { Object.assign(internals, originalInternals) })
  25. interface Script {
  26. before?(session: Session): void
  27. afterPrompt(session: Session, message: UserMessage, agent: Agent): Promise<void> | void
  28. }
  29. /** Observation stub returned by the `--session-id` query path. */
  30. interface ObservationStub {
  31. header: { cwd?: string; origin?: string; parentSession?: string; agentPreset?: string }
  32. events: readonly { type: string; data: unknown }[]
  33. [Symbol.dispose](): void
  34. }
  35. /** Runner invocation options layered over the scripted Agent factory. */
  36. interface BenchOptions {
  37. /** Provider-resolved cwd, which can differ from the harness process directory. */
  38. filesystemCwd?: string
  39. task?: string
  40. useStdin?: boolean
  41. readStdin?: () => Promise<string>
  42. sessionId?: string
  43. json?: boolean
  44. observe?: () => Promise<ObservationStub>
  45. /** Leave the query service unmounted to exercise the fail-loud path. */
  46. omitSessionQuery?: boolean
  47. /** Leave the persistence service unmounted to exercise the fail-loud path. */
  48. omitPersistence?: boolean
  49. /** Register a live Agent under `sessionId` before the runner starts. */
  50. prelive?: boolean
  51. /** Header facts for that pre-registered live Agent. */
  52. preliveMeta?: { cwd?: string; origin?: 'subagent'; agentPreset?: string }
  53. /** Run when the runner awaits idle, e.g. to append to the attached log. */
  54. onWhenIdle?: (agent: Agent) => void
  55. }
  56. const frameStates = new WeakMap<Agent, { attemptId: ReturnType<typeof LlmAttemptId>; revision: number; index: number }>()
  57. function startFrames(agent: Agent, turn = 1, step = 1): void {
  58. const state = { attemptId: LlmAttemptId(`${agent.id}:test`), revision: 1, index: 0 }
  59. frameStates.set(agent, state)
  60. agent.ctx.emit('agent/assistant-stream', {
  61. agent,
  62. frame: {
  63. type: 'start', attemptId: state.attemptId, revision: state.revision, turn, step,
  64. },
  65. })
  66. }
  67. function emitChunk(agent: Agent, chunk: StreamChunk): void {
  68. const state = frameStates.get(agent)
  69. if (state === undefined) throw new Error('test Assistant frames have not started')
  70. const frame: AssistantStreamFrame = {
  71. type: 'chunk', attemptId: state.attemptId, revision: ++state.revision,
  72. index: state.index++, time: Date.now(), chunk,
  73. }
  74. agent.ctx.emit('agent/assistant-stream', { agent, frame })
  75. }
  76. function appendTurn(
  77. session: Session,
  78. turn: number,
  79. message: UserMessage,
  80. text: string | undefined,
  81. completed: boolean,
  82. ): void {
  83. session.append('turn/start', { turn })
  84. session.append('step/start', { turn, step: 1 })
  85. session.append('user/message', message, { surfaceOp: 'append' })
  86. if (text !== undefined) {
  87. session.append('assistant/message', {
  88. stream: [],
  89. turn,
  90. step: 1,
  91. message: createAssistantMessage({
  92. content: [{ type: 'text', text }],
  93. source: { provider: 'test-provider', model: 'test-model' },
  94. }),
  95. }, { surfaceOp: 'append' })
  96. }
  97. session.append('step/end', { turn, step: 1 })
  98. session.append('turn/end', {
  99. turn,
  100. reason: completed
  101. ? { kind: 'completed' }
  102. : { kind: 'aborted', reason: { kind: 'user' } },
  103. })
  104. }
  105. /** Append the preset-selection event owned by dsh-agent-presets. */
  106. function selectPreset(session: Session, agentPreset: string): void {
  107. const target = session as unknown as { append(type: string, data: unknown): void }
  108. target.append('agent-preset/selected', { agentPreset })
  109. }
  110. /** Mount the real registries around a small scripted Agent factory. */
  111. async function bench(script: Script, options: BenchOptions = {}): Promise<{
  112. ctx: Context
  113. output(): { out: string; err: string; order: string[] }
  114. run(): Promise<{ code: number; out: string; err: string; order: string[] }>
  115. }> {
  116. const ctx = new Context()
  117. if (options.filesystemCwd !== undefined) {
  118. const cwd = options.filesystemCwd
  119. ctx.provide('fs', {
  120. resolve: async () => ({ targetKey: cwd, displayPath: cwd }),
  121. processPath: () => cwd,
  122. } as never)
  123. }
  124. let out = ''
  125. let err = ''
  126. const order: string[] = []
  127. const mount = async (
  128. ownerCtx: Context,
  129. session: Session,
  130. createOptions: CreateAgentOptions | ResumeAgentOptions,
  131. ): Promise<Agent> => {
  132. const inbox = createInboxStub()
  133. let idle = Promise.resolve()
  134. const agent: Agent = {
  135. id: session.id,
  136. options: createOptions.agentOptions ?? {},
  137. session,
  138. inbox,
  139. status: 'idle',
  140. ctx: ownerCtx,
  141. cancel: () => {},
  142. runMaintenance: () => Promise.reject(new Error('not used')),
  143. send: () => {},
  144. followup: (message: UserMessage) => {
  145. agent.inbox.append('next-turn', message)
  146. idle = Promise.resolve().then(() => script.afterPrompt(session, message, agent))
  147. },
  148. steer: () => {},
  149. inject: () => {},
  150. whenIdle: () => {
  151. options.onWhenIdle?.(agent)
  152. return idle
  153. },
  154. }
  155. await createOptions.setup?.(ownerCtx, agent)
  156. await ctx.agents.register(agent)
  157. return agent
  158. }
  159. await ctx.plugin(SessionStore)
  160. await ctx.plugin(SessionProjectionRegistry)
  161. await ctx.plugin(AgentRegistry)
  162. await ctx.plugin(AgentDefaultModelConfig, { provider: 'test-provider', model: 'test-model' })
  163. ctx.agents.setFactory({
  164. async createAgent(ownerCtx: Context, createOptions: CreateAgentOptions): Promise<AgentHandle> {
  165. const session = ctx.sessions.create(createOptions.sessionId, {
  166. ...createOptions.meta === undefined ? {} : { meta: createOptions.meta },
  167. })
  168. script.before?.(session)
  169. const agent = await mount(ownerCtx, session, createOptions)
  170. return { agent, dispose: () => Promise.resolve() }
  171. },
  172. async resume(ownerCtx: Context, resumeOptions: ResumeAgentOptions): Promise<AgentHandle> {
  173. const session = ctx.sessions.get(resumeOptions.resumeSessionId)
  174. if (session === undefined) throw new Error(`no attached Session ${resumeOptions.resumeSessionId}`)
  175. const agent = await mount(ownerCtx, session, resumeOptions)
  176. return { agent, dispose: () => Promise.resolve() }
  177. },
  178. })
  179. if (options.omitSessionQuery !== true && (options.sessionId !== undefined || options.observe !== undefined)) {
  180. const observe = options.observe ?? (() => Promise.reject(new SessionQueryError('missing', 'SESSION_QUERY_SESSION_NOT_FOUND')))
  181. ctx.provide('sessionQuery', { observeSession: () => observe() } as never)
  182. }
  183. if (options.omitPersistence !== true) {
  184. ctx.provide('sessionPersistence', {} as never)
  185. }
  186. return {
  187. ctx,
  188. output: () => ({ out, err, order: [...order] }),
  189. run: async () => {
  190. ctx.on('session/flush', () => { order.push('flush') })
  191. internals.stdout = { write: (chunk: string) => { out += chunk; return true } }
  192. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  193. if (options.readStdin !== undefined) internals.readStdin = options.readStdin
  194. const exited = new Promise<number>((resolve) => {
  195. ctx.provide('appExit', (code: number) => { order.push('exit'); resolve(code) })
  196. })
  197. if (options.prelive === true || options.preliveMeta !== undefined) {
  198. await ctx.agents.create({
  199. sessionId: brandString<SessionId>(options.sessionId ?? 'session-exact'),
  200. meta: { cwd: process.cwd(), ...options.preliveMeta },
  201. })
  202. }
  203. apply(ctx, {
  204. ...options.useStdin === true ? {} : { task: options.task ?? 'do the thing' },
  205. ...options.sessionId === undefined ? {} : { sessionId: options.sessionId },
  206. ...options.json === undefined ? {} : { json: options.json },
  207. })
  208. return { code: await exited, out, err, order }
  209. },
  210. }
  211. }
  212. describe('headless runner', () => {
  213. it('records a fresh Session in the filesystem provider working directory', async () => {
  214. const cwd = '/remote/workspace'
  215. const test = await bench({
  216. before(session) { expect(session.header.cwd).toBe(cwd) },
  217. afterPrompt(session, message) { appendTurn(session, 1, message, 'remote answer', true) },
  218. }, { filesystemCwd: cwd })
  219. try { expect(await test.run()).toMatchObject({ code: 0, out: 'remote answer\n' }) }
  220. finally { await test.ctx.fiber.dispose() }
  221. })
  222. it('reports the provider cwd in its opening JSON event', async () => {
  223. const cwd = '/remote/workspace'
  224. const test = await bench({
  225. afterPrompt(session, message) { appendTurn(session, 1, message, 'remote answer', true) },
  226. }, { filesystemCwd: cwd, json: true })
  227. try {
  228. const result = await test.run()
  229. expect(result.code).toBe(0)
  230. expect(JSON.parse(result.out.split('\n')[0] as string)).toMatchObject({ type: 'session', cwd })
  231. } finally { await test.ctx.fiber.dispose() }
  232. })
  233. it('resumes against the provider cwd instead of the host launch directory', async () => {
  234. const cwd = '/remote/workspace'
  235. const test = await bench({
  236. afterPrompt(session, message) { appendTurn(session, 1, message, 'remote resumed', true) },
  237. }, {
  238. filesystemCwd: cwd, sessionId: 'session-exact',
  239. observe: async () => ({ header: { cwd, origin: 'user' }, events: [], [Symbol.dispose]() {} }),
  240. })
  241. test.ctx.sessions.create(brandString<SessionId>('session-exact'), { meta: { cwd } })
  242. try { expect(await test.run()).toMatchObject({ code: 0, out: 'remote resumed\n' }) }
  243. finally { await test.ctx.fiber.dispose() }
  244. })
  245. it('aggregates the final text across the complete idle-to-idle interval and flushes before exit', async () => {
  246. const test = await bench({
  247. before(session) {
  248. const setupMessage = {
  249. role: 'user', content: [{ type: 'text', text: 'setup' }], source: { kind: 'user' }, id: 'setup',
  250. } as UserMessage
  251. appendTurn(session, 0, setupMessage, 'pre-task noise', true)
  252. },
  253. async afterPrompt(session, message) {
  254. await Promise.resolve()
  255. appendTurn(session, 1, message, '', true)
  256. appendTurn(session, 2, message, 'final answer', true)
  257. },
  258. })
  259. const result = await test.run()
  260. expect(result).toEqual({
  261. code: 0,
  262. out: 'final answer\n',
  263. err: '',
  264. order: ['flush', 'exit'],
  265. })
  266. await test.ctx.fiber.dispose()
  267. })
  268. it('ignores durable inbox events before the first owned turn', async () => {
  269. const test = await bench({
  270. afterPrompt(session, message) {
  271. session.append('agent/inbox/spliced', {
  272. target: 'next-turn',
  273. start: 0,
  274. inserted: [message],
  275. })
  276. appendTurn(session, 1, message, 'answer after inbox activity', true)
  277. },
  278. })
  279. expect(await test.run()).toMatchObject({
  280. code: 0,
  281. out: 'answer after inbox activity\n',
  282. err: '',
  283. })
  284. await test.ctx.fiber.dispose()
  285. })
  286. it('waits for asynchronously appended events instead of racing Agent idleness', async () => {
  287. const test = await bench({
  288. afterPrompt: async (session, message) => {
  289. await new Promise(resolve => setTimeout(resolve, 5))
  290. appendTurn(session, 1, message, 'race-free answer', true)
  291. },
  292. })
  293. expect(await test.run()).toMatchObject({ code: 0, out: 'race-free answer\n', err: '' })
  294. await test.ctx.fiber.dispose()
  295. })
  296. it('streams reasoning before the Agent becomes idle and terminates its stderr line', async () => {
  297. const reasoningAppended = Promise.withResolvers<undefined>()
  298. const release = Promise.withResolvers<undefined>()
  299. const test = await bench({
  300. async afterPrompt(session, message, agent) {
  301. session.append('turn/start', { turn: 1 })
  302. session.append('step/start', { turn: 1, step: 1 })
  303. session.append('user/message', message, { surfaceOp: 'append' })
  304. startFrames(agent)
  305. emitChunk(agent, { type: 'block-start', index: 0, blockType: 'reasoning' })
  306. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: '' })
  307. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'checking the workspace' })
  308. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: ' safely\n' })
  309. emitChunk(agent, { type: 'block-end', index: 0, block: { type: 'reasoning', text: 'checking the workspace safely\n' } })
  310. emitChunk(agent, { type: 'usage', usage: { inputTokens: 1, outputTokens: 2, reasoningTokens: 2 } })
  311. emitChunk(agent, { type: 'block-start', index: 1, blockType: 'reasoning' })
  312. emitChunk(agent, { type: 'reasoning-delta', index: 1, text: 'second pass\n' })
  313. reasoningAppended.resolve(undefined)
  314. await release.promise
  315. emitChunk(agent, { type: 'block-start', index: 2, blockType: 'text' })
  316. emitChunk(agent, { type: 'text-delta', index: 2, text: 'done' })
  317. emitChunk(agent, { type: 'block-end', index: 2, block: { type: 'text', text: 'done' } })
  318. session.append('assistant/message', {
  319. stream: [],
  320. turn: 1,
  321. step: 1,
  322. message: createAssistantMessage({
  323. content: [{ type: 'text', text: 'done' }],
  324. source: { provider: 'test-provider', model: 'test-model' },
  325. }),
  326. }, { surfaceOp: 'append' })
  327. session.append('step/end', { turn: 1, step: 1 })
  328. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  329. },
  330. })
  331. const running = test.run()
  332. await reasoningAppended.promise
  333. const other = test.ctx.sessions.create()
  334. other.append('turn/start', { turn: 1 })
  335. other.append('step/start', { turn: 1, step: 1 })
  336. test.ctx.emit('agent/assistant-stream', {
  337. agent: { session: other } as Agent,
  338. frame: {
  339. type: 'chunk', attemptId: LlmAttemptId('other'), revision: 1,
  340. index: 0, time: Date.now(), chunk: { type: 'reasoning-delta', index: 0, text: 'other session' },
  341. },
  342. })
  343. const streamed = test.output()
  344. release.resolve(undefined)
  345. const result = await running
  346. expect(streamed).toEqual({
  347. out: '',
  348. err: 'dsh: reasoning:\nchecking the workspace safely\nsecond pass\n',
  349. order: [],
  350. })
  351. expect(result).toEqual({
  352. code: 0,
  353. out: 'done\n',
  354. err: 'dsh: reasoning:\nchecking the workspace safely\nsecond pass\n',
  355. order: ['flush', 'exit'],
  356. })
  357. await test.ctx.fiber.dispose()
  358. })
  359. it('closes an unterminated reasoning line as soon as the attempt ends', async () => {
  360. const reasoningAppended = Promise.withResolvers<undefined>()
  361. const releaseEnd = Promise.withResolvers<undefined>()
  362. const ended = Promise.withResolvers<undefined>()
  363. const finish = Promise.withResolvers<undefined>()
  364. const test = await bench({
  365. async afterPrompt(session, message, agent) {
  366. session.append('turn/start', { turn: 1 })
  367. session.append('step/start', { turn: 1, step: 1 })
  368. session.append('user/message', message, { surfaceOp: 'append' })
  369. startFrames(agent)
  370. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'unfinished reasoning' })
  371. reasoningAppended.resolve(undefined)
  372. await releaseEnd.promise
  373. const state = frameStates.get(agent)
  374. if (state === undefined) throw new Error('test Assistant frames have not started')
  375. agent.ctx.emit('agent/assistant-stream', {
  376. agent,
  377. frame: {
  378. type: 'end', attemptId: state.attemptId, revision: ++state.revision,
  379. index: state.index, outcome: { kind: 'abandoned' },
  380. },
  381. })
  382. ended.resolve(undefined)
  383. await finish.promise
  384. session.append('step/end', { turn: 1, step: 1 })
  385. session.append('turn/end', {
  386. turn: 1, reason: { kind: 'aborted', reason: { kind: 'user' } },
  387. })
  388. },
  389. })
  390. const running = test.run()
  391. await reasoningAppended.promise
  392. expect(test.output().err).toBe('dsh: reasoning:\nunfinished reasoning')
  393. releaseEnd.resolve(undefined)
  394. await ended.promise
  395. expect(test.output().err).toBe('dsh: reasoning:\nunfinished reasoning\n')
  396. finish.resolve(undefined)
  397. await expect(running).resolves.toMatchObject({ code: 1 })
  398. await test.ctx.fiber.dispose()
  399. })
  400. it('exits 1 when the final turn does not complete', async () => {
  401. const test = await bench({
  402. afterPrompt(session, message) { appendTurn(session, 1, message, undefined, false) },
  403. })
  404. expect(await test.run()).toMatchObject({ code: 1, out: '\n', err: '' })
  405. await test.ctx.fiber.dispose()
  406. })
  407. it('prints the durable model failure when the final turn ends in error', async () => {
  408. const test = await bench({
  409. afterPrompt(session, message) {
  410. session.append('turn/start', { turn: 1 })
  411. session.append('step/start', { turn: 1, step: 1 })
  412. session.append('user/message', message, { surfaceOp: 'append' })
  413. session.append('step/end', { turn: 1, step: 1 })
  414. session.append('turn/end', {
  415. turn: 1,
  416. reason: { kind: 'error', error: { code: 'SERVER', message: 'provider unavailable' } },
  417. })
  418. },
  419. })
  420. expect(await test.run()).toMatchObject({
  421. code: 1,
  422. out: '\n',
  423. err: 'dsh: SERVER: provider unavailable\n',
  424. })
  425. await test.ctx.fiber.dispose()
  426. })
  427. it('separates an unterminated reasoning prefix from the terminal model failure', async () => {
  428. const test = await bench({
  429. afterPrompt(session, message, agent) {
  430. session.append('turn/start', { turn: 1 })
  431. session.append('step/start', { turn: 1, step: 1 })
  432. session.append('user/message', message, { surfaceOp: 'append' })
  433. startFrames(agent)
  434. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'trying recovery' })
  435. session.append('step/end', { turn: 1, step: 1 })
  436. session.append('turn/end', {
  437. turn: 1,
  438. reason: { kind: 'error', error: { code: 'SERVER', message: 'provider unavailable' } },
  439. })
  440. },
  441. })
  442. expect(await test.run()).toMatchObject({
  443. code: 1,
  444. out: '\n',
  445. err: 'dsh: reasoning:\ntrying recovery\ndsh: SERVER: provider unavailable\n',
  446. })
  447. await test.ctx.fiber.dispose()
  448. })
  449. it('exits 1 when the owned interval contains no turn', async () => {
  450. const test = await bench({ afterPrompt: () => {} })
  451. expect(await test.run()).toMatchObject({ code: 1, out: '\n', err: '' })
  452. await test.ctx.fiber.dispose()
  453. })
  454. it('fails when an event below the captured Session length cannot be read', async () => {
  455. let capturedLength = 0
  456. const test = await bench({
  457. afterPrompt(session, message) {
  458. appendTurn(session, 1, message, 'unreachable', true)
  459. capturedLength = session.seq
  460. Object.defineProperty(session, 'eventAt', { value: () => undefined })
  461. },
  462. })
  463. const result = await test.run()
  464. expect(capturedLength).toBeGreaterThan(0)
  465. expect(result).toMatchObject({
  466. code: 1,
  467. out: '',
  468. err: `dsh: headless summary cannot read seq 0 below captured length ${String(capturedLength)}\n`,
  469. })
  470. await test.ctx.fiber.dispose()
  471. })
  472. it('reads the task from stdin when the invocation omits one', async () => {
  473. const test = await bench({
  474. afterPrompt(session, message) { appendTurn(session, 1, message, 'stdin answer', true) },
  475. }, {
  476. useStdin: true,
  477. readStdin: () => Promise.resolve('task from stdin'),
  478. })
  479. expect(await test.run()).toMatchObject({ code: 0, out: 'stdin answer\n', err: '' })
  480. await test.ctx.fiber.dispose()
  481. })
  482. it('rejects an empty stdin task', async () => {
  483. const test = await bench({ afterPrompt: () => {} }, {
  484. useStdin: true,
  485. readStdin: () => Promise.resolve(' \n'),
  486. })
  487. expect(await test.run()).toMatchObject({
  488. code: 1,
  489. err: 'dsh: a task is required, for example: dsh --profile headless "run the tests"\n',
  490. })
  491. await test.ctx.fiber.dispose()
  492. })
  493. it('reads the default process stdin when no override is installed', async () => {
  494. const original = Object.getOwnPropertyDescriptor(process, 'stdin')
  495. Object.defineProperty(process, 'stdin', {
  496. value: Readable.from([Buffer.from('piped'), Buffer.from(' task')]),
  497. configurable: true,
  498. })
  499. try {
  500. await expect(originalInternals.readStdin()).resolves.toBe('piped task')
  501. } finally {
  502. if (original !== undefined) Object.defineProperty(process, 'stdin', original)
  503. }
  504. })
  505. it('treats a bare dash positional as the stdin marker', async () => {
  506. const test = await bench({
  507. afterPrompt(session, message) { appendTurn(session, 1, message, 'dash answer', true) },
  508. }, {
  509. task: '-',
  510. readStdin: () => Promise.resolve('piped dash task'),
  511. })
  512. expect(await test.run()).toMatchObject({ code: 0, out: 'dash answer\n', err: '' })
  513. await test.ctx.fiber.dispose()
  514. })
  515. it('rejects --session-id when the query reports the id missing', async () => {
  516. const seen: string[] = []
  517. const test = await bench({
  518. afterPrompt(session, message) {
  519. seen.push(session.id)
  520. appendTurn(session, 1, message, 'created', true)
  521. },
  522. }, {
  523. sessionId: 'session-exact',
  524. observe: () => Promise.reject(new SessionQueryError('missing', 'SESSION_QUERY_SESSION_NOT_FOUND')),
  525. })
  526. const result = await test.run()
  527. expect(result.code).toBe(1)
  528. expect(result.err).toContain('session "session-exact" does not exist; omit --session-id to start a new Session')
  529. expect(result.out).toBe('')
  530. expect(seen).toEqual([])
  531. await test.ctx.fiber.dispose()
  532. })
  533. it('rejects --session-id when persistence is not mounted', async () => {
  534. const test = await bench({
  535. afterPrompt(session, message) { appendTurn(session, 1, message, 'created', true) },
  536. }, {
  537. sessionId: 'session-exact',
  538. observe: () => Promise.reject(new SessionQueryError('missing', 'SESSION_QUERY_SESSION_NOT_FOUND')),
  539. omitPersistence: true,
  540. })
  541. const result = await test.run()
  542. expect(result.code).toBe(1)
  543. expect(result.err).toContain('requires the sessionPersistence service')
  544. expect(result.out).toBe('')
  545. await test.ctx.fiber.dispose()
  546. })
  547. it('rejects adopting a live Session when persistence is not mounted', async () => {
  548. const test = await bench({
  549. afterPrompt(session, message) { appendTurn(session, 1, message, 'live', true) },
  550. }, {
  551. sessionId: 'session-exact',
  552. prelive: true,
  553. omitPersistence: true,
  554. })
  555. const result = await test.run()
  556. expect(result.code).toBe(1)
  557. expect(result.err).toContain('requires the sessionPersistence service')
  558. expect(result.out).toBe('')
  559. await test.ctx.fiber.dispose()
  560. })
  561. it('resumes the persisted Session when the query finds it', async () => {
  562. const test = await bench({
  563. afterPrompt(session, message) { appendTurn(session, 1, message, 'resumed answer', true) },
  564. }, {
  565. sessionId: 'session-exact',
  566. observe: () => Promise.resolve({
  567. header: { cwd: process.cwd(), origin: 'user' },
  568. events: [],
  569. [Symbol.dispose]() {},
  570. }),
  571. })
  572. const session = test.ctx.sessions.create(brandString<SessionId>('session-exact'), { meta: { cwd: process.cwd() } })
  573. const history = {
  574. role: 'user', content: [{ type: 'text', text: 'earlier' }], source: { kind: 'user' }, id: 'history',
  575. } as UserMessage
  576. appendTurn(session, 0, history, 'earlier answer', true)
  577. const before = session.seq
  578. expect(await test.run()).toMatchObject({ code: 0, out: 'resumed answer\n', err: '' })
  579. expect(session.seq).toBeGreaterThan(before)
  580. await test.ctx.fiber.dispose()
  581. })
  582. it('rejects a persisted Session recorded in another working directory', async () => {
  583. const test = await bench({ afterPrompt: () => {} }, {
  584. sessionId: 'session-exact',
  585. observe: () => Promise.resolve({
  586. header: { cwd: '/somewhere/else', origin: 'user' },
  587. events: [],
  588. [Symbol.dispose]() {},
  589. }),
  590. })
  591. const result = await test.run()
  592. expect(result.code).toBe(1)
  593. expect(result.err).toContain('was recorded in "/somewhere/else"')
  594. await test.ctx.fiber.dispose()
  595. })
  596. it('rejects a persisted Session created under an agent preset', async () => {
  597. const test = await bench({ afterPrompt: () => {} }, {
  598. sessionId: 'session-exact',
  599. observe: () => Promise.resolve({
  600. header: { cwd: process.cwd(), agentPreset: 'minimal' },
  601. events: [],
  602. [Symbol.dispose]() {},
  603. }),
  604. })
  605. const result = await test.run()
  606. expect(result.code).toBe(1)
  607. expect(result.err).toContain('runs under agent preset "minimal"')
  608. await test.ctx.fiber.dispose()
  609. })
  610. it('rejects a persisted Session that switched to an agent preset after creation', async () => {
  611. const test = await bench({ afterPrompt: () => {} }, {
  612. sessionId: 'session-exact',
  613. observe: () => Promise.resolve({
  614. header: { cwd: process.cwd() },
  615. events: [{ type: 'agent-preset/selected', data: { agentPreset: 'minimal' } }],
  616. [Symbol.dispose]() {},
  617. }),
  618. })
  619. const result = await test.run()
  620. expect(result.code).toBe(1)
  621. expect(result.err).toContain('runs under agent preset "minimal"')
  622. await test.ctx.fiber.dispose()
  623. })
  624. it('rejects a persisted Session whose preset record names no preset', async () => {
  625. const test = await bench({ afterPrompt: () => {} }, {
  626. sessionId: 'session-exact',
  627. observe: () => Promise.resolve({
  628. header: { cwd: process.cwd() },
  629. events: [{ type: 'agent-preset/selected', data: {} }],
  630. [Symbol.dispose]() {},
  631. }),
  632. })
  633. const result = await test.run()
  634. expect(result.code).toBe(1)
  635. expect(result.err).toContain('malformed agent-preset/selected event')
  636. await test.ctx.fiber.dispose()
  637. })
  638. it('rejects a persisted Session that recorded no working directory', async () => {
  639. const test = await bench({ afterPrompt: () => {} }, {
  640. sessionId: 'session-exact',
  641. observe: () => Promise.resolve({
  642. header: {},
  643. events: [],
  644. [Symbol.dispose]() {},
  645. }),
  646. })
  647. const result = await test.run()
  648. expect(result.code).toBe(1)
  649. expect(result.err).toContain('recorded no working directory')
  650. await test.ctx.fiber.dispose()
  651. })
  652. it('rejects a persisted Session owned by a subagent', async () => {
  653. const test = await bench({ afterPrompt: () => {} }, {
  654. sessionId: 'session-exact',
  655. observe: () => Promise.resolve({
  656. header: { cwd: process.cwd(), origin: 'subagent' },
  657. events: [],
  658. [Symbol.dispose]() {},
  659. }),
  660. })
  661. const result = await test.run()
  662. expect(result.code).toBe(1)
  663. expect(result.err).toContain('is a subagent or forked session')
  664. await test.ctx.fiber.dispose()
  665. })
  666. it('requires the Session query service for an exact Session identity', async () => {
  667. const test = await bench({ afterPrompt: () => {} }, { sessionId: 'session-exact', omitSessionQuery: true })
  668. const result = await test.run()
  669. expect(result.code).toBe(1)
  670. expect(result.err).toContain('requires the sessionQuery service')
  671. await test.ctx.fiber.dispose()
  672. })
  673. it('requires the Session query service even when a live Agent holds the identity', async () => {
  674. const test = await bench({
  675. afterPrompt(session, message) { appendTurn(session, 1, message, 'live', true) },
  676. }, {
  677. sessionId: 'session-exact',
  678. prelive: true,
  679. omitSessionQuery: true,
  680. })
  681. const result = await test.run()
  682. expect(result.code).toBe(1)
  683. expect(result.err).toContain('requires the sessionQuery service')
  684. expect(result.out).toBe('')
  685. await test.ctx.fiber.dispose()
  686. })
  687. it('rejects a whitespace-only session identity from configuration', async () => {
  688. const test = await bench({ afterPrompt: () => {} }, { sessionId: ' ' })
  689. const result = await test.run()
  690. expect(result.code).toBe(1)
  691. expect(result.err).toContain('sessionId must not be blank')
  692. expect(result.out).toBe('')
  693. await test.ctx.fiber.dispose()
  694. })
  695. it('refuses a live Agent identity it cannot own exclusively', async () => {
  696. const test = await bench({
  697. afterPrompt(session, message) { appendTurn(session, 1, message, 'live answer', true) },
  698. }, {
  699. sessionId: 'session-exact',
  700. prelive: true,
  701. })
  702. const result = await test.run()
  703. expect(result.code).toBe(1)
  704. expect(result.err).toContain('is live in this process, so the one-shot runner cannot own an exclusive run interval')
  705. expect(result.out).toBe('')
  706. await test.ctx.fiber.dispose()
  707. })
  708. it('rejects a live Agent recorded in another working directory', async () => {
  709. const test = await bench({ afterPrompt: () => {} }, {
  710. sessionId: 'session-exact',
  711. preliveMeta: { cwd: '/somewhere/else' },
  712. })
  713. const result = await test.run()
  714. expect(result.code).toBe(1)
  715. expect(result.err).toContain('was recorded in "/somewhere/else"')
  716. await test.ctx.fiber.dispose()
  717. })
  718. it('rejects a live Agent owned by a subagent', async () => {
  719. const test = await bench({ afterPrompt: () => {} }, {
  720. sessionId: 'session-exact',
  721. preliveMeta: { origin: 'subagent' },
  722. })
  723. const result = await test.run()
  724. expect(result.code).toBe(1)
  725. expect(result.err).toContain('is a subagent or forked session')
  726. await test.ctx.fiber.dispose()
  727. })
  728. it('rejects a live Agent created under an agent preset', async () => {
  729. const test = await bench({ afterPrompt: () => {} }, {
  730. sessionId: 'session-exact',
  731. preliveMeta: { agentPreset: 'minimal' },
  732. })
  733. const result = await test.run()
  734. expect(result.code).toBe(1)
  735. expect(result.err).toContain('runs under agent preset "minimal"')
  736. await test.ctx.fiber.dispose()
  737. })
  738. it('rejects a live Agent that switched to an agent preset while blank', async () => {
  739. const test = await bench({
  740. before(session) {
  741. const history = {
  742. role: 'user', content: [{ type: 'text', text: 'earlier' }], source: { kind: 'user' }, id: 'history',
  743. } as UserMessage
  744. appendTurn(session, 0, history, 'earlier answer', true)
  745. selectPreset(session, 'minimal')
  746. },
  747. afterPrompt: () => {},
  748. }, {
  749. sessionId: 'session-exact',
  750. prelive: true,
  751. })
  752. const result = await test.run()
  753. expect(result.code).toBe(1)
  754. expect(result.err).toContain('runs under agent preset "minimal"')
  755. await test.ctx.fiber.dispose()
  756. })
  757. it('rejects a preset appended after the observation snapshot was taken', async () => {
  758. const test = await bench({ afterPrompt: () => {} }, {
  759. sessionId: 'session-exact',
  760. observe: () => Promise.resolve({
  761. header: { cwd: process.cwd() },
  762. events: [],
  763. [Symbol.dispose]() {},
  764. }),
  765. })
  766. const session = test.ctx.sessions.create(brandString<SessionId>('session-exact'), { meta: { cwd: process.cwd() } })
  767. selectPreset(session, 'minimal')
  768. const result = await test.run()
  769. expect(result.code).toBe(1)
  770. expect(result.err).toContain('runs under agent preset "minimal"')
  771. await test.ctx.fiber.dispose()
  772. })
  773. it('rejects a preset an overlay appends while the runner awaits idle', async () => {
  774. const test = await bench({ afterPrompt: () => {} }, {
  775. sessionId: 'session-exact',
  776. observe: () => Promise.resolve({
  777. header: { cwd: process.cwd() },
  778. events: [],
  779. [Symbol.dispose]() {},
  780. }),
  781. onWhenIdle: (agent) => { selectPreset(agent.session, 'minimal') },
  782. })
  783. test.ctx.sessions.create(brandString<SessionId>('session-exact'), { meta: { cwd: process.cwd() } })
  784. const result = await test.run()
  785. expect(result.code).toBe(1)
  786. expect(result.err).toContain('runs under agent preset "minimal"')
  787. expect(result.out).toBe('')
  788. await test.ctx.fiber.dispose()
  789. })
  790. it('fails when a live event below the captured Session length cannot be read', async () => {
  791. let capturedLength = 0
  792. const test = await bench({
  793. before(session) {
  794. const history = {
  795. role: 'user', content: [{ type: 'text', text: 'earlier' }], source: { kind: 'user' }, id: 'history',
  796. } as UserMessage
  797. appendTurn(session, 0, history, 'earlier answer', true)
  798. capturedLength = session.seq
  799. Object.defineProperty(session, 'eventAt', { value: () => undefined })
  800. },
  801. afterPrompt: () => {},
  802. }, {
  803. sessionId: 'session-exact',
  804. prelive: true,
  805. })
  806. const result = await test.run()
  807. expect(capturedLength).toBeGreaterThan(0)
  808. expect(result.code).toBe(1)
  809. expect(result.err).toContain(`headless adoption cannot read seq 0 below captured length ${String(capturedLength)}`)
  810. await test.ctx.fiber.dispose()
  811. })
  812. it('bounds the error event message in --json mode', async () => {
  813. const test = await bench({ afterPrompt: () => {} }, {
  814. sessionId: 'session-exact',
  815. json: true,
  816. observe: () => Promise.reject(new SessionQueryError('x'.repeat(9 * 1024), 'SESSION_QUERY_CORRUPT_SESSION')),
  817. })
  818. const result = await test.run()
  819. const event = JSON.parse(result.out.trim()) as { message: string; truncated?: boolean }
  820. expect(event.truncated).toBe(true)
  821. expect(event.message.length).toBe(8 * 1024)
  822. await test.ctx.fiber.dispose()
  823. })
  824. it('rejects a persisted Session linked to a parent', async () => {
  825. const test = await bench({ afterPrompt: () => {} }, {
  826. sessionId: 'session-exact',
  827. observe: () => Promise.resolve({
  828. header: { cwd: process.cwd(), origin: 'user', parentSession: 'parent-1' },
  829. events: [],
  830. [Symbol.dispose]() {},
  831. }),
  832. })
  833. const result = await test.run()
  834. expect(result.code).toBe(1)
  835. expect(result.err).toContain('is a subagent or forked session')
  836. await test.ctx.fiber.dispose()
  837. })
  838. it('propagates a Session query failure that is not a missing log', async () => {
  839. const test = await bench({ afterPrompt: () => {} }, {
  840. sessionId: 'session-exact',
  841. observe: () => Promise.reject(new SessionQueryError('log is corrupt', 'SESSION_QUERY_CORRUPT_SESSION')),
  842. })
  843. const result = await test.run()
  844. expect(result.code).toBe(1)
  845. expect(result.err).toBe('dsh: log is corrupt\n')
  846. await test.ctx.fiber.dispose()
  847. })
  848. it('projects the run as ordered newline-delimited events in --json mode', async () => {
  849. const test = await bench({
  850. afterPrompt(session, message, agent) {
  851. session.append('turn/start', { turn: 1 })
  852. session.append('step/start', { turn: 1, step: 1 })
  853. session.append('user/message', message, { surfaceOp: 'append' })
  854. startFrames(agent)
  855. // A live attempt that never commits must not reach the projection.
  856. emitChunk(agent, { type: 'reasoning-delta', index: 0, text: 'discarded attempt' })
  857. emitChunk(agent, { type: 'text-delta', index: 0, text: 'discarded answer' })
  858. session.append('assistant/message', {
  859. stream: [],
  860. turn: 1,
  861. step: 1,
  862. usage: { inputTokens: 3, outputTokens: 4 },
  863. message: createAssistantMessage({
  864. content: [
  865. { type: 'reasoning', text: 'thinking hard' },
  866. { type: 'text', text: 'answer' },
  867. { type: 'tool-call', id: ToolCallId('call-1'), name: 'bash', arguments: '{"command":"ls"}' },
  868. ],
  869. source: { provider: 'test-provider', model: 'test-model' },
  870. }),
  871. }, { surfaceOp: 'append' })
  872. session.append('tool/call', {
  873. turn: 1, step: 1, callId: ToolCallId('call-1'), name: 'bash', arguments: '{"command":"ls"}',
  874. })
  875. session.append('tool/result', {
  876. turn: 1,
  877. step: 1,
  878. message: createToolResultMessage({
  879. callId: ToolCallId('call-1'),
  880. content: [{ type: 'text', text: 'a.txt' }],
  881. isError: false,
  882. }),
  883. }, { surfaceOp: 'append' })
  884. session.append('step/end', { turn: 1, step: 1 })
  885. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  886. },
  887. }, { json: true })
  888. const result = await test.run()
  889. const events = result.out.trim().split('\n').map(line => JSON.parse(line) as Record<string, unknown>)
  890. expect(events.map(event => event.type)).toEqual([
  891. 'session', 'status', 'status', 'thinking', 'text',
  892. 'tool_call', 'tool_result', 'status', 'status', 'final',
  893. ])
  894. expect(events[0]).toMatchObject({ type: 'session', cwd: process.cwd() })
  895. expect(typeof events[0]?.sessionId).toBe('string')
  896. expect(events[1]).toMatchObject({ type: 'status', phase: 'turn_start', turn: 1 })
  897. expect(events[3]).toMatchObject({ type: 'thinking', text: 'thinking hard' })
  898. expect(events[4]).toMatchObject({ type: 'text', text: 'answer' })
  899. expect(result.out).not.toContain('discarded')
  900. expect(events[5]).toMatchObject({ type: 'tool_call', callId: 'call-1', tool: 'bash', input: { command: 'ls' } })
  901. expect(events[6]).toMatchObject({ type: 'tool_result', callId: 'call-1', status: 'completed', result: 'a.txt' })
  902. expect(events[7]).toMatchObject({ type: 'status', phase: 'step_end', usage: { inputTokens: 3, outputTokens: 4 } })
  903. expect(events[8]).toMatchObject({ type: 'status', phase: 'turn_end', reason: { kind: 'completed' } })
  904. expect(events[9]).toMatchObject({ type: 'final', text: 'answer' })
  905. expect(result.err).toBe('')
  906. expect(result.code).toBe(0)
  907. await test.ctx.fiber.dispose()
  908. })
  909. it('reports a direct failure as an error event in --json mode', async () => {
  910. const test = await bench({ afterPrompt: () => {} }, {
  911. useStdin: true,
  912. readStdin: () => Promise.resolve(''),
  913. json: true,
  914. })
  915. const result = await test.run()
  916. expect(result.code).toBe(1)
  917. expect(JSON.parse(result.out.trim())).toMatchObject({ type: 'error' })
  918. expect(result.err).toContain('a task is required')
  919. await test.ctx.fiber.dispose()
  920. })
  921. it('reports a direct Agent creation failure', async () => {
  922. const ctx = new Context()
  923. let err = ''
  924. internals.stdout = { write: () => true }
  925. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  926. const exited = new Promise<number>((resolve) => {
  927. ctx.provide('appExit', resolve)
  928. })
  929. ctx.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  930. ctx.provide('sessions', { flush: () => Promise.resolve(true) } as never)
  931. ctx.provide('agents', { create: () => Promise.reject(new Error('factory exploded')) } as never)
  932. apply(ctx, { task: 't' })
  933. expect(await exited).toBe(1)
  934. expect(err).toBe('dsh: factory exploded\n')
  935. await ctx.fiber.dispose()
  936. })
  937. it('stringifies a non-Error Agent creation failure', async () => {
  938. const ctx = new Context()
  939. let err = ''
  940. internals.stdout = { write: () => true }
  941. internals.stderr = { write: (chunk: string) => { err += chunk; return true } }
  942. const exited = new Promise<number>((resolve) => {
  943. ctx.provide('appExit', resolve)
  944. })
  945. ctx.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  946. ctx.provide('sessions', { flush: () => Promise.resolve(true) } as never)
  947. const rejected = {
  948. then(_resolve: (value: never) => void, reject: (reason: unknown) => void): void {
  949. reject('factory exploded')
  950. },
  951. }
  952. ctx.provide('agents', { create: () => rejected } as never)
  953. apply(ctx, { task: 't' })
  954. expect(await exited).toBe(1)
  955. expect(err).toBe('dsh: factory exploded\n')
  956. await ctx.fiber.dispose()
  957. })
  958. it('abandons a run when the tree is disposed during Loader settlement', async () => {
  959. const ctx = new Context()
  960. let exited = false
  961. internals.stdout = { write: () => true }
  962. internals.stderr = { write: () => true }
  963. ctx.provide('appExit', () => { exited = true })
  964. const services = ctx.plugin((child: Context) => {
  965. child.provide('agentDefaultModel', { currentSelection: () => ({ provider: 'p', model: 'm' }) } as never)
  966. child.provide('sessions', {} as never)
  967. child.provide('agents', {} as never)
  968. })
  969. await services
  970. let release: () => void
  971. const settlement = new Promise<void>((resolve) => { release = resolve })
  972. ctx.provide('loader', { await: () => settlement } as never)
  973. apply(ctx, { task: 't' })
  974. await services.dispose()
  975. release!()
  976. await new Promise(resolve => setTimeout(resolve, 10))
  977. expect(exited).toBe(false)
  978. await ctx.fiber.dispose()
  979. })
  980. it('fails loud without the launcher-provided exit request', () => {
  981. const ctx = new Context()
  982. expect(() => { apply(ctx, { task: 't' }) }).toThrow('must provide ctx.appExit')
  983. })
  984. it('validates config: the task and run options are optional', () => {
  985. expect(new Config({})).toEqual({})
  986. expect(new Config({ task: 'x', sessionId: 'session-x', json: true }))
  987. .toEqual({ task: 'x', sessionId: 'session-x', json: true })
  988. })
  989. })