wire.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374
  1. /**
  2. * Minimal Codex app-server 0.146.0 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. type JsonObject = Record<string, unknown>
  14. function object(value: unknown, label: string): JsonObject {
  15. if (value === null || typeof value !== 'object' || Array.isArray(value)) {
  16. throw new Error(`subagent-codex: app-server returned invalid ${label}`)
  17. }
  18. return value as JsonObject
  19. }
  20. function string(value: unknown, label: string): string {
  21. if (typeof value !== 'string' || value.length === 0) {
  22. throw new Error(`subagent-codex: app-server returned invalid ${label}`)
  23. }
  24. return value
  25. }
  26. function unattendedDecision(params: JsonObject): 'cancel' | 'decline' {
  27. const available = params.availableDecisions
  28. if (available === undefined || available === null) return 'decline'
  29. if (Array.isArray(available)) {
  30. if (available.includes('cancel')) return 'cancel'
  31. if (available.includes('decline')) return 'decline'
  32. }
  33. throw new Error('subagent-codex: app-server offered no unattended approval decision')
  34. }
  35. function isContextWindowExceeded(turn: JsonObject): boolean {
  36. if (turn.status !== 'failed') return false
  37. const error = turn.error
  38. return error !== null
  39. && typeof error === 'object'
  40. && !Array.isArray(error)
  41. && (error as JsonObject).codexErrorInfo === 'contextWindowExceeded'
  42. }
  43. function thrown(value: unknown): Error {
  44. /* v8 ignore next -- typed protocol and stream failures reject with Error. */
  45. return value instanceof Error ? value : new Error(String(value))
  46. }
  47. function abortError(signal: AbortSignal): Error {
  48. return signal.reason instanceof Error
  49. ? signal.reason
  50. : new Error(`subagent-codex: app-server request aborted: ${String(signal.reason)}`)
  51. }
  52. async function raceAbort<T>(pending: Promise<T>, signal: AbortSignal): Promise<T> {
  53. if (signal.aborted) {
  54. void pending.catch(() => {})
  55. throw abortError(signal)
  56. }
  57. let rejectAbort!: (error: Error) => void
  58. const aborted = new Promise<never>((_resolve, reject) => { rejectAbort = reject })
  59. const onAbort = (): void => { rejectAbort(abortError(signal)) }
  60. signal.addEventListener('abort', onAbort, { once: true })
  61. try {
  62. return await Promise.race([pending, aborted])
  63. } finally {
  64. signal.removeEventListener('abort', onAbort)
  65. }
  66. }
  67. /**
  68. * One app-server connection and its single ephemeral thread/turn.
  69. *
  70. * The class deliberately exposes no generic request surface. Supporting
  71. * another product method must first become part of the provider contract.
  72. */
  73. export class CodexAppServerWire {
  74. private readonly transport: JsonRpcLineTransport
  75. private readonly fatal = Promise.withResolvers<never>()
  76. private threadId: string | undefined
  77. private turnId: string | undefined
  78. private pendingTurnId: string | undefined
  79. private turnCompleted: PromiseWithResolvers<JsonObject> | undefined
  80. private readonly earlyTurnNotifications: Array<{
  81. readonly method: string
  82. readonly params: JsonObject
  83. }> = []
  84. private lastFinalAnswer: string | undefined
  85. private lastUnphasedAnswer: string | undefined
  86. private closed = false
  87. constructor(
  88. private readonly input: Readable,
  89. output: Writable,
  90. ) {
  91. this.transport = new JsonRpcLineTransport(input, output)
  92. // Fatal protocol state can arrive after the current guarded operation has
  93. // already settled. Keep the shared rejection observed without inserting
  94. // another promise-adoption hop into active races.
  95. void this.fatal.promise.catch(() => {})
  96. this.transport.onRequest((method, params) => this.handleServerRequest(method, params))
  97. this.transport.onNotification((method, params) => {
  98. try {
  99. this.handleNotification(method, params)
  100. } catch (error: unknown) {
  101. this.fail(thrown(error))
  102. }
  103. })
  104. this.input.on('error', this.onInputError)
  105. this.input.on('end', this.onInputEnd)
  106. // Pipe errors can race protocol closure and process teardown. Retain both
  107. // error listeners for the lifetime of their per-run streams so no late
  108. // EPIPE or read failure becomes an unhandled EventEmitter error.
  109. output.on('error', this.onOutputError)
  110. }
  111. /** Start reading app-server frames. */
  112. start(): void {
  113. this.transport.start()
  114. }
  115. /**
  116. * Perform the required app-server initialize/initialized handshake.
  117. * @param signal - unpublished-start cancellation.
  118. */
  119. async initialize(signal: AbortSignal): Promise<void> {
  120. object(await this.guarded(this.transport.request('initialize', {
  121. clientInfo: {
  122. name: 'deepseek-harness',
  123. title: 'DeepSeek Harness',
  124. version: '0.0.1',
  125. },
  126. capabilities: {
  127. experimentalApi: false,
  128. requestAttestation: false,
  129. },
  130. }, signal), signal), 'initialize response')
  131. this.transport.notify('initialized')
  132. await this.guarded(this.transport.flush(), signal)
  133. }
  134. /**
  135. * Create the run's private ephemeral thread and retain its identity.
  136. * @param cwd - parent Session workspace.
  137. * @param signal - unpublished-start cancellation.
  138. */
  139. async startThread(cwd: string, signal: AbortSignal): Promise<void> {
  140. const response = object(await this.guarded(this.transport.request('thread/start', {
  141. cwd,
  142. ephemeral: true,
  143. }, signal), signal), 'thread/start response')
  144. const thread = object(response.thread, 'thread/start thread')
  145. const id = string(thread.id, 'thread/start thread id')
  146. if (thread.ephemeral !== true) {
  147. throw new Error('subagent-codex: app-server did not create an ephemeral thread')
  148. }
  149. this.threadId = id
  150. }
  151. /**
  152. * Submit the one text-only task and wait for this thread/turn's authoritative
  153. * terminal notification.
  154. * @param texts - already validated task text blocks.
  155. * @param signal - local cancellation for the published run.
  156. * @returns the shared subagent result.
  157. */
  158. async runTurn(
  159. texts: readonly string[],
  160. signal: AbortSignal,
  161. ): Promise<SubagentResult> {
  162. const completion = Promise.withResolvers<JsonObject>()
  163. this.turnCompleted = completion
  164. const threadId = this.threadId as string
  165. const response = object(await this.guarded(this.transport.request('turn/start', {
  166. threadId,
  167. input: texts.map(text => ({ type: 'text', text, text_elements: [] })),
  168. }, signal), signal), 'turn/start response')
  169. const turn = object(response.turn, 'turn/start turn')
  170. this.commitTurnId(string(turn.id, 'turn/start turn id'))
  171. const completed = await this.guarded(completion.promise, signal)
  172. const terminal = object(completed.turn, 'turn/completed turn')
  173. const status = terminal.status
  174. if (isContextWindowExceeded(terminal)) {
  175. return { output: this.collectOutput(), stopReason: 'max-tokens' }
  176. }
  177. if (status !== 'completed') {
  178. const detail = status === 'failed'
  179. ? `: ${JSON.stringify(terminal.error)}`
  180. : ''
  181. throw new Error(`subagent-codex: Codex turn ended with status ${String(status)}${detail}`)
  182. }
  183. const output = this.collectOutput()
  184. if (output.length === 0) {
  185. throw new Error('subagent-codex: Codex completed without a final answer')
  186. }
  187. return { output, stopReason: 'completed' }
  188. }
  189. /**
  190. * Best-effort remote cancellation. Local settlement and process teardown
  191. * remain authoritative when the child no longer accepts protocol requests.
  192. */
  193. interrupt(): void {
  194. if (this.threadId === undefined || this.turnId === undefined || this.closed) return
  195. void this.transport.request('turn/interrupt', {
  196. threadId: this.threadId,
  197. turnId: this.turnId,
  198. }).catch(() => {})
  199. }
  200. /**
  201. * The best non-commentary answer observed so far, preserving exact bytes.
  202. * @returns the selected final or nullable-phase text block, if any.
  203. */
  204. collectOutput(): ContentBlock[] {
  205. const selected = this.lastFinalAnswer ?? this.lastUnphasedAnswer
  206. return selected !== undefined && selected.trim().length > 0
  207. ? [{ type: 'text', text: selected }]
  208. : []
  209. }
  210. /** Detach JSON-RPC listeners and reject outstanding requests. Idempotent. */
  211. close(): void {
  212. if (this.closed) return
  213. this.closed = true
  214. this.input.off('end', this.onInputEnd)
  215. this.transport.close()
  216. }
  217. private async guarded<T>(pending: Promise<T>, signal: AbortSignal): Promise<T> {
  218. const withFatal = Promise.race([this.fatal.promise, pending])
  219. return raceAbort(withFatal, signal)
  220. }
  221. private fail(error: Error): void {
  222. this.fatal.reject(error)
  223. }
  224. private readonly onInputError = (error: Error): void => {
  225. this.fail(error)
  226. }
  227. private readonly onOutputError = (error: Error): void => {
  228. this.fail(error)
  229. }
  230. private readonly onInputEnd = (): void => {
  231. this.fail(new Error('subagent-codex: app-server protocol stream closed'))
  232. }
  233. private observePendingTurnId(id: string): void {
  234. if (this.turnCompleted === undefined) {
  235. throw new Error('subagent-codex: app-server referenced a turn before turn/start')
  236. }
  237. if (this.pendingTurnId !== undefined && this.pendingTurnId !== id) {
  238. throw new Error('subagent-codex: app-server referenced conflicting turns')
  239. }
  240. this.pendingTurnId = id
  241. }
  242. private commitTurnId(id: string): void {
  243. if (this.pendingTurnId !== undefined && this.pendingTurnId !== id) {
  244. throw new Error('subagent-codex: turn/start response did not match the active turn')
  245. }
  246. this.turnId = id
  247. const notifications = this.earlyTurnNotifications.splice(0)
  248. for (const notification of notifications) {
  249. this.handleNotification(notification.method, notification.params)
  250. }
  251. }
  252. private validateRunIds(params: JsonObject, nullableTurn = false): void {
  253. if (params.threadId !== this.threadId) {
  254. throw new Error('subagent-codex: app-server request referenced another thread')
  255. }
  256. if (nullableTurn && params.turnId === null) return
  257. const id = string(params.turnId, 'server request turn id')
  258. if (this.turnId === undefined) {
  259. this.observePendingTurnId(id)
  260. return
  261. }
  262. if (id !== this.turnId) {
  263. throw new Error('subagent-codex: app-server request referenced another turn')
  264. }
  265. }
  266. private handleServerRequest(method: string, params: JsonObject): Promise<unknown> {
  267. try {
  268. switch (method) {
  269. case 'item/commandExecution/requestApproval':
  270. case 'item/fileChange/requestApproval':
  271. this.validateRunIds(params)
  272. return Promise.resolve({ decision: unattendedDecision(params) })
  273. case 'item/permissions/requestApproval':
  274. this.validateRunIds(params)
  275. return Promise.resolve({ permissions: {}, scope: 'turn' })
  276. case 'item/tool/requestUserInput':
  277. this.validateRunIds(params)
  278. return Promise.resolve({ answers: {} })
  279. case 'mcpServer/elicitation/request':
  280. this.validateRunIds(params, true)
  281. return Promise.resolve({ action: 'decline', content: null, _meta: null })
  282. default:
  283. throw new Error(`subagent-codex: unsupported app-server request ${JSON.stringify(method)}`)
  284. }
  285. } catch (error: unknown) {
  286. const normalized = thrown(error)
  287. this.fail(normalized)
  288. return Promise.reject(normalized)
  289. }
  290. }
  291. private handleNotification(method: string, params: JsonObject): void {
  292. if (method === 'turn/started') {
  293. const threadId = string(params.threadId, 'turn/started thread id')
  294. if (threadId !== this.threadId) return
  295. const turn = object(params.turn, 'turn/started turn')
  296. if (this.turnCompleted !== undefined && this.turnId === undefined) {
  297. this.observePendingTurnId(string(turn.id, 'turn/started turn id'))
  298. }
  299. return
  300. }
  301. if (method === 'item/completed') {
  302. const threadId = string(params.threadId, 'item/completed thread id')
  303. if (threadId !== this.threadId) return
  304. const id = string(params.turnId, 'item/completed turn id')
  305. if (this.turnId === undefined) {
  306. if (this.turnCompleted !== undefined) {
  307. this.observePendingTurnId(id)
  308. this.earlyTurnNotifications.push({ method, params })
  309. }
  310. return
  311. }
  312. if (id !== this.turnId) return
  313. const item = object(params.item, 'item/completed item')
  314. if (item.type !== 'agentMessage') return
  315. const text = typeof item.text === 'string'
  316. ? item.text
  317. : (() => { throw new Error('subagent-codex: app-server returned an invalid agent message') })()
  318. if (item.phase === 'final_answer') {
  319. this.lastFinalAnswer = text
  320. } else if (item.phase === null) {
  321. this.lastUnphasedAnswer = text
  322. } else if (item.phase !== 'commentary') {
  323. throw new Error(`subagent-codex: app-server returned an unknown agent message phase ${JSON.stringify(item.phase)}`)
  324. }
  325. return
  326. }
  327. if (method !== 'turn/completed') return
  328. const threadId = string(params.threadId, 'turn/completed thread id')
  329. if (threadId !== this.threadId) return
  330. const turn = object(params.turn, 'turn/completed turn')
  331. const id = string(turn.id, 'turn/completed turn id')
  332. const turnCompleted = this.turnCompleted
  333. if (turnCompleted === undefined) return
  334. if (this.turnId === undefined) {
  335. this.observePendingTurnId(id)
  336. this.earlyTurnNotifications.push({ method, params })
  337. return
  338. }
  339. if (id !== this.turnId) return
  340. if (!['completed', 'interrupted', 'failed'].includes(String(turn.status))) {
  341. throw new Error(`subagent-codex: app-server returned invalid terminal turn status ${String(turn.status)}`)
  342. }
  343. turnCompleted.resolve(params)
  344. }
  345. }