wire.ts 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703
  1. /**
  2. * Minimal Codex app-server 0.149.1 protocol adapter. The shared JSON-RPC
  3. * transport owns framing and request correlation; this module owns only the
  4. * product methods, current thread/turn association, unattended approval
  5. * responses, and terminal-answer selection.
  6. *
  7. * @module @deepseek-ai/dsh-subagent-codex/wire
  8. */
  9. import type { Readable, Writable } from 'node:stream'
  10. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  11. import type { SubagentResult } from '@deepseek-ai/dsh-subagent'
  12. import { JsonRpcLineTransport } from '@deepseek-ai/dsh-sdk-protocol'
  13. import type { CodexPermissionMode } from './run.ts'
  14. type JsonObject = Record<string, unknown>
  15. /** Product facts owned by the Codex wire after publication. */
  16. export interface CodexWireFailureFacts {
  17. readonly stage: 'turn-start' | 'turn'
  18. readonly category:
  19. | 'limit'
  20. | 'access-policy'
  21. | 'service'
  22. | 'transport'
  23. | 'product-error'
  24. | 'invalid-result'
  25. | 'unknown'
  26. readonly httpStatus?: number | undefined
  27. }
  28. const THREAD_PERMISSION_PARAMS: Readonly<Record<CodexPermissionMode, JsonObject>> = {
  29. never: { approvalPolicy: 'never' },
  30. 'approve-for-me': {
  31. approvalPolicy: 'on-request',
  32. approvalsReviewer: 'auto_review',
  33. sandbox: 'workspace-write',
  34. },
  35. 'dangerously-bypass-approvals-and-sandbox': {
  36. approvalPolicy: 'never',
  37. sandbox: 'danger-full-access',
  38. },
  39. }
  40. function object(value: unknown, label: string): JsonObject {
  41. if (value === null || typeof value !== 'object' || Array.isArray(value)) {
  42. throw new Error(`subagent-codex: app-server returned invalid ${label}`)
  43. }
  44. return value as JsonObject
  45. }
  46. function string(value: unknown, label: string): string {
  47. if (typeof value !== 'string' || value.length === 0) {
  48. throw new Error(`subagent-codex: app-server returned invalid ${label}`)
  49. }
  50. return value
  51. }
  52. function unattendedDecision(params: JsonObject): 'cancel' | 'decline' {
  53. const available = params.availableDecisions
  54. if (available === undefined || available === null) return 'decline'
  55. if (Array.isArray(available)) {
  56. if (available.includes('cancel')) return 'cancel'
  57. if (available.includes('decline')) return 'decline'
  58. }
  59. throw new Error('subagent-codex: app-server offered no unattended approval decision')
  60. }
  61. function numericHttpStatus(value: unknown): number | undefined {
  62. return typeof value === 'number'
  63. && Number.isInteger(value)
  64. && value >= 0
  65. && value <= 65_535
  66. ? value
  67. : undefined
  68. }
  69. interface ParsedFailureInfo {
  70. readonly category: CodexWireFailureFacts['category']
  71. readonly httpStatus?: number | undefined
  72. readonly maxTokens?: true
  73. readonly sandboxFailure?: true
  74. }
  75. function objectFailureInfo(value: JsonObject): ParsedFailureInfo {
  76. const keys = Object.keys(value)
  77. const category = keys[0]
  78. if (keys.length !== 1 || category === undefined) {
  79. return { category: 'unknown' }
  80. }
  81. const detail = value[category]
  82. if (detail === null || typeof detail !== 'object' || Array.isArray(detail)) {
  83. return { category: 'unknown' }
  84. }
  85. const fields = detail as JsonObject
  86. switch (category) {
  87. case 'httpConnectionFailed':
  88. case 'responseStreamConnectionFailed':
  89. case 'responseStreamDisconnected':
  90. case 'responseTooManyFailedAttempts':
  91. {
  92. const httpStatus = numericHttpStatus(fields.httpStatusCode)
  93. return httpStatus === undefined
  94. ? { category: 'transport' }
  95. : { category: 'transport', httpStatus }
  96. }
  97. case 'activeTurnNotSteerable':
  98. return { category: 'product-error' }
  99. default:
  100. return { category: 'unknown' }
  101. }
  102. }
  103. function failureInfo(turn: JsonObject): ParsedFailureInfo {
  104. if (turn.status !== 'failed') return { category: 'unknown' }
  105. const error = turn.error
  106. if (error === null || typeof error !== 'object' || Array.isArray(error)) {
  107. return { category: 'unknown' }
  108. }
  109. const info = (error as JsonObject).codexErrorInfo
  110. if (typeof info === 'string') {
  111. switch (info) {
  112. case 'contextWindowExceeded':
  113. return { category: 'limit', maxTokens: true }
  114. case 'sessionBudgetExceeded':
  115. case 'usageLimitExceeded':
  116. return { category: 'limit' }
  117. case 'serverOverloaded':
  118. case 'internalServerError':
  119. return { category: 'service' }
  120. case 'cyberPolicy':
  121. case 'misalignmentPolicyViolation':
  122. case 'unauthorized':
  123. return { category: 'access-policy' }
  124. case 'badRequest':
  125. case 'threadRollbackFailed':
  126. case 'other':
  127. return { category: 'product-error' }
  128. case 'sandboxError':
  129. return { category: 'access-policy', sandboxFailure: true }
  130. default:
  131. return { category: 'unknown' }
  132. }
  133. }
  134. return info !== null && typeof info === 'object' && !Array.isArray(info)
  135. ? objectFailureInfo(info as JsonObject)
  136. : { category: 'unknown' }
  137. }
  138. function unattendedDiagnostic(
  139. mode: CodexPermissionMode,
  140. request: 'command approval' | 'file approval' | 'permission grant' | 'user input' | 'MCP elicitation' | 'command execution' | 'file change' | 'sandbox execution',
  141. decision: 'cancelled' | 'declined' | 'denied' | 'empty response' | 'failed',
  142. reason: string,
  143. ): string {
  144. return `Codex unattended decision (mode: ${mode}; request: ${request}; decision: ${decision}): ${reason}`
  145. }
  146. function thrown(value: unknown): Error {
  147. /* v8 ignore next -- typed protocol and stream failures reject with Error. */
  148. return value instanceof Error ? value : new Error(String(value))
  149. }
  150. function abortError(signal: AbortSignal): Error {
  151. return signal.reason instanceof Error
  152. ? signal.reason
  153. : new Error(`subagent-codex: app-server request aborted: ${String(signal.reason)}`)
  154. }
  155. async function raceAbort<T>(pending: Promise<T>, signal: AbortSignal): Promise<T> {
  156. if (signal.aborted) {
  157. void pending.catch(() => {})
  158. throw abortError(signal)
  159. }
  160. let rejectAbort!: (error: Error) => void
  161. const aborted = new Promise<never>((_resolve, reject) => { rejectAbort = reject })
  162. const onAbort = (): void => { rejectAbort(abortError(signal)) }
  163. signal.addEventListener('abort', onAbort, { once: true })
  164. try {
  165. return await Promise.race([pending, aborted])
  166. } finally {
  167. signal.removeEventListener('abort', onAbort)
  168. }
  169. }
  170. /**
  171. * One app-server connection and its single ephemeral thread/turn.
  172. *
  173. * The class deliberately exposes no generic request surface. Supporting
  174. * another product method must first become part of the provider contract.
  175. */
  176. export class CodexAppServerWire {
  177. private readonly transport: JsonRpcLineTransport
  178. private readonly fatal = Promise.withResolvers<never>()
  179. private threadId: string | undefined
  180. private turnId: string | undefined
  181. private pendingTurnId: string | undefined
  182. private turnCompleted: PromiseWithResolvers<{
  183. readonly params: JsonObject
  184. readonly order: number
  185. }> | undefined
  186. private readonly earlyTurnNotifications: Array<{
  187. readonly method: string
  188. readonly params: JsonObject
  189. readonly order: number
  190. }> = []
  191. private lastFinalAnswer: string | undefined
  192. private lastUnphasedAnswer: string | undefined
  193. private diagnostic: string | undefined
  194. private failure: CodexWireFailureFacts | undefined
  195. private diagnosticOrder = 0
  196. private observationOrder = 0
  197. private pendingDiagnostic: {
  198. readonly order: number
  199. readonly request: Parameters<typeof unattendedDiagnostic>[1]
  200. readonly decision: Parameters<typeof unattendedDiagnostic>[2]
  201. readonly reason: string
  202. } | undefined
  203. private inputEnded = false
  204. private terminalObserved = false
  205. private closed = false
  206. constructor(
  207. private readonly input: Readable,
  208. output: Writable,
  209. private readonly permissionMode: CodexPermissionMode,
  210. private readonly model?: string,
  211. ) {
  212. this.transport = new JsonRpcLineTransport(input, output)
  213. // Fatal protocol state can arrive after the current guarded operation has
  214. // already settled. Keep the shared rejection observed without inserting
  215. // another promise-adoption hop into active races.
  216. void this.fatal.promise.catch(() => {})
  217. this.transport.onRequest((method, params) => this.handleServerRequest(method, params))
  218. this.transport.onNotification((method, params) => {
  219. try {
  220. this.handleNotification(method, params)
  221. } catch (error: unknown) {
  222. this.fail(thrown(error))
  223. }
  224. })
  225. this.input.on('error', this.onInputError)
  226. this.input.on('end', this.onInputEnd)
  227. // Pipe errors can race protocol closure and process teardown. Retain both
  228. // error listeners for the lifetime of their per-run streams so no late
  229. // EPIPE or read failure becomes an unhandled EventEmitter error.
  230. output.on('error', this.onOutputError)
  231. }
  232. /** Start reading app-server frames. */
  233. start(): void {
  234. this.transport.start()
  235. }
  236. /**
  237. * Whether protocol output ended before a terminal turn notification.
  238. * @returns `true` only for an early protocol close without a terminal turn.
  239. */
  240. endedBeforeTerminal(): boolean {
  241. return this.inputEnded && !this.terminalObserved
  242. }
  243. /**
  244. * Perform the required app-server initialize/initialized handshake.
  245. * @param signal - unpublished-start cancellation.
  246. */
  247. async initialize(signal: AbortSignal): Promise<void> {
  248. object(await this.guarded(this.transport.request('initialize', {
  249. clientInfo: {
  250. name: 'deepseek-harness',
  251. title: 'DeepSeek Harness',
  252. version: '0.0.1',
  253. },
  254. capabilities: {
  255. experimentalApi: false,
  256. requestAttestation: false,
  257. },
  258. }, signal), signal), 'initialize response')
  259. this.transport.notify('initialized')
  260. await this.guarded(this.transport.flush(), signal)
  261. }
  262. /**
  263. * Create the run's private ephemeral thread and retain its identity.
  264. * @param cwd - parent Session workspace.
  265. * @param signal - unpublished-start cancellation.
  266. */
  267. async startThread(cwd: string, signal: AbortSignal): Promise<void> {
  268. const response = object(await this.guarded(this.transport.request('thread/start', {
  269. cwd,
  270. ephemeral: true,
  271. ...this.model === undefined ? {} : { model: this.model },
  272. ...THREAD_PERMISSION_PARAMS[this.permissionMode],
  273. }, signal), signal), 'thread/start response')
  274. const thread = object(response.thread, 'thread/start thread')
  275. const id = string(thread.id, 'thread/start thread id')
  276. if (thread.ephemeral !== true) {
  277. throw new Error('subagent-codex: app-server did not create an ephemeral thread')
  278. }
  279. this.threadId = id
  280. }
  281. /**
  282. * Submit the one text-only task and wait for this thread/turn's authoritative
  283. * terminal notification.
  284. * @param texts - already validated task text blocks.
  285. * @param signal - local cancellation for the published run.
  286. * @returns the shared subagent result.
  287. */
  288. async runTurn(
  289. texts: readonly string[],
  290. signal: AbortSignal,
  291. ): Promise<SubagentResult> {
  292. const completion = Promise.withResolvers<{
  293. readonly params: JsonObject
  294. readonly order: number
  295. }>()
  296. this.turnCompleted = completion
  297. const threadId = this.threadId as string
  298. try {
  299. const response = object(await this.guarded(this.transport.request('turn/start', {
  300. threadId,
  301. input: texts.map(text => ({ type: 'text', text, text_elements: [] })),
  302. }, signal), signal), 'turn/start response')
  303. const turn = object(response.turn, 'turn/start turn')
  304. this.commitTurnId(string(turn.id, 'turn/start turn id'))
  305. } catch (error: unknown) {
  306. this.recordFailure({ stage: 'turn-start', category: 'unknown' })
  307. throw error
  308. }
  309. let completed: {
  310. readonly params: JsonObject
  311. readonly order: number
  312. }
  313. let terminal: JsonObject
  314. try {
  315. completed = await this.guarded(completion.promise, signal)
  316. terminal = object(completed.params.turn, 'turn/completed turn')
  317. } catch (error: unknown) {
  318. this.recordFailure({ stage: 'turn', category: 'unknown' })
  319. throw error
  320. }
  321. const status = terminal.status
  322. if (status !== 'completed') {
  323. const parsed = failureInfo(terminal)
  324. this.recordFailure(parsed.httpStatus === undefined
  325. ? { stage: 'turn', category: parsed.category }
  326. : {
  327. stage: 'turn',
  328. category: parsed.category,
  329. httpStatus: parsed.httpStatus,
  330. })
  331. if (parsed.sandboxFailure) {
  332. this.recordDiagnostic(
  333. 'sandbox execution',
  334. 'failed',
  335. 'Codex reported a sandbox failure',
  336. completed.order,
  337. )
  338. }
  339. if (parsed.maxTokens) {
  340. return { output: this.collectOutput(), stopReason: 'max-tokens' }
  341. }
  342. const detail = status === 'failed' ? `: ${parsed.category}` : ''
  343. throw new Error(`subagent-codex: Codex turn ended with status ${String(status)}${detail}`)
  344. }
  345. const output = this.collectOutput()
  346. if (output.length === 0) {
  347. this.recordFailure({ stage: 'turn', category: 'invalid-result' })
  348. throw new Error('subagent-codex: Codex completed without a final answer')
  349. }
  350. return { output, stopReason: 'completed' }
  351. }
  352. /**
  353. * Best-effort remote cancellation. Local settlement and process teardown
  354. * remain authoritative when the child no longer accepts protocol requests.
  355. */
  356. interrupt(): void {
  357. if (this.threadId === undefined || this.turnId === undefined || this.closed) return
  358. void this.transport.request('turn/interrupt', {
  359. threadId: this.threadId,
  360. turnId: this.turnId,
  361. }).catch(() => {})
  362. }
  363. /**
  364. * The best non-commentary answer observed so far, preserving exact bytes.
  365. * @returns the selected final or nullable-phase text block, if any.
  366. */
  367. collectOutput(): ContentBlock[] {
  368. const selected = this.lastFinalAnswer ?? this.lastUnphasedAnswer
  369. return selected !== undefined && selected.trim().length > 0
  370. ? [{ type: 'text', text: selected }]
  371. : []
  372. }
  373. /**
  374. * The latest safe unattended permission fact observed for this run.
  375. * @returns provider-authored diagnostic text, when one was observed.
  376. */
  377. collectDiagnostic(): string | undefined {
  378. return this.diagnostic
  379. }
  380. /**
  381. * The structured failure fact observed for this published turn.
  382. * Call only after a non-completed return or rejection from {@link runTurn}.
  383. * @returns the fixed stage/category pair and optional HTTP status.
  384. */
  385. collectFailure(): CodexWireFailureFacts {
  386. return this.failure as CodexWireFailureFacts
  387. }
  388. /** Detach JSON-RPC listeners and reject outstanding requests. Idempotent. */
  389. close(): void {
  390. if (this.closed) return
  391. this.closed = true
  392. this.input.off('end', this.onInputEnd)
  393. this.transport.close()
  394. }
  395. private async guarded<T>(pending: Promise<T>, signal: AbortSignal): Promise<T> {
  396. const withFatal = Promise.race([this.fatal.promise, pending])
  397. return raceAbort(withFatal, signal)
  398. }
  399. private fail(error: Error): void {
  400. this.fatal.reject(error)
  401. }
  402. private readonly onInputError = (error: Error): void => {
  403. this.fail(error)
  404. }
  405. private readonly onOutputError = (error: Error): void => {
  406. this.fail(error)
  407. }
  408. private readonly onInputEnd = (): void => {
  409. this.inputEnded = true
  410. this.fail(new Error('subagent-codex: app-server protocol stream closed'))
  411. }
  412. private observePendingTurnId(id: string): void {
  413. if (this.turnCompleted === undefined) {
  414. throw new Error('subagent-codex: app-server referenced a turn before turn/start')
  415. }
  416. if (this.pendingTurnId !== undefined && this.pendingTurnId !== id) {
  417. throw new Error('subagent-codex: app-server referenced conflicting turns')
  418. }
  419. this.pendingTurnId = id
  420. }
  421. private commitTurnId(id: string): void {
  422. if (this.pendingTurnId !== undefined && this.pendingTurnId !== id) {
  423. throw new Error('subagent-codex: turn/start response did not match the active turn')
  424. }
  425. this.turnId = id
  426. const pendingDiagnostic = this.pendingDiagnostic
  427. this.pendingDiagnostic = undefined
  428. if (pendingDiagnostic !== undefined) {
  429. this.recordDiagnostic(
  430. pendingDiagnostic.request,
  431. pendingDiagnostic.decision,
  432. pendingDiagnostic.reason,
  433. pendingDiagnostic.order,
  434. )
  435. }
  436. const notifications = this.earlyTurnNotifications.splice(0)
  437. for (const notification of notifications) {
  438. this.handleNotification(
  439. notification.method,
  440. notification.params,
  441. notification.order,
  442. )
  443. }
  444. }
  445. /**
  446. * Validate the request's thread and turn association.
  447. * @returns `true` when the matching turn is still provisional, so the caller
  448. * defers its diagnostic until `commitTurnId()`.
  449. */
  450. private validateRunIds(
  451. params: JsonObject,
  452. nullableTurn = false,
  453. ): boolean {
  454. if (params.threadId !== this.threadId) {
  455. throw new Error('subagent-codex: app-server request referenced another thread')
  456. }
  457. if (nullableTurn && params.turnId === null) return false
  458. const id = string(params.turnId, 'server request turn id')
  459. if (this.turnId === undefined) {
  460. this.observePendingTurnId(id)
  461. return true
  462. }
  463. if (id !== this.turnId) {
  464. throw new Error('subagent-codex: app-server request referenced another turn')
  465. }
  466. return false
  467. }
  468. private recordRequestDiagnostic(
  469. provisional: boolean,
  470. request: Parameters<typeof unattendedDiagnostic>[1],
  471. decision: Parameters<typeof unattendedDiagnostic>[2],
  472. reason: string,
  473. ): void {
  474. const order = this.nextObservationOrder()
  475. if (provisional) {
  476. this.pendingDiagnostic = {
  477. order,
  478. request,
  479. decision,
  480. reason,
  481. }
  482. return
  483. }
  484. this.recordDiagnostic(request, decision, reason, order)
  485. }
  486. private recordDiagnostic(
  487. request: Parameters<typeof unattendedDiagnostic>[1],
  488. decision: Parameters<typeof unattendedDiagnostic>[2],
  489. reason: string,
  490. order = this.nextObservationOrder(),
  491. ): void {
  492. if (order < this.diagnosticOrder) return
  493. this.diagnosticOrder = order
  494. this.diagnostic = unattendedDiagnostic(
  495. this.permissionMode,
  496. request,
  497. decision,
  498. reason,
  499. )
  500. }
  501. private recordFailure(facts: CodexWireFailureFacts): void {
  502. this.failure = facts
  503. }
  504. private nextObservationOrder(): number {
  505. this.observationOrder += 1
  506. return this.observationOrder
  507. }
  508. private recordDeclinedItem(item: JsonObject, order?: number): boolean {
  509. if (item.type === 'commandExecution' && item.status === 'declined') {
  510. this.recordDiagnostic(
  511. 'command execution',
  512. 'declined',
  513. 'Codex declined the command under the selected permission mode',
  514. order,
  515. )
  516. return true
  517. }
  518. if (item.type === 'fileChange' && item.status === 'declined') {
  519. this.recordDiagnostic(
  520. 'file change',
  521. 'declined',
  522. 'Codex declined the file change under the selected permission mode',
  523. order,
  524. )
  525. return true
  526. }
  527. return false
  528. }
  529. private handleServerRequest(method: string, params: JsonObject): Promise<unknown> {
  530. try {
  531. switch (method) {
  532. case 'item/commandExecution/requestApproval':
  533. {
  534. const provisional = this.validateRunIds(params)
  535. const decision = unattendedDecision(params)
  536. this.recordRequestDiagnostic(
  537. provisional,
  538. 'command approval',
  539. decision === 'cancel' ? 'cancelled' : 'declined',
  540. 'the provider does not grant interactive approval',
  541. )
  542. return Promise.resolve({ decision })
  543. }
  544. case 'item/fileChange/requestApproval':
  545. {
  546. const provisional = this.validateRunIds(params)
  547. const decision = unattendedDecision(params)
  548. this.recordRequestDiagnostic(
  549. provisional,
  550. 'file approval',
  551. decision === 'cancel' ? 'cancelled' : 'declined',
  552. 'the provider does not grant interactive approval',
  553. )
  554. return Promise.resolve({ decision })
  555. }
  556. case 'item/permissions/requestApproval':
  557. this.recordRequestDiagnostic(
  558. this.validateRunIds(params),
  559. 'permission grant',
  560. 'denied',
  561. 'the provider grants no additional turn permissions',
  562. )
  563. return Promise.resolve({ permissions: {}, scope: 'turn' })
  564. case 'item/tool/requestUserInput':
  565. this.recordRequestDiagnostic(
  566. this.validateRunIds(params),
  567. 'user input',
  568. 'empty response',
  569. 'the provider does not collect interactive answers',
  570. )
  571. return Promise.resolve({ answers: {} })
  572. case 'mcpServer/elicitation/request':
  573. this.recordRequestDiagnostic(
  574. this.validateRunIds(params, true),
  575. 'MCP elicitation',
  576. 'declined',
  577. 'the provider does not collect interactive MCP input',
  578. )
  579. return Promise.resolve({ action: 'decline', content: null, _meta: null })
  580. default:
  581. throw new Error(`subagent-codex: unsupported app-server request ${JSON.stringify(method)}`)
  582. }
  583. } catch (error: unknown) {
  584. const normalized = thrown(error)
  585. this.fail(normalized)
  586. return Promise.reject(normalized)
  587. }
  588. }
  589. private handleNotification(
  590. method: string,
  591. params: JsonObject,
  592. order?: number,
  593. ): void {
  594. if (method === 'turn/started') {
  595. const threadId = string(params.threadId, 'turn/started thread id')
  596. if (threadId !== this.threadId) return
  597. const turn = object(params.turn, 'turn/started turn')
  598. if (this.turnCompleted !== undefined && this.turnId === undefined) {
  599. this.observePendingTurnId(string(turn.id, 'turn/started turn id'))
  600. }
  601. return
  602. }
  603. if (method === 'item/completed') {
  604. const threadId = string(params.threadId, 'item/completed thread id')
  605. if (threadId !== this.threadId) return
  606. const id = string(params.turnId, 'item/completed turn id')
  607. if (this.turnId === undefined) {
  608. if (this.turnCompleted !== undefined) {
  609. this.observePendingTurnId(id)
  610. this.earlyTurnNotifications.push({
  611. method,
  612. params,
  613. order: this.nextObservationOrder(),
  614. })
  615. }
  616. return
  617. }
  618. if (id !== this.turnId) return
  619. const item = object(params.item, 'item/completed item')
  620. if (this.recordDeclinedItem(item, order)) return
  621. if (item.type !== 'agentMessage') return
  622. const text = typeof item.text === 'string'
  623. ? item.text
  624. : (() => { throw new Error('subagent-codex: app-server returned an invalid agent message') })()
  625. if (item.phase === 'final_answer') {
  626. this.lastFinalAnswer = text
  627. } else if (item.phase === null) {
  628. this.lastUnphasedAnswer = text
  629. } else if (item.phase !== 'commentary') {
  630. throw new Error(`subagent-codex: app-server returned an unknown agent message phase ${JSON.stringify(item.phase)}`)
  631. }
  632. return
  633. }
  634. if (method !== 'turn/completed') return
  635. const threadId = string(params.threadId, 'turn/completed thread id')
  636. if (threadId !== this.threadId) return
  637. const turn = object(params.turn, 'turn/completed turn')
  638. const id = string(turn.id, 'turn/completed turn id')
  639. const turnCompleted = this.turnCompleted
  640. if (turnCompleted === undefined) return
  641. if (this.turnId === undefined) {
  642. this.observePendingTurnId(id)
  643. this.earlyTurnNotifications.push({
  644. method,
  645. params,
  646. order: this.nextObservationOrder(),
  647. })
  648. return
  649. }
  650. if (id !== this.turnId) return
  651. this.terminalObserved = true
  652. if (!['completed', 'interrupted', 'failed'].includes(String(turn.status))) {
  653. throw new Error(`subagent-codex: app-server returned invalid terminal turn status ${String(turn.status)}`)
  654. }
  655. turnCompleted.resolve({
  656. params,
  657. order: order ?? this.nextObservationOrder(),
  658. })
  659. }
  660. }