diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 08cf5916d..b3005e664 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -127,7 +127,7 @@ - Fixed authoritative providers (e.g. `openai-codex`) keeping unsupported bundled models selectable when a fresh model cache and an expired OAuth token coincided: built-in discovery now forces the OAuth refresh so the provider's model manager is constructed and prunes stale bundled entries (e.g. `gpt-5.4-nano`) instead of waiting out the cache TTL. ([#5364](https://github.com/can1357/oh-my-pi/issues/5364)) ### Fixed -- Bounded RPC JSONL frames to 1 MiB, compacted `agent_end` to retain only messages not already streamed, and guaranteed worker reaping plus pending-request rejection after output-reader failures or explicit stops ([#5405](https://github.com/can1357/oh-my-pi/issues/5405)). +- Bounded RPC JSONL frames to 1 MiB, compacted oversized `agent_end` frames to retain only messages not already streamed, and guaranteed worker reaping plus pending-request rejection after output-reader failures or explicit stops ([#5405](https://github.com/can1357/oh-my-pi/issues/5405)). ## [17.0.4] - 2026-07-18 diff --git a/packages/coding-agent/src/modes/rpc/rpc-frame.ts b/packages/coding-agent/src/modes/rpc/rpc-frame.ts index 4d893202f..7e32bad3c 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-frame.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-frame.ts @@ -121,13 +121,16 @@ 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 { - const compacted = compactTerminalFrame(frame, streamedMessageCount, streamedMessages); - let json = JSON.stringify(compacted); + let json = JSON.stringify(frame); if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`; - if (isRecord(compacted) && compacted.type === "response") { - return `${JSON.stringify(overflowFrame(compacted))}\n`; + if (isRecord(frame) && frame.type === "response") { + return `${JSON.stringify(overflowFrame(frame))}\n`; } + const compacted = compactTerminalFrame(frame, streamedMessageCount, streamedMessages); + json = JSON.stringify(compacted); + if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`; + for (const pass of SHRINK_PASSES) { json = JSON.stringify(shrinkValue(compacted, pass)); if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`; diff --git a/packages/coding-agent/test/rpc-frame.test.ts b/packages/coding-agent/test/rpc-frame.test.ts index c9848f4ba..bdc07f9ca 100644 --- a/packages/coding-agent/test/rpc-frame.test.ts +++ b/packages/coding-agent/test/rpc-frame.test.ts @@ -5,6 +5,14 @@ function decode(frame: string): Record { return JSON.parse(frame) as Record; } +function oversizedMessageHistory(prefix: string) { + const payload = "x".repeat(1024); + return Array.from({ length: 1024 }, (_, index) => ({ + role: "assistant", + content: [{ type: "text", text: `${prefix}-${index}-${payload}` }], + })); +} + describe("RPC frame encoding", () => { it("preserves frames that already fit", () => { const frame = { id: "request-1", type: "response", command: "get_state", success: true, data: { ok: true } }; @@ -39,93 +47,77 @@ describe("RPC frame encoding", () => { expect(decoded).toEqual({ type: "agent_end", messages: [aborted], - messageCount: 1, }); }); - it("compacts stateful agent_end messages that match earlier message events", () => { + it("preserves terminal histories that fit for clients reading agent_end messages", () => { const streamed = { role: "assistant", content: [{ type: "text", text: "done" }] }; const encoder = new RpcFrameEncoder(); encoder.encode({ type: "agent_start" }); encoder.encode({ type: "message_end", message: streamed }); - const decoded = decode( - encoder.encode({ - type: "agent_end", - messages: [{ role: "assistant", content: [{ type: "text", text: "done" }] }], - }), - ); + const frame = { type: "agent_end", messages: [streamed] }; - expect(decoded).toEqual({ - type: "agent_end", - messages: [], - messageCount: 1, - }); + expect(encoder.encode(frame)).toBe(`${JSON.stringify(frame)}\n`); }); - it("matches terminal messages in the JSON shape sent by message_end", () => { + it("matches oversized terminal messages in the JSON shape sent by message_end", () => { + const messages = oversizedMessageHistory("wire-shape"); const encoder = new RpcFrameEncoder(); encoder.encode({ type: "agent_start" }); - encoder.encode({ - type: "message_end", - message: { - role: "assistant", - content: [{ type: "text", text: "done" }], - disabledFeatures: undefined, - toolCallAbortMessages: undefined, - }, - }); - const decoded = decode( + for (const message of messages) { encoder.encode({ - type: "agent_end", - messages: [{ role: "assistant", content: [{ type: "text", text: "done" }] }], - }), - ); + type: "message_end", + message: { + ...message, + disabledFeatures: undefined, + toolCallAbortMessages: undefined, + }, + }); + } + const encoded = encoder.encode({ type: "agent_end", messages }); - expect(decoded).toEqual({ + expect(Buffer.byteLength(encoded, "utf8")).toBeLessThanOrEqual(MAX_RPC_FRAME_BYTES); + expect(decode(encoded)).toEqual({ type: "agent_end", messages: [], - messageCount: 1, + messageCount: messages.length, }); }); it("does not let later mutation rewrite the message_end snapshot", () => { - const streamed = { role: "assistant", content: [{ type: "text", text: "before" }] }; + const messages = oversizedMessageHistory("before"); const encoder = new RpcFrameEncoder(); encoder.encode({ type: "agent_start" }); - encoder.encode({ type: "message_end", message: streamed }); - streamed.content[0].text = "after"; - const decoded = decode(encoder.encode({ type: "agent_end", messages: [streamed] })); + for (const message of messages) encoder.encode({ type: "message_end", message }); + messages[0].content[0].text = "after"; + const decoded = decode(encoder.encode({ type: "agent_end", messages })); - expect(decoded).toEqual({ - type: "agent_end", - messages: [streamed], - messageCount: 1, - }); + expect(decoded.messageCount).toBe(messages.length); + expect(Array.isArray(decoded.messages)).toBe(true); + expect((decoded.messages as unknown[]).length).toBeGreaterThan(0); }); it("keeps the active run snapshot when a continuing agent_end arrives late", () => { - const active = { role: "assistant", content: [{ type: "text", text: "active" }] }; + const active = oversizedMessageHistory("active"); const stale = { role: "assistant", content: [{ type: "text", text: "stale" }] }; const encoder = new RpcFrameEncoder(); encoder.encode({ type: "agent_start" }); - encoder.encode({ type: "message_end", message: active }); + for (const message of active) encoder.encode({ type: "message_end", message }); expect(decode(encoder.encode({ type: "agent_end", messages: [stale], willContinue: true }))).toEqual({ type: "agent_end", messages: [stale], - messageCount: 1, willContinue: true, }); - expect(decode(encoder.encode({ type: "agent_end", messages: [active] }))).toEqual({ + expect(decode(encoder.encode({ type: "agent_end", messages: active }))).toEqual({ type: "agent_end", messages: [], - messageCount: 1, - }); - expect(decode(encoder.encode({ type: "agent_end", messages: [active] }))).toEqual({ - type: "agent_end", - messages: [active], - messageCount: 1, + messageCount: active.length, }); + const replayed = decode(encoder.encode({ type: "agent_end", messages: active })); + expect(replayed.messageCount).toBe(active.length); + expect(Array.isArray(replayed.messages)).toBe(true); + expect((replayed.messages as unknown[]).length).toBeGreaterThan(0); }); it("bounds a single multi-byte message without losing its event discriminator", () => {