subagent-codex.spec.ts 38 KB

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