From 0ff3c236ecb6384302c20f729452af4d0cc9da5d Mon Sep 17 00:00:00 2001 From: pku-xht Date: Sat, 15 Aug 2026 20:09:04 +0800 Subject: [PATCH] refactor(subagent): simplify Codex diagnostic handoff --- packages/subagent/subagent-codex/src/run.ts | 14 +---- packages/subagent/subagent-codex/src/wire.ts | 24 ++++---- .../tests/subagent-codex.spec.ts | 58 ++----------------- 3 files changed, 19 insertions(+), 77 deletions(-) diff --git a/packages/subagent/subagent-codex/src/run.ts b/packages/subagent/subagent-codex/src/run.ts index 78f9a6fd9c..fdf467c876 100644 --- a/packages/subagent/subagent-codex/src/run.ts +++ b/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. } diff --git a/packages/subagent/subagent-codex/src/wire.ts b/packages/subagent/subagent-codex/src/wire.ts index a777b05331..cfb24a7481 100644 --- a/packages/subagent/subagent-codex/src/wire.ts +++ b/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[1] readonly decision: Parameters[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[1], decision: Parameters[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', diff --git a/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts b/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts index 1c63e11592..497303237a 100644 --- a/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts +++ b/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() 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()