Merge PR #8150: perf(rpc): reuse serialized output frames (@MikeeI)
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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", () => {
|
||||
|
||||
Reference in New Issue
Block a user