subagent-codex.spec.ts 39 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140
  1. import { PassThrough } from 'node:stream'
  2. import { Context } from 'cordis'
  3. import Loader from '@cordisjs/plugin-loader'
  4. import { describe, expect, it, vi } from 'vitest'
  5. import type { Agent } from '@deepseek-ai/dsh-agent'
  6. import type { InvariantInstaller } from '@deepseek-ai/dsh-invariants'
  7. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  8. import SubagentService from '@deepseek-ai/dsh-subagent'
  9. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  10. import type {
  11. SubprocessHandle,
  12. SubprocessOutcome,
  13. } from '@deepseek-ai/dsh-subprocess'
  14. import LocalSubprocessService from '@deepseek-ai/dsh-subprocess-local'
  15. import * as codex from '../src/index.ts'
  16. import * as invariant from '../src/invariant.ts'
  17. import {
  18. codexAppServerArgv,
  19. DEFAULT_DISPOSE_GRACE_MS,
  20. disposeCodexChild,
  21. startCodexRun,
  22. textTask,
  23. type CodexRunSpec,
  24. } from '../src/run.ts'
  25. import { CodexAppServerWire } from '../src/wire.ts'
  26. type JsonObject = Record<string, unknown>
  27. const fakeParent = {
  28. id: 'parent',
  29. session: { header: { cwd: process.cwd() } },
  30. } as unknown as Agent
  31. function request(
  32. prompt: ContentBlock[] = [{ type: 'text', text: 'do the task' }],
  33. signal = new AbortController().signal,
  34. ) {
  35. return { prompt, parent: fakeParent, signal }
  36. }
  37. async function nextTask(): Promise<void> {
  38. await new Promise<void>((resolve) => { setImmediate(resolve) })
  39. }
  40. class ProtocolPeer {
  41. private buffer = ''
  42. private readonly frames: JsonObject[] = []
  43. private readonly wakeups = new Set<() => void>()
  44. constructor(
  45. input: PassThrough,
  46. private readonly output: PassThrough,
  47. ) {
  48. input.on('data', (chunk: Buffer | string) => {
  49. this.buffer += chunk.toString()
  50. for (;;) {
  51. const newline = this.buffer.indexOf('\n')
  52. if (newline < 0) break
  53. const line = this.buffer.slice(0, newline)
  54. this.buffer = this.buffer.slice(newline + 1)
  55. if (line.trim().length > 0) this.frames.push(JSON.parse(line) as JsonObject)
  56. }
  57. for (const wake of this.wakeups) wake()
  58. this.wakeups.clear()
  59. })
  60. }
  61. async next(predicate: (frame: JsonObject) => boolean): Promise<JsonObject> {
  62. for (;;) {
  63. const index = this.frames.findIndex(predicate)
  64. if (index >= 0) return this.frames.splice(index, 1)[0]!
  65. await new Promise<void>((resolve) => { this.wakeups.add(resolve) })
  66. }
  67. }
  68. nextMethod(method: string): Promise<JsonObject> {
  69. return this.next(frame => frame.method === method)
  70. }
  71. nextResponse(id: unknown): Promise<JsonObject> {
  72. return this.next(frame => frame.id === id && frame.method === undefined)
  73. }
  74. send(...frames: readonly JsonObject[]): void {
  75. this.output.write(`${frames.map(frame => JSON.stringify(frame)).join('\n')}\n`)
  76. }
  77. respond(requestFrame: JsonObject, result: unknown): void {
  78. this.send({ id: requestFrame.id, result })
  79. }
  80. }
  81. interface FakeChildOptions {
  82. readonly pid?: number
  83. readonly exitOnTerminate?: boolean
  84. readonly doneError?: Error
  85. }
  86. interface FakeChild {
  87. readonly handle: SubprocessHandle
  88. readonly peer: ProtocolPeer
  89. readonly fromChild: PassThrough
  90. readonly toChild: PassThrough
  91. readonly settle: (outcome?: SubprocessOutcome) => void
  92. readonly fail: (error: Error) => void
  93. readonly terminate: () => void
  94. readonly waitForExit: (signal?: AbortSignal) => Promise<boolean>
  95. }
  96. function fakeChild(options: FakeChildOptions = {}): FakeChild {
  97. const fromChild = new PassThrough()
  98. const toChild = new PassThrough()
  99. const peer = new ProtocolPeer(toChild, fromChild)
  100. let exited = false
  101. let resolveDone!: (outcome: SubprocessOutcome) => void
  102. let rejectDone!: (error: Error) => void
  103. const done = new Promise<SubprocessOutcome>((resolve, reject) => {
  104. resolveDone = resolve
  105. rejectDone = reject
  106. })
  107. const settle = (
  108. outcome: SubprocessOutcome = { exitCode: 0, signal: null },
  109. ): void => {
  110. if (exited) return
  111. exited = true
  112. resolveDone(outcome)
  113. }
  114. const fail = (error: Error): void => {
  115. if (exited) return
  116. exited = true
  117. rejectDone(error)
  118. }
  119. if (options.doneError !== undefined) fail(options.doneError)
  120. const terminate = vi.fn(() => {
  121. if (options.exitOnTerminate !== false) settle()
  122. })
  123. const waitForExit = vi.fn(async (signal?: AbortSignal) => {
  124. if (exited) return true
  125. if (signal === undefined) {
  126. await done.catch(() => {})
  127. return true
  128. }
  129. return await new Promise<boolean>((resolve) => {
  130. const onAbort = (): void => { resolve(false) }
  131. signal.addEventListener('abort', onAbort, { once: true })
  132. void done.then(
  133. () => {
  134. signal.removeEventListener('abort', onAbort)
  135. resolve(true)
  136. },
  137. () => {
  138. signal.removeEventListener('abort', onAbort)
  139. resolve(true)
  140. },
  141. )
  142. })
  143. })
  144. const handle: SubprocessHandle = {
  145. pid: options.pid ?? 1234,
  146. stdin: toChild,
  147. stdout: fromChild,
  148. stderr: undefined,
  149. collected: {},
  150. done,
  151. terminate,
  152. waitForExit,
  153. }
  154. return {
  155. handle,
  156. peer,
  157. fromChild,
  158. toChild,
  159. settle,
  160. fail,
  161. terminate,
  162. waitForExit,
  163. }
  164. }
  165. function runSpec(
  166. child: FakeChild,
  167. overrides: Partial<CodexRunSpec> = {},
  168. ): CodexRunSpec {
  169. return {
  170. cwd: process.cwd(),
  171. env: {},
  172. disposeGraceMs: DEFAULT_DISPOSE_GRACE_MS,
  173. spawn: () => child.handle,
  174. ...overrides,
  175. }
  176. }
  177. async function initializeWire(): Promise<{
  178. readonly child: FakeChild
  179. readonly wire: CodexAppServerWire
  180. }> {
  181. const child = fakeChild()
  182. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  183. wire.start()
  184. const initializing = wire.initialize(new AbortController().signal)
  185. const initialize = await child.peer.nextMethod('initialize')
  186. child.peer.respond(initialize, { userAgent: 'codex-cli 0.146.0' })
  187. await initializing
  188. expect(await child.peer.nextMethod('initialized')).toEqual({
  189. jsonrpc: '2.0',
  190. method: 'initialized',
  191. })
  192. const starting = wire.startThread(process.cwd(), new AbortController().signal)
  193. const threadStart = await child.peer.nextMethod('thread/start')
  194. child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } })
  195. await starting
  196. return { child, wire }
  197. }
  198. async function publishRun(
  199. child = fakeChild(),
  200. signal = new AbortController().signal,
  201. specOverrides: Partial<CodexRunSpec> = {},
  202. ) {
  203. const starting = startCodexRun(request(undefined, signal), runSpec(child, specOverrides))
  204. const initialize = await child.peer.nextMethod('initialize')
  205. child.peer.respond(initialize, { userAgent: 'codex-cli 0.146.0' })
  206. await child.peer.nextMethod('initialized')
  207. const threadStart = await child.peer.nextMethod('thread/start')
  208. child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } })
  209. const run = await starting
  210. const turnStart = await child.peer.nextMethod('turn/start')
  211. return { child, run, turnStart }
  212. }
  213. function agentMessage(
  214. text: unknown,
  215. phase: unknown,
  216. turnId = 'turn-1',
  217. threadId = 'thread-1',
  218. ): JsonObject {
  219. return {
  220. method: 'item/completed',
  221. params: {
  222. threadId,
  223. turnId,
  224. item: { type: 'agentMessage', text, phase },
  225. },
  226. }
  227. }
  228. function turnCompleted(
  229. status: unknown,
  230. turnId = 'turn-1',
  231. threadId = 'thread-1',
  232. error: unknown = null,
  233. ): JsonObject {
  234. return {
  235. method: 'turn/completed',
  236. params: {
  237. threadId,
  238. turn: { id: turnId, status, error },
  239. },
  240. }
  241. }
  242. describe('task admission and package contracts', () => {
  243. it('resolves the fixed app-server command through the Windows npm shim boundary', () => {
  244. expect(codexAppServerArgv('win32')).toEqual([
  245. 'cmd.exe',
  246. '/d',
  247. '/s',
  248. '/c',
  249. 'codex',
  250. 'app-server',
  251. '--stdio',
  252. ])
  253. expect(codexAppServerArgv('linux')).toEqual(['codex', 'app-server', '--stdio'])
  254. })
  255. it('accepts one or more text blocks and rejects empty or non-text tasks', () => {
  256. expect(textTask([
  257. { type: 'text', text: 'one' },
  258. { type: 'text', text: 'two' },
  259. ])).toEqual(['one', 'two'])
  260. expect(() => textTask([])).toThrow('only text blocks')
  261. expect(() => textTask([{ type: 'reasoning', text: 'hidden' }]))
  262. .toThrow('only text blocks')
  263. expect(() => textTask([{ type: 'text', text: ' \n ' }]))
  264. .toThrow('must not be empty')
  265. })
  266. it('registers one fixed descriptor, validates config, and unregisters on HMR', async () => {
  267. const ctx = new Context()
  268. await ctx.plugin(SubagentService)
  269. await ctx.plugin(LocalSubprocessService)
  270. const fiber = await ctx.plugin(codex, {})
  271. const provider = ctx.subagents.getProvider('codex')!
  272. expect(provider).toMatchObject({
  273. name: 'codex',
  274. capabilities: {
  275. outputSchema: false,
  276. depthLimit: false,
  277. toolFilter: false,
  278. persona: false,
  279. },
  280. inheritsParentContext: false,
  281. })
  282. expect(ctx.subagents.list()).toEqual(['codex'])
  283. await fiber.dispose()
  284. expect(ctx.subagents.list()).toEqual([])
  285. for (const disposeGraceMs of [0, -1, Number.NaN, Number.POSITIVE_INFINITY]) {
  286. await expect(ctx.plugin(codex, { disposeGraceMs }))
  287. .rejects.toThrow('disposeGraceMs must be a positive finite number')
  288. }
  289. await expect(ctx.plugin(codex, { disposeGraceMs: MAX_TIMER_DELAY_MS + 1 }))
  290. .rejects.toThrow(`disposeGraceMs must be no greater than ${MAX_TIMER_DELAY_MS}`)
  291. await ctx.fiber.dispose()
  292. })
  293. it('requires a parent session cwd without suggesting unsupported config', async () => {
  294. const ctx = new Context()
  295. await ctx.plugin(SubagentService)
  296. await ctx.plugin(LocalSubprocessService)
  297. const spawn = vi.spyOn(ctx.subprocess, 'spawn')
  298. await ctx.plugin(codex, {})
  299. await expect(ctx.subagents.start('codex', {
  300. prompt: [{ type: 'text', text: 'task' }],
  301. parent: {
  302. id: 'parent-without-cwd',
  303. session: { header: {} },
  304. } as unknown as Agent,
  305. signal: new AbortController().signal,
  306. })).rejects.toThrow(
  307. 'subagent-codex: no working directory for the child — delegate from a parent session that has one',
  308. )
  309. expect(spawn).not.toHaveBeenCalled()
  310. await ctx.fiber.dispose()
  311. })
  312. it('keeps the namespace export shape and package-owned empty invariant', async () => {
  313. expect('default' in codex).toBe(false)
  314. expect(codex.name).toBe('subagent-codex')
  315. expect(codex.inject).toEqual(['subagents', 'subprocess'])
  316. const loader = Object.create(Loader.prototype) as Loader
  317. expect(loader.unwrapExports(codex)).toBe(codex)
  318. const dispose = vi.fn()
  319. const register = vi.fn((
  320. _packageName: string,
  321. _installer: InvariantInstaller,
  322. ) => dispose)
  323. const ctx = { invariants: { register } } as unknown as Context
  324. await expect(invariant.apply(ctx)).resolves.toBe(dispose)
  325. expect(register).toHaveBeenCalledWith(
  326. '@deepseek-ai/dsh-subagent-codex',
  327. expect.any(Function),
  328. )
  329. const install = register.mock.calls[0]![1]
  330. await install(new Context(), (message) => { throw new Error(message) })
  331. expect(invariant.name).toBe('subagent-codex-invariant')
  332. expect(invariant.inject).toEqual(['invariants'])
  333. })
  334. })
  335. describe('CodexAppServerWire', () => {
  336. it('sends the fixed handshake, thread, and turn payloads and keeps final_answer', async () => {
  337. const child = fakeChild()
  338. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  339. expect(wire.collectOutput()).toEqual([])
  340. wire.start()
  341. const initializing = wire.initialize(new AbortController().signal)
  342. const initialize = await child.peer.nextMethod('initialize')
  343. expect(initialize.params).toEqual({
  344. clientInfo: {
  345. name: 'deepseek-harness',
  346. title: 'DeepSeek Harness',
  347. version: '0.0.1',
  348. },
  349. capabilities: {
  350. experimentalApi: false,
  351. requestAttestation: false,
  352. },
  353. })
  354. child.peer.respond(initialize, { userAgent: 'codex-cli 0.146.0' })
  355. await initializing
  356. await child.peer.nextMethod('initialized')
  357. const starting = wire.startThread('/workspace', new AbortController().signal)
  358. const threadStart = await child.peer.nextMethod('thread/start')
  359. expect(threadStart.params).toEqual({ cwd: '/workspace', ephemeral: true })
  360. child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } })
  361. await starting
  362. const result = wire.runTurn(
  363. ['first', 'second'],
  364. new AbortController().signal,
  365. () => false,
  366. )
  367. const turnStart = await child.peer.nextMethod('turn/start')
  368. expect(turnStart.params).toEqual({
  369. threadId: 'thread-1',
  370. input: [
  371. { type: 'text', text: 'first', text_elements: [] },
  372. { type: 'text', text: 'second', text_elements: [] },
  373. ],
  374. })
  375. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  376. await nextTask()
  377. child.peer.send(
  378. {
  379. method: 'turn/started',
  380. params: { threadId: 'thread-1', turn: { id: 'turn-1' } },
  381. },
  382. agentMessage('other thread', 'final_answer', 'turn-1', 'thread-2'),
  383. agentMessage('other turn', 'final_answer', 'turn-2'),
  384. {
  385. method: 'item/completed',
  386. params: {
  387. threadId: 'thread-1',
  388. turnId: 'turn-1',
  389. item: { type: 'reasoning', text: 'not output' },
  390. },
  391. },
  392. agentMessage('commentary', 'commentary'),
  393. agentMessage('unphased', null),
  394. agentMessage('first final', 'final_answer'),
  395. agentMessage('last final', 'final_answer'),
  396. turnCompleted('completed'),
  397. )
  398. await expect(result).resolves.toEqual({
  399. output: [{ type: 'text', text: 'last final' }],
  400. stopReason: 'completed',
  401. })
  402. expect(wire.collectOutput()).toEqual([{ type: 'text', text: 'last final' }])
  403. wire.close()
  404. wire.close()
  405. })
  406. it('uses the last nullable-phase answer when no explicit final exists', async () => {
  407. const { child, wire } = await initializeWire()
  408. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  409. const turnStart = await child.peer.nextMethod('turn/start')
  410. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  411. child.peer.send(
  412. agentMessage('first', null),
  413. agentMessage('fallback', null),
  414. turnCompleted('completed'),
  415. )
  416. await expect(result).resolves.toEqual({
  417. output: [{ type: 'text', text: 'fallback' }],
  418. stopReason: 'completed',
  419. })
  420. wire.close()
  421. })
  422. it('maps only an explicit context-window failure to max-tokens', async () => {
  423. const { child, wire } = await initializeWire()
  424. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  425. const turnStart = await child.peer.nextMethod('turn/start')
  426. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  427. child.peer.send(
  428. agentMessage('partial answer', null),
  429. turnCompleted('failed', 'turn-1', 'thread-1', {
  430. message: 'too much context',
  431. codexErrorInfo: 'contextWindowExceeded',
  432. }),
  433. )
  434. await expect(result).resolves.toEqual({
  435. output: [{ type: 'text', text: 'partial answer' }],
  436. stopReason: 'max-tokens',
  437. })
  438. wire.close()
  439. })
  440. it('rejects invalid handshake, thread, and turn response shapes', async () => {
  441. {
  442. const child = fakeChild()
  443. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  444. wire.start()
  445. const pending = wire.initialize(new AbortController().signal)
  446. const frame = await child.peer.nextMethod('initialize')
  447. child.peer.respond(frame, null)
  448. await expect(pending).rejects.toThrow('invalid initialize response')
  449. wire.close()
  450. }
  451. {
  452. const child = fakeChild()
  453. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  454. wire.start()
  455. const pending = wire.startThread('/workspace', new AbortController().signal)
  456. const frame = await child.peer.nextMethod('thread/start')
  457. child.peer.respond(frame, { thread: { id: 'thread-1', ephemeral: false } })
  458. await expect(pending).rejects.toThrow('did not create an ephemeral thread')
  459. wire.close()
  460. }
  461. {
  462. const { child, wire } = await initializeWire()
  463. const pending = wire.runTurn(['task'], new AbortController().signal, () => false)
  464. const frame = await child.peer.nextMethod('turn/start')
  465. child.peer.respond(frame, { turn: { id: '' } })
  466. await expect(pending).rejects.toThrow('turn/start turn id')
  467. wire.close()
  468. }
  469. })
  470. it('fails closed for empty output, malformed messages, phases, and terminal status', async () => {
  471. const scenarios: Array<{
  472. readonly frames: JsonObject[]
  473. readonly message: string
  474. }> = [
  475. {
  476. frames: [turnCompleted('completed')],
  477. message: 'without a final answer',
  478. },
  479. {
  480. frames: [
  481. agentMessage('fallback', null),
  482. agentMessage(' \n ', 'final_answer'),
  483. turnCompleted('completed'),
  484. ],
  485. message: 'without a final answer',
  486. },
  487. {
  488. frames: [agentMessage(42, 'final_answer')],
  489. message: 'invalid agent message',
  490. },
  491. {
  492. frames: [agentMessage('answer', 'future_phase')],
  493. message: 'unknown agent message phase',
  494. },
  495. {
  496. frames: [turnCompleted('failed', 'turn-1', 'thread-1', { message: 'no' })],
  497. message: 'status failed',
  498. },
  499. {
  500. frames: [turnCompleted('interrupted')],
  501. message: 'status interrupted',
  502. },
  503. {
  504. frames: [turnCompleted('inProgress')],
  505. message: 'invalid terminal turn status',
  506. },
  507. ]
  508. for (const scenario of scenarios) {
  509. const { child, wire } = await initializeWire()
  510. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  511. const turnStart = await child.peer.nextMethod('turn/start')
  512. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  513. child.peer.send(...scenario.frames)
  514. await expect(result).rejects.toThrow(scenario.message)
  515. wire.close()
  516. }
  517. })
  518. it('fails closed when terminal notification params are not an object', async () => {
  519. const { child, wire } = await initializeWire()
  520. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  521. const turnStart = await child.peer.nextMethod('turn/start')
  522. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  523. child.peer.send({ method: 'turn/completed', params: null })
  524. await expect(result).rejects.toThrow('invalid turn/completed thread id')
  525. wire.close()
  526. })
  527. it('keeps an unsupported request authoritative over an early terminal in the same chunk', async () => {
  528. const { child, wire } = await initializeWire()
  529. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  530. const turnStart = await child.peer.nextMethod('turn/start')
  531. child.peer.send(
  532. { id: turnStart.id, result: { turn: { id: 'turn-1' } } },
  533. { id: 'future-request', method: 'future/request', params: {} },
  534. agentMessage('early answer', 'final_answer'),
  535. turnCompleted('completed'),
  536. )
  537. await expect(result).rejects.toThrow('unsupported app-server request')
  538. wire.close()
  539. })
  540. it('gives local cancellation precedence over a remote completed turn', async () => {
  541. const { child, wire } = await initializeWire()
  542. let cancelled = false
  543. const result = wire.runTurn(
  544. ['task'],
  545. new AbortController().signal,
  546. () => cancelled,
  547. )
  548. const turnStart = await child.peer.nextMethod('turn/start')
  549. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  550. cancelled = true
  551. child.peer.send(agentMessage('late', 'final_answer'), turnCompleted('completed'))
  552. await expect(result).resolves.toEqual({
  553. output: [{ type: 'text', text: 'late' }],
  554. stopReason: 'aborted',
  555. })
  556. wire.close()
  557. })
  558. it('answers all five unattended request classes without granting authority', async () => {
  559. const { child, wire } = await initializeWire()
  560. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  561. const turnStart = await child.peer.nextMethod('turn/start')
  562. child.peer.send({
  563. id: 'command',
  564. method: 'item/commandExecution/requestApproval',
  565. params: {
  566. threadId: 'thread-1',
  567. turnId: 'turn-1',
  568. availableDecisions: ['decline', 'cancel'],
  569. },
  570. })
  571. expect(await child.peer.nextResponse('command')).toMatchObject({
  572. result: { decision: 'cancel' },
  573. })
  574. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  575. await nextTask()
  576. const requests = [
  577. {
  578. id: 'file',
  579. method: 'item/fileChange/requestApproval',
  580. params: {
  581. threadId: 'thread-1',
  582. turnId: 'turn-1',
  583. availableDecisions: ['decline'],
  584. },
  585. result: { decision: 'decline' },
  586. },
  587. {
  588. id: 'file-default',
  589. method: 'item/fileChange/requestApproval',
  590. params: { threadId: 'thread-1', turnId: 'turn-1' },
  591. result: { decision: 'decline' },
  592. },
  593. {
  594. id: 'permissions',
  595. method: 'item/permissions/requestApproval',
  596. params: { threadId: 'thread-1', turnId: 'turn-1' },
  597. result: { permissions: {}, scope: 'turn' },
  598. },
  599. {
  600. id: 'user-input',
  601. method: 'item/tool/requestUserInput',
  602. params: { threadId: 'thread-1', turnId: 'turn-1', questions: [] },
  603. result: { answers: {} },
  604. },
  605. {
  606. id: 'mcp',
  607. method: 'mcpServer/elicitation/request',
  608. params: { threadId: 'thread-1', turnId: null },
  609. result: { action: 'decline', content: null, _meta: null },
  610. },
  611. ] as const
  612. for (const serverRequest of requests) {
  613. child.peer.send(serverRequest)
  614. expect(await child.peer.nextResponse(serverRequest.id)).toMatchObject({
  615. result: serverRequest.result,
  616. })
  617. }
  618. child.peer.send(agentMessage('answer', 'final_answer'), turnCompleted('completed'))
  619. await expect(result).resolves.toMatchObject({ stopReason: 'completed' })
  620. wire.close()
  621. })
  622. it('fails the run on unknown requests or wrong request association', async () => {
  623. for (const serverRequest of [
  624. {
  625. id: 'unknown',
  626. method: 'future/request',
  627. params: { threadId: 'thread-1', turnId: 'turn-1' },
  628. },
  629. {
  630. id: 'approval',
  631. method: 'item/commandExecution/requestApproval',
  632. params: {
  633. threadId: 'thread-1',
  634. turnId: 'turn-1',
  635. availableDecisions: ['accept'],
  636. },
  637. },
  638. {
  639. id: 'malformed-approval',
  640. method: 'item/fileChange/requestApproval',
  641. params: {
  642. threadId: 'thread-1',
  643. turnId: 'turn-1',
  644. availableDecisions: 'decline',
  645. },
  646. },
  647. {
  648. id: 'thread',
  649. method: 'item/fileChange/requestApproval',
  650. params: { threadId: 'thread-2', turnId: 'turn-1' },
  651. },
  652. {
  653. id: 'turn',
  654. method: 'item/fileChange/requestApproval',
  655. params: { threadId: 'thread-1', turnId: 'turn-2' },
  656. },
  657. ]) {
  658. const { child, wire } = await initializeWire()
  659. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  660. const turnStart = await child.peer.nextMethod('turn/start')
  661. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  662. await nextTask()
  663. child.peer.send(serverRequest)
  664. const response = await child.peer.nextResponse(serverRequest.id)
  665. expect(response.error).toMatchObject({ code: -32603 })
  666. await expect(result).rejects.toThrow()
  667. wire.close()
  668. }
  669. })
  670. it('rejects conflicting early turn identities before accepting output', async () => {
  671. const { child, wire } = await initializeWire()
  672. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  673. const turnStart = await child.peer.nextMethod('turn/start')
  674. child.peer.send({
  675. method: 'turn/started',
  676. params: { threadId: 'thread-1', turn: { id: 'turn-early' } },
  677. })
  678. child.peer.respond(turnStart, { turn: { id: 'turn-response' } })
  679. await expect(result).rejects.toThrow('did not match the active turn')
  680. wire.close()
  681. })
  682. it('rejects conflicting early notifications and requests before turn/start', async () => {
  683. {
  684. const { child, wire } = await initializeWire()
  685. child.peer.send({
  686. id: 'too-early',
  687. method: 'item/fileChange/requestApproval',
  688. params: { threadId: 'thread-1', turnId: 'turn-1' },
  689. })
  690. const response = await child.peer.nextResponse('too-early')
  691. expect(response.error).toMatchObject({ code: -32603 })
  692. wire.close()
  693. }
  694. {
  695. const { child, wire } = await initializeWire()
  696. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  697. await child.peer.nextMethod('turn/start')
  698. child.peer.send(
  699. {
  700. method: 'turn/started',
  701. params: { threadId: 'thread-1', turn: { id: 'turn-1' } },
  702. },
  703. agentMessage('wrong', 'final_answer', 'turn-2'),
  704. )
  705. await expect(result).rejects.toThrow('conflicting turns')
  706. wire.close()
  707. }
  708. })
  709. it('interrupts only an active open turn and contains remote interrupt failure', async () => {
  710. const { child, wire } = await initializeWire()
  711. wire.interrupt()
  712. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  713. const turnStart = await child.peer.nextMethod('turn/start')
  714. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  715. await nextTask()
  716. wire.interrupt()
  717. const interrupt = await child.peer.nextMethod('turn/interrupt')
  718. expect(interrupt.params).toEqual({ threadId: 'thread-1', turnId: 'turn-1' })
  719. child.peer.send({
  720. id: interrupt.id,
  721. error: { code: -32000, message: 'already done' },
  722. })
  723. child.peer.send(agentMessage('answer', 'final_answer'), turnCompleted('completed'))
  724. await expect(result).resolves.toMatchObject({ stopReason: 'completed' })
  725. wire.close()
  726. wire.interrupt()
  727. })
  728. it('ignores unrelated and out-of-window notifications', async () => {
  729. const { child, wire } = await initializeWire()
  730. child.peer.send(
  731. {
  732. method: 'turn/started',
  733. params: { threadId: 'thread-2', turn: { id: 'turn-other' } },
  734. },
  735. {
  736. method: 'turn/started',
  737. params: { threadId: 'thread-1', turn: { id: 'turn-before' } },
  738. },
  739. agentMessage('before', 'final_answer'),
  740. { method: 'future/notification', params: {} },
  741. turnCompleted('completed'),
  742. turnCompleted('completed', 'turn-other', 'thread-2'),
  743. )
  744. await nextTask()
  745. const result = wire.runTurn(['task'], new AbortController().signal, () => false)
  746. const turnStart = await child.peer.nextMethod('turn/start')
  747. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  748. await nextTask()
  749. child.peer.send(
  750. agentMessage('wrong turn', 'final_answer', 'turn-2'),
  751. turnCompleted('completed', 'turn-2'),
  752. agentMessage('answer', 'final_answer'),
  753. turnCompleted('completed'),
  754. )
  755. await expect(result).resolves.toEqual({
  756. output: [{ type: 'text', text: 'answer' }],
  757. stopReason: 'completed',
  758. })
  759. wire.close()
  760. })
  761. it('rejects pending work on abort, EOF, and stream error', async () => {
  762. {
  763. const child = fakeChild()
  764. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  765. wire.start()
  766. const controller = new AbortController()
  767. controller.abort('pre-aborted')
  768. await expect(wire.initialize(controller.signal))
  769. .rejects.toThrow('app-server request aborted: pre-aborted')
  770. wire.close()
  771. }
  772. {
  773. const child = fakeChild()
  774. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  775. wire.start()
  776. const controller = new AbortController()
  777. const pending = wire.initialize(controller.signal)
  778. await child.peer.nextMethod('initialize')
  779. controller.abort(new Error('cancel initialize'))
  780. await expect(pending).rejects.toThrow('cancel initialize')
  781. wire.close()
  782. }
  783. {
  784. const child = fakeChild()
  785. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  786. wire.start()
  787. const pending = wire.initialize(new AbortController().signal)
  788. await child.peer.nextMethod('initialize')
  789. child.fromChild.end()
  790. await expect(pending).rejects.toThrow(/(?:protocol stream|JSON-RPC input) closed/)
  791. wire.close()
  792. }
  793. {
  794. const child = fakeChild()
  795. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  796. wire.start()
  797. const pending = wire.initialize(new AbortController().signal)
  798. await child.peer.nextMethod('initialize')
  799. child.fromChild.emit('error', new Error('stdout broke'))
  800. await expect(pending).rejects.toThrow('stdout broke')
  801. wire.close()
  802. }
  803. {
  804. const child = fakeChild()
  805. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  806. wire.start()
  807. const pending = wire.initialize(new AbortController().signal)
  808. await child.peer.nextMethod('initialize')
  809. child.toChild.emit('error', new Error('stdin broke'))
  810. await expect(pending).rejects.toThrow('stdin broke')
  811. wire.close()
  812. child.toChild.emit('error', new Error('late stdin close'))
  813. }
  814. })
  815. })
  816. describe('run lifecycle and quiescence', () => {
  817. it('spawns the fixed app-server, publishes after thread creation, and disposes once', async () => {
  818. const child = fakeChild()
  819. const spawn = vi.fn(() => child.handle)
  820. const starting = startCodexRun(
  821. request([{ type: 'text', text: 'task' }]),
  822. runSpec(child, { env: { OPENAI_API_KEY: 'fake' }, spawn }),
  823. )
  824. let published = false
  825. void starting.then(() => { published = true })
  826. const initialize = await child.peer.nextMethod('initialize')
  827. expect(published).toBe(false)
  828. child.peer.respond(initialize, { userAgent: 'codex-cli 0.146.0' })
  829. await child.peer.nextMethod('initialized')
  830. const threadStart = await child.peer.nextMethod('thread/start')
  831. expect(published).toBe(false)
  832. child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } })
  833. const run = await starting
  834. expect(spawn).toHaveBeenCalledWith({
  835. argv: codexAppServerArgv(),
  836. cwd: process.cwd(),
  837. stdio: { stdin: 'pipe', stdout: 'pipe', stderr: 'inherit' },
  838. graceMs: DEFAULT_DISPOSE_GRACE_MS,
  839. env: { OPENAI_API_KEY: 'fake' },
  840. })
  841. expect(run.localAgent).toBeUndefined()
  842. const turnStart = await child.peer.nextMethod('turn/start')
  843. child.peer.send(
  844. { id: turnStart.id, result: { turn: { id: 'turn-1' } } },
  845. agentMessage('answer', 'final_answer'),
  846. turnCompleted('completed'),
  847. )
  848. await expect(run.result).resolves.toEqual({
  849. output: [{ type: 'text', text: 'answer' }],
  850. stopReason: 'completed',
  851. })
  852. const disposal = run.dispose()
  853. expect(run.dispose()).toBe(disposal)
  854. await disposal
  855. await nextTask()
  856. expect(child.terminate).toHaveBeenCalledTimes(1)
  857. expect(child.waitForExit).toHaveBeenCalledTimes(1)
  858. })
  859. it('settles local cancellation immediately and sends best-effort interrupt', async () => {
  860. const controller = new AbortController()
  861. const { child, run, turnStart } = await publishRun(
  862. fakeChild(),
  863. controller.signal,
  864. )
  865. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  866. await nextTask()
  867. controller.abort(new Error('stop'))
  868. await expect(run.result).resolves.toEqual({
  869. output: [],
  870. stopReason: 'aborted',
  871. })
  872. expect(await child.peer.nextMethod('turn/interrupt')).toMatchObject({
  873. params: { threadId: 'thread-1', turnId: 'turn-1' },
  874. })
  875. await run.dispose()
  876. })
  877. it('flattens child exit and protocol failures after publication', async () => {
  878. const errors: string[] = []
  879. {
  880. const child = fakeChild({ exitOnTerminate: false })
  881. const { run } = await publishRun(child, undefined, {
  882. onError: (error) => { errors.push(error.message) },
  883. })
  884. child.settle({ exitCode: 9, signal: null })
  885. await expect(run.result).resolves.toEqual({ output: [], stopReason: 'error' })
  886. expect(errors.at(-1)).toContain('code 9')
  887. await run.dispose().catch(() => {})
  888. }
  889. {
  890. const child = fakeChild()
  891. const { run, turnStart } = await publishRun(child, undefined, {
  892. onError: () => { throw new Error('diagnostic sink') },
  893. })
  894. child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
  895. child.fromChild.end()
  896. await expect(run.result).resolves.toEqual({ output: [], stopReason: 'error' })
  897. await run.dispose()
  898. }
  899. })
  900. it('rejects before spawn when pre-aborted and rolls back startup failures', async () => {
  901. const controller = new AbortController()
  902. controller.abort()
  903. const spawn = vi.fn()
  904. await expect(startCodexRun(
  905. request(undefined, controller.signal),
  906. {
  907. cwd: process.cwd(),
  908. env: {},
  909. disposeGraceMs: 10,
  910. spawn,
  911. },
  912. )).rejects.toThrow('aborted before app-server startup')
  913. expect(spawn).not.toHaveBeenCalled()
  914. const child = fakeChild()
  915. const starting = startCodexRun(request(), runSpec(child))
  916. const initialize = await child.peer.nextMethod('initialize')
  917. child.peer.respond(initialize, null)
  918. await expect(starting).rejects.toThrow('invalid initialize response')
  919. expect(child.terminate).toHaveBeenCalledTimes(1)
  920. })
  921. it('rolls back an abort that wins immediately after thread creation', async () => {
  922. const controller = new AbortController()
  923. const child = fakeChild()
  924. const starting = startCodexRun(
  925. request(undefined, controller.signal),
  926. runSpec(child),
  927. )
  928. const initialize = await child.peer.nextMethod('initialize')
  929. child.peer.respond(initialize, { userAgent: 'codex-cli 0.146.0' })
  930. await child.peer.nextMethod('initialized')
  931. const threadStart = await child.peer.nextMethod('thread/start')
  932. child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } })
  933. controller.abort('startup race')
  934. await expect(starting).rejects.toThrow('aborted before run publication')
  935. expect(child.terminate).toHaveBeenCalledTimes(1)
  936. })
  937. it('rolls back a subprocess done rejection during startup', async () => {
  938. const child = fakeChild({ doneError: new Error('spawn observer failed') })
  939. const error: unknown = await startCodexRun(request(), runSpec(child)).then(
  940. () => undefined,
  941. (failure: unknown) => failure,
  942. )
  943. expect(error).toBeInstanceOf(AggregateError)
  944. if (!(error instanceof AggregateError)) {
  945. throw new Error('expected startup and rollback failures')
  946. }
  947. expect(error.errors).toEqual([
  948. expect.objectContaining({ message: 'spawn observer failed' }),
  949. expect.objectContaining({ message: 'spawn observer failed' }),
  950. ])
  951. expect(child.terminate).toHaveBeenCalledTimes(1)
  952. })
  953. it('keeps overlapping runs isolated', async () => {
  954. const first = fakeChild()
  955. const second = fakeChild()
  956. const runs = await Promise.all([
  957. publishRun(first),
  958. publishRun(second),
  959. ])
  960. for (const [index, entry] of runs.entries()) {
  961. const id = `turn-${index + 1}`
  962. entry.child.peer.send(
  963. { id: entry.turnStart.id, result: { turn: { id } } },
  964. agentMessage(`answer-${index + 1}`, 'final_answer', id),
  965. turnCompleted('completed', id),
  966. )
  967. }
  968. const results = await Promise.all(runs.map(entry => entry.run.result))
  969. expect(results.map(result => result.output)).toEqual([
  970. [{ type: 'text', text: 'answer-1' }],
  971. [{ type: 'text', text: 'answer-2' }],
  972. ])
  973. expect(runs[0].run.id).not.toBe(runs[1].run.id)
  974. await Promise.all(runs.map(entry => entry.run.dispose()))
  975. })
  976. it('uses the registered provider config and logs flattened errors', async () => {
  977. const ctx = new Context()
  978. await ctx.plugin(SubagentService)
  979. await ctx.plugin(LocalSubprocessService)
  980. const child = fakeChild()
  981. const spawn = vi.spyOn(ctx.subprocess, 'spawn').mockReturnValue(child.handle)
  982. const warnings: string[] = []
  983. ctx.logger.warn = ((message: unknown) => {
  984. warnings.push(String(message))
  985. }) as typeof ctx.logger.warn
  986. await ctx.plugin(codex, {
  987. env: { OPENAI_API_KEY: 'fake' },
  988. disposeGraceMs: 25,
  989. })
  990. const starting = ctx.subagents.start('codex', {
  991. prompt: [{ type: 'text', text: 'task' }],
  992. parent: fakeParent,
  993. signal: new AbortController().signal,
  994. })
  995. const initialize = await child.peer.nextMethod('initialize')
  996. child.peer.respond(initialize, { userAgent: 'codex-cli 0.146.0' })
  997. await child.peer.nextMethod('initialized')
  998. const threadStart = await child.peer.nextMethod('thread/start')
  999. child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } })
  1000. const run = await starting
  1001. await child.peer.nextMethod('turn/start')
  1002. child.settle({ exitCode: 1, signal: null })
  1003. await expect(run.result).resolves.toMatchObject({ stopReason: 'error' })
  1004. expect(spawn).toHaveBeenCalledWith(expect.objectContaining({
  1005. env: { OPENAI_API_KEY: 'fake' },
  1006. graceMs: 25,
  1007. cwd: process.cwd(),
  1008. }))
  1009. expect(warnings).toEqual([
  1010. expect.stringContaining('subagent-codex: child run failed (error):'),
  1011. ])
  1012. await run.dispose().catch(() => {})
  1013. await ctx.fiber.dispose()
  1014. })
  1015. })
  1016. describe('disposeCodexChild', () => {
  1017. it('closes stdin, terminates, and waits for the managed tree', async () => {
  1018. const child = fakeChild()
  1019. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  1020. const end = vi.spyOn(child.toChild, 'end')
  1021. await disposeCodexChild(wire, child.handle)
  1022. expect(end).toHaveBeenCalled()
  1023. expect(child.terminate).toHaveBeenCalledTimes(1)
  1024. expect(child.waitForExit).toHaveBeenCalledTimes(1)
  1025. expect(child.waitForExit).toHaveBeenCalledWith()
  1026. })
  1027. it('does not finish disposal before the managed tree exits', async () => {
  1028. const child = fakeChild({ exitOnTerminate: false })
  1029. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  1030. let disposed = false
  1031. const disposal = disposeCodexChild(wire, child.handle).then(() => {
  1032. disposed = true
  1033. })
  1034. await new Promise<void>((resolve) => { setImmediate(resolve) })
  1035. expect(disposed).toBe(false)
  1036. child.settle()
  1037. await disposal
  1038. expect(disposed).toBe(true)
  1039. })
  1040. it('contains a concurrently closed stdin error', async () => {
  1041. const child = fakeChild()
  1042. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  1043. vi.spyOn(child.toChild, 'end').mockImplementation(() => {
  1044. throw new Error('already closed')
  1045. })
  1046. await expect(disposeCodexChild(wire, child.handle))
  1047. .resolves.toBeUndefined()
  1048. })
  1049. it('handles a spawn-level failure with no process tree', async () => {
  1050. const child = fakeChild({
  1051. pid: -1,
  1052. doneError: new Error('spawn failed'),
  1053. })
  1054. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  1055. await expect(disposeCodexChild(wire, child.handle))
  1056. .resolves.toBeUndefined()
  1057. expect(child.terminate).not.toHaveBeenCalled()
  1058. expect(child.waitForExit).not.toHaveBeenCalled()
  1059. })
  1060. it('reports direct-child observer failure and accepts absent stdin', async () => {
  1061. {
  1062. const child = fakeChild({
  1063. doneError: new Error('close observer failed'),
  1064. })
  1065. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  1066. await expect(disposeCodexChild(wire, child.handle))
  1067. .rejects.toThrow('close observer failed')
  1068. }
  1069. {
  1070. const child = fakeChild()
  1071. const handle = { ...child.handle, stdin: undefined }
  1072. const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!)
  1073. await expect(disposeCodexChild(wire, handle)).resolves.toBeUndefined()
  1074. }
  1075. })
  1076. })