瀏覽代碼

refactor(subagent): simplify Codex diagnostic handoff

pku-xht 1 月之前
父節點
當前提交
0ff3c236ec

+ 2 - 12
packages/subagent/subagent-codex/src/run.ts

@@ -8,7 +8,7 @@
  */
 
 import { randomUUID } from 'node:crypto'
-import { writeSync } from 'node:fs'
+import { writeFileSync } from 'node:fs'
 import type { ContentBlock } from '@deepseek-ai/dsh-llm'
 import { SessionId } from '@deepseek-ai/dsh-session'
 import {
@@ -158,17 +158,7 @@ export async function startCodexRun(
     const bytes = typeof chunk === 'string' ? Buffer.from(chunk) : chunk
     wire.observeStderr(bytes.toString())
     try {
-      let offset = 0
-      while (offset < bytes.byteLength) {
-        const written = writeSync(
-          process.stderr.fd,
-          bytes,
-          offset,
-          bytes.byteLength - offset,
-        )
-        if (written <= 0) throw new Error('subagent-codex: host stderr made no write progress')
-        offset += written
-      }
+      writeFileSync(process.stderr.fd, bytes)
     } catch {
       // Host stderr is an observation sink, not a child-run failure authority.
     }

+ 11 - 13
packages/subagent/subagent-codex/src/wire.ts

@@ -167,7 +167,6 @@ export class CodexAppServerWire {
   private diagnosticOrder = 0
   private observationOrder = 0
   private pendingDiagnostic: {
-    readonly turnId: string
     readonly order: number
     readonly request: Parameters<typeof unattendedDiagnostic>[1]
     readonly decision: Parameters<typeof unattendedDiagnostic>[2]
@@ -399,7 +398,7 @@ export class CodexAppServerWire {
     this.turnId = id
     const pendingDiagnostic = this.pendingDiagnostic
     this.pendingDiagnostic = undefined
-    if (pendingDiagnostic?.turnId === id) {
+    if (pendingDiagnostic !== undefined) {
       this.recordDiagnostic(
         pendingDiagnostic.request,
         pendingDiagnostic.decision,
@@ -420,32 +419,31 @@ export class CodexAppServerWire {
   private validateRunIds(
     params: JsonObject,
     nullableTurn = false,
-  ): string | undefined {
+  ): boolean {
     if (params.threadId !== this.threadId) {
       throw new Error('subagent-codex: app-server request referenced another thread')
     }
-    if (nullableTurn && params.turnId === null) return undefined
+    if (nullableTurn && params.turnId === null) return false
     const id = string(params.turnId, 'server request turn id')
     if (this.turnId === undefined) {
       this.observePendingTurnId(id)
-      return id
+      return true
     }
     if (id !== this.turnId) {
       throw new Error('subagent-codex: app-server request referenced another turn')
     }
-    return undefined
+    return false
   }
 
   private recordRequestDiagnostic(
-    provisionalTurnId: string | undefined,
+    provisional: boolean,
     request: Parameters<typeof unattendedDiagnostic>[1],
     decision: Parameters<typeof unattendedDiagnostic>[2],
     reason: string,
   ): void {
     const order = this.nextObservationOrder()
-    if (provisionalTurnId !== undefined) {
+    if (provisional) {
       this.pendingDiagnostic = {
-        turnId: provisionalTurnId,
         order,
         request,
         decision,
@@ -504,10 +502,10 @@ export class CodexAppServerWire {
       switch (method) {
         case 'item/commandExecution/requestApproval':
         {
-          const provisionalTurnId = this.validateRunIds(params)
+          const provisional = this.validateRunIds(params)
           const decision = unattendedDecision(params)
           this.recordRequestDiagnostic(
-            provisionalTurnId,
+            provisional,
             'command approval',
             decision === 'cancel' ? 'cancelled' : 'declined',
             'the provider does not grant interactive approval',
@@ -516,10 +514,10 @@ export class CodexAppServerWire {
         }
         case 'item/fileChange/requestApproval':
         {
-          const provisionalTurnId = this.validateRunIds(params)
+          const provisional = this.validateRunIds(params)
           const decision = unattendedDecision(params)
           this.recordRequestDiagnostic(
-            provisionalTurnId,
+            provisional,
             'file approval',
             decision === 'cancel' ? 'cancelled' : 'declined',
             'the provider does not grant interactive approval',

+ 6 - 52
packages/subagent/subagent-codex/tests/subagent-codex.spec.ts

@@ -30,8 +30,6 @@ const { hostStderrWrite } = vi.hoisted(() => ({
   hostStderrWrite: {
     capture: false,
     failNext: false,
-    zeroNext: false,
-    maxBytesPerWrite: undefined as number | undefined,
     chunks: [] as Buffer[],
   },
 }))
@@ -40,17 +38,11 @@ vi.mock('node:fs', async (importOriginal) => {
   const actual = await importOriginal<typeof import('node:fs')>()
   return {
     ...actual,
-    writeSync(
+    writeFileSync(
       fd: number,
       value: string | Uint8Array,
-      offset?: number | null,
-      length?: number | null,
-    ): number {
+    ): void {
       if (fd === 2 && hostStderrWrite.capture) {
-        if (hostStderrWrite.zeroNext) {
-          hostStderrWrite.zeroNext = false
-          return 0
-        }
         if (hostStderrWrite.failNext) {
           hostStderrWrite.failNext = false
           throw Object.assign(new Error('host stderr broke'), { code: 'EIO' })
@@ -58,26 +50,10 @@ vi.mock('node:fs', async (importOriginal) => {
         const bytes = typeof value === 'string'
           ? Buffer.from(value)
           : Buffer.from(value.buffer, value.byteOffset, value.byteLength)
-        const start = typeof value === 'string' ? 0 : offset ?? 0
-        const requested = typeof value === 'string'
-          ? bytes.byteLength
-          : length ?? bytes.byteLength - start
-        const written = Math.min(
-          requested,
-          hostStderrWrite.maxBytesPerWrite ?? requested,
-        )
-        hostStderrWrite.chunks.push(Buffer.from(bytes.subarray(start, start + written)))
-        return written
+        hostStderrWrite.chunks.push(bytes)
+        return
       }
-      return typeof value === 'string'
-        ? actual.writeSync(fd, value, null, 'utf8')
-        : actual.writeSync(
-          fd,
-          value,
-          offset ?? 0,
-          length ?? value.byteLength - (offset ?? 0),
-          null,
-        )
+      actual.writeFileSync(fd, value)
     },
   }
 })
@@ -1390,7 +1366,6 @@ describe('run lifecycle and quiescence', () => {
   it('forwards stderr while extracting only a fixed safe permission signature', async () => {
     const child = fakeChild()
     hostStderrWrite.capture = true
-    hostStderrWrite.maxBytesPerWrite = 3
     hostStderrWrite.chunks.length = 0
     const { run, turnStart } = await publishRun(child)
     child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
@@ -1407,10 +1382,9 @@ describe('run lifecycle and quiescence', () => {
       stopReason: 'error',
     })
     expect(Buffer.concat(hostStderrWrite.chunks).toString()).toContain('SECRET_TOKEN')
-    expect(hostStderrWrite.chunks.length).toBeGreaterThan(3)
+    expect(hostStderrWrite.chunks).toHaveLength(3)
     await run.dispose()
     expect(child.stderr.listenerCount('data')).toBe(0)
-    hostStderrWrite.maxBytesPerWrite = undefined
     hostStderrWrite.capture = false
   })
 
@@ -1434,26 +1408,6 @@ describe('run lifecycle and quiescence', () => {
     hostStderrWrite.capture = false
   })
 
-  it('contains a zero-progress host stderr write without losing the diagnostic', async () => {
-    const child = fakeChild()
-    hostStderrWrite.capture = true
-    hostStderrWrite.zeroNext = true
-    const { run, turnStart } = await publishRun(child)
-    child.peer.respond(turnStart, { turn: { id: 'turn-1' } })
-    child.stderr.write('approval policy is Never; reject command')
-    child.peer.send(turnCompleted('failed', 'turn-1', 'thread-1', {
-      message: 'fixture terminal failure',
-      codexErrorInfo: 'badRequest',
-    }))
-    await expect(run.result).resolves.toEqual({
-      output: [],
-      diagnostic: 'Codex unattended decision (mode: never; request: command execution; decision: denied): Codex rejected an escalation because the selected policy never asks for approval',
-      stopReason: 'error',
-    })
-    await run.dispose()
-    hostStderrWrite.capture = false
-  })
-
   it('rejects before spawn when pre-aborted and rolls back startup failures', async () => {
     const controller = new AbortController()
     controller.abort()