diff --git a/packages/coding-agent/src/modes/rpc/rpc-frame.ts b/packages/coding-agent/src/modes/rpc/rpc-frame.ts index f4a86a507..4d893202f 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-frame.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-frame.ts @@ -147,7 +147,7 @@ export class RpcFrameEncoder { if (frame.type === "message_end") { const snapshot = encodedMessageSnapshot(encoded); if (snapshot) this.#streamedMessages.push(snapshot.message); - } else if (frame.type === "agent_end") this.#streamedMessages = []; + } else if (frame.type === "agent_end" && frame.willContinue !== true) this.#streamedMessages = []; return encoded; } } diff --git a/packages/coding-agent/test/rpc-frame.test.ts b/packages/coding-agent/test/rpc-frame.test.ts index e08117416..c9848f4ba 100644 --- a/packages/coding-agent/test/rpc-frame.test.ts +++ b/packages/coding-agent/test/rpc-frame.test.ts @@ -103,6 +103,31 @@ describe("RPC frame encoding", () => { }); }); + it("keeps the active run snapshot when a continuing agent_end arrives late", () => { + const active = { role: "assistant", content: [{ type: "text", text: "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 }); + + 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({ + type: "agent_end", + messages: [], + messageCount: 1, + }); + expect(decode(encoder.encode({ type: "agent_end", messages: [active] }))).toEqual({ + type: "agent_end", + messages: [active], + messageCount: 1, + }); + }); + it("bounds a single multi-byte message without losing its event discriminator", () => { const encoded = encodeRpcFrame({ type: "message_end",