diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 980d80b70..6abf7d8a9 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -17,6 +17,7 @@ - Fixed Agent Hub lineage registration timestamps displaying in UTC instead of the user's local timezone. - Fixed Python/Julia/Ruby eval kernels failing to start after their staged runner script was cleared mid-session (e.g. a macOS tmpdir sweep): the memoized runner path is now re-validated so a long-lived process self-heals instead of only recovering on restart ([#8140](https://github.com/can1357/oh-my-pi/issues/8140)). - Fixed session resume fully reading and parsing the journal twice by reusing the entries already loaded by `SessionManager.open()` ([#8117](https://github.com/can1357/oh-my-pi/issues/8117)). +- Fixed RPC `message_end` frames being serialized more than once before output while preserving v1 and v2 wire bytes ([#8118](https://github.com/can1357/oh-my-pi/issues/8118)). ## [17.2.12] - 2026-08-08 diff --git a/packages/coding-agent/src/modes/rpc/rpc-frame.ts b/packages/coding-agent/src/modes/rpc/rpc-frame.ts index 862b654c8..6163806d3 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-frame.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-frame.ts @@ -239,9 +239,12 @@ function overflowFrame(frame: object): object { }; } -/** Serialize a complete JSONL frame while enforcing the transport byte ceiling. */ -export function encodeRpcFrame(frame: object, streamedMessageCount = 0, streamedMessages?: readonly unknown[]): string { - let json = JSON.stringify(frame); +function encodeRpcFrameFromJson( + frame: object, + json: string, + streamedMessageCount: number, + streamedMessages?: readonly unknown[], +): string { if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`; if (isRecord(frame) && frame.type === "response") { return `${JSON.stringify(overflowFrame(frame))}\n`; @@ -259,6 +262,11 @@ export function encodeRpcFrame(frame: object, streamedMessageCount = 0, streamed return `${JSON.stringify(overflowFrame(compacted))}\n`; } +/** Serialize a complete JSONL frame while enforcing the transport byte ceiling. */ +export function encodeRpcFrame(frame: object, streamedMessageCount = 0, streamedMessages?: readonly unknown[]): string { + return encodeRpcFrameFromJson(frame, JSON.stringify(frame), streamedMessageCount, streamedMessages); +} + /** Stateful encoder that tracks which messages a client has already received. */ export class RpcFrameEncoder { #streamedMessages: unknown[] = []; @@ -292,14 +300,14 @@ export class RpcFrameEncoder { frames = [singleFrame]; } } else { - singleFrame = encodeRpcFrame(frame, this.#streamedMessages.length, this.#streamedMessages); + singleFrame = encodeRpcFrameFromJson(frame, json, this.#streamedMessages.length, this.#streamedMessages); frames = [singleFrame]; } if (!isRecord(frame)) return frames; if (frame.type === "message_end") { const snapshot = this.#protocolVersion === 2 && Object.hasOwn(frame, "message") - ? { message: jsonSnapshot(frame.message) } + ? (encodedMessageSnapshot(json) ?? { message: jsonSnapshot(frame.message) }) : singleFrame !== undefined ? encodedMessageSnapshot(singleFrame) : undefined; diff --git a/packages/coding-agent/test/rpc-frame.test.ts b/packages/coding-agent/test/rpc-frame.test.ts index 2117054ba..8babaeccf 100644 --- a/packages/coding-agent/test/rpc-frame.test.ts +++ b/packages/coding-agent/test/rpc-frame.test.ts @@ -20,9 +20,26 @@ function oversizedMessageHistory(prefix: string) { } describe("RPC frame encoding", () => { - it("preserves frames that already fit", () => { + it("preserves fitting frames and serializes stateful message frames once", () => { const frame = { id: "request-1", type: "response", command: "get_state", success: true, data: { ok: true } }; expect(encodeRpcFrame(frame)).toBe(`${JSON.stringify(frame)}\n`); + + for (const version of [1, 2] as const) { + let messageReads = 0; + const message = { role: "assistant", content: [{ type: "text", text: "done" }] }; + const event = { + type: "message_end", + get message() { + messageReads++; + return message; + }, + }; + const encoder = new RpcFrameEncoder(); + encoder.setProtocolVersion(version); + + expect(decode(encoder.encode(event))).toEqual({ type: "message_end", message }); + expect(messageReads).toBe(1); + } }); it("compacts agent_end after message events have streamed", () => {