feat(rpc): add negotiated lossless output framing
This commit is contained in:
+40
-3
@@ -25,15 +25,48 @@ Behavior notes:
|
||||
- RPC mode disables automatic session title generation by default to avoid an extra model call.
|
||||
- RPC mode resets workflow-altering `todo.*`, `task.*`, `memory.backend`/`memories.enabled`, `advisor.*`, `async.*`, and `bash.autoBackground.*` settings to their built-in defaults instead of inheriting user overrides.
|
||||
- The process reads stdin as JSONL (`readJsonl(Bun.stdin.stream())`).
|
||||
- At startup it writes `{ "type": "ready" }` before processing commands.
|
||||
- At startup it writes a `ready` frame before processing commands. The frame advertises supported protocol versions and transport limits.
|
||||
- When stdin closes, pending host-tool calls and host-URI requests are rejected and the process exits with code `0`.
|
||||
- Responses/events are written as one JSON object per line.
|
||||
|
||||
## Transport and Framing
|
||||
|
||||
Each frame is a single JSON object followed by `\n`.
|
||||
Protocol v1 frames are a single JSON object followed by `\n`. Every physical JSONL frame is limited to 1 MiB.
|
||||
|
||||
There is no envelope beyond the object shape itself.
|
||||
The initial ready frame uses protocol v1 and advertises the opt-in lossless transport:
|
||||
|
||||
```json
|
||||
{
|
||||
"type": "ready",
|
||||
"protocolVersion": 1,
|
||||
"supportedProtocolVersions": [1, 2],
|
||||
"maxFrameBytes": 1048576,
|
||||
"maxReassembledFrameBytes": 67108864
|
||||
}
|
||||
```
|
||||
|
||||
Clients that support protocol v2 SHOULD immediately send:
|
||||
|
||||
```json
|
||||
{ "id": "protocol-1", "type": "negotiate_protocol", "protocolVersion": 2 }
|
||||
```
|
||||
|
||||
After the success response, oversized stdout objects are emitted losslessly as an uninterrupted sequence of `rpc_chunk` frames. Each chunk carries a base64 segment of the original UTF-8 JSON object:
|
||||
|
||||
```json
|
||||
{
|
||||
"type": "rpc_chunk",
|
||||
"chunkId": "rpc-1",
|
||||
"index": 0,
|
||||
"count": 7,
|
||||
"byteLength": 1600042,
|
||||
"data": "eyJ0eXBlIjoicmVzcG9uc2UiLC4uLn0="
|
||||
}
|
||||
```
|
||||
|
||||
Clients MUST validate `chunkId`, `index`, `count`, and `byteLength`, reject interleaved or interrupted sequences, enforce the advertised reassembly limit, concatenate decoded bytes in index order, decode them as strict UTF-8, and parse the result as one JSON object. The exported `RpcFrameDecoder` implements this validation. `RpcClient` negotiates v2 automatically when the ready frame advertises it.
|
||||
|
||||
Legacy clients may ignore the added ready fields and remain on v1. V1 retains its bounded fallback behavior for oversized output. Frames above the v2 reassembly ceiling still fail explicitly; large history APIs should use pagination rather than depending on arbitrarily large logical frames.
|
||||
|
||||
### Outbound frame categories (stdout)
|
||||
|
||||
@@ -84,6 +117,10 @@ Important edge behavior from runtime:
|
||||
- `{ id?, type: "abort_and_prompt", message: string, images?: ImageContent[] }`
|
||||
- `{ id?, type: "new_session", parentSession?: string }`
|
||||
|
||||
### Protocol
|
||||
|
||||
- `{ id?, type: "negotiate_protocol", protocolVersion: 2 }`
|
||||
|
||||
### State
|
||||
|
||||
- `{ id?, type: "get_state" }`
|
||||
|
||||
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
|
||||
- Added opt-in RPC protocol v2 negotiation with bounded, lossless chunking for stdout objects up to 64 MiB. Legacy JSONL clients remain on protocol v1, while the TypeScript RPC client negotiates and reassembles v2 automatically.
|
||||
|
||||
## [17.0.8] - 2026-07-22
|
||||
|
||||
### Added
|
||||
|
||||
@@ -12,6 +12,7 @@ import { isRecord, ptree, readJsonl } from "@oh-my-pi/pi-utils";
|
||||
import type { FileSink } from "bun";
|
||||
import type { BashResult } from "../../exec/bash-executor";
|
||||
import type { AgentSessionEvent, SessionStats } from "../../session/agent-session";
|
||||
import { MAX_RPC_FRAME_BYTES, MAX_RPC_REASSEMBLED_BYTES, RpcFrameDecoder } from "./rpc-frame";
|
||||
import type {
|
||||
RpcAvailableCommandsUpdateFrame,
|
||||
RpcAvailableSlashCommand,
|
||||
@@ -136,6 +137,16 @@ function isRpcResponse(value: unknown): value is RpcResponse {
|
||||
return true;
|
||||
}
|
||||
|
||||
function supportsRpcProtocolV2(value: Record<string, unknown>): boolean {
|
||||
return (
|
||||
value.type === "ready" &&
|
||||
Array.isArray(value.supportedProtocolVersions) &&
|
||||
value.supportedProtocolVersions.includes(2) &&
|
||||
value.maxFrameBytes === MAX_RPC_FRAME_BYTES &&
|
||||
value.maxReassembledFrameBytes === MAX_RPC_REASSEMBLED_BYTES
|
||||
);
|
||||
}
|
||||
|
||||
function isAgentEvent(value: unknown): value is AgentEvent {
|
||||
if (!isRecord(value)) return false;
|
||||
const type = value.type;
|
||||
@@ -269,6 +280,9 @@ export class RpcClient {
|
||||
// Wait for the "ready" signal or process exit
|
||||
const { promise: readyPromise, resolve: readyResolve, reject: readyReject } = Promise.withResolvers<void>();
|
||||
let readySettled = false;
|
||||
let protocolV2Supported = false;
|
||||
let protocolV2Enabled = false;
|
||||
const frameDecoder = new RpcFrameDecoder();
|
||||
|
||||
const reapAfterOutputFailure = async (error: Error) => {
|
||||
if (this.#process !== child) return;
|
||||
@@ -294,11 +308,15 @@ export class RpcClient {
|
||||
void (async () => {
|
||||
for await (const line of lines) {
|
||||
if (!readySettled && isRecord(line) && line.type === "ready") {
|
||||
protocolV2Supported = supportsRpcProtocolV2(line);
|
||||
readySettled = true;
|
||||
readyResolve();
|
||||
continue;
|
||||
}
|
||||
this.#handleLine(line);
|
||||
if (isRecord(line) && line.type === "rpc_chunk" && !protocolV2Enabled)
|
||||
throw new Error("RPC chunk received before protocol negotiation");
|
||||
const decoded = frameDecoder.push(line);
|
||||
if (decoded) this.#handleLine(decoded);
|
||||
}
|
||||
// A closed stdout is terminal even if the child remains alive. Startup
|
||||
// failures are reaped by the readyPromise catch below; established
|
||||
@@ -359,6 +377,17 @@ export class RpcClient {
|
||||
|
||||
try {
|
||||
await readyPromise;
|
||||
if (protocolV2Supported) {
|
||||
protocolV2Enabled = true;
|
||||
const response = await this.#send({ type: "negotiate_protocol", protocolVersion: 2 });
|
||||
if (
|
||||
!response.success ||
|
||||
response.command !== "negotiate_protocol" ||
|
||||
!isRecord(response.data) ||
|
||||
response.data.protocolVersion !== 2
|
||||
)
|
||||
throw new Error("RPC protocol v2 negotiation failed");
|
||||
}
|
||||
if (this.#customTools.length > 0) {
|
||||
await this.setCustomTools(this.#customTools);
|
||||
}
|
||||
|
||||
@@ -1,8 +1,24 @@
|
||||
import { isDeepStrictEqual } from "node:util";
|
||||
import { isRecord } from "@oh-my-pi/pi-utils";
|
||||
import type { RpcChunkFrame } from "./rpc-types";
|
||||
|
||||
/** Maximum UTF-8 size of one newline-delimited RPC frame, including the newline. */
|
||||
export const MAX_RPC_FRAME_BYTES = 1024 * 1024;
|
||||
/** Maximum UTF-8 size of one logical frame reassembled by protocol v2. */
|
||||
export const MAX_RPC_REASSEMBLED_BYTES = 64 * 1024 * 1024;
|
||||
|
||||
const RPC_CHUNK_PAYLOAD_BYTES = 256 * 1024;
|
||||
|
||||
export type RpcProtocolVersion = 1 | 2;
|
||||
|
||||
interface PendingRpcChunks {
|
||||
chunkId: string;
|
||||
count: number;
|
||||
byteLength: number;
|
||||
nextIndex: number;
|
||||
chunks: Buffer[];
|
||||
receivedBytes: number;
|
||||
}
|
||||
|
||||
interface ShrinkPass {
|
||||
stringCap: number;
|
||||
@@ -68,6 +84,103 @@ function encodedMessageSnapshot(encoded: string): { message: unknown } | undefin
|
||||
: undefined;
|
||||
}
|
||||
|
||||
function encodeChunkedRpcFrame(frame: object, chunkId: string): string {
|
||||
const json = JSON.stringify(frame);
|
||||
const bytes = Buffer.from(json, "utf8");
|
||||
if (bytes.byteLength > MAX_RPC_REASSEMBLED_BYTES) return `${JSON.stringify(overflowFrame(frame))}\n`;
|
||||
const count = Math.ceil(bytes.byteLength / RPC_CHUNK_PAYLOAD_BYTES);
|
||||
let encoded = "";
|
||||
for (let index = 0; index < count; index++) {
|
||||
const chunk: RpcChunkFrame = {
|
||||
type: "rpc_chunk",
|
||||
chunkId,
|
||||
index,
|
||||
count,
|
||||
byteLength: bytes.byteLength,
|
||||
data: bytes
|
||||
.subarray(index * RPC_CHUNK_PAYLOAD_BYTES, (index + 1) * RPC_CHUNK_PAYLOAD_BYTES)
|
||||
.toString("base64"),
|
||||
};
|
||||
const line = `${JSON.stringify(chunk)}\n`;
|
||||
if (serializedFrameBytes(line.slice(0, -1)) > MAX_RPC_FRAME_BYTES)
|
||||
throw new Error("RPC chunk exceeded the transport limit");
|
||||
encoded += line;
|
||||
}
|
||||
return encoded;
|
||||
}
|
||||
|
||||
function isRpcChunkFrame(value: unknown): value is RpcChunkFrame {
|
||||
return isRecord(value) && value.type === "rpc_chunk";
|
||||
}
|
||||
|
||||
function decodeBase64(data: unknown): Buffer {
|
||||
if (
|
||||
typeof data !== "string" ||
|
||||
data.length === 0 ||
|
||||
!/^(?:[A-Za-z0-9+/]{4})*(?:[A-Za-z0-9+/]{2}==|[A-Za-z0-9+/]{3}=)?$/.test(data)
|
||||
)
|
||||
throw new Error("invalid rpc chunk data");
|
||||
const bytes = Buffer.from(data, "base64");
|
||||
if (bytes.toString("base64") !== data) throw new Error("invalid rpc chunk data");
|
||||
return bytes;
|
||||
}
|
||||
|
||||
/** Reassemble protocol v2 chunk frames after each JSONL line has been parsed. */
|
||||
export class RpcFrameDecoder {
|
||||
#pending?: PendingRpcChunks;
|
||||
|
||||
push(value: unknown): object | undefined {
|
||||
if (!isRpcChunkFrame(value)) {
|
||||
if (this.#pending) throw new Error("rpc chunk sequence interrupted");
|
||||
if (!isRecord(value)) throw new Error("rpc frame must be an object");
|
||||
return value;
|
||||
}
|
||||
const { chunkId, index, count, byteLength } = value;
|
||||
if (
|
||||
typeof chunkId !== "string" ||
|
||||
chunkId.length === 0 ||
|
||||
chunkId.length > 128 ||
|
||||
!Number.isSafeInteger(index) ||
|
||||
!Number.isSafeInteger(count) ||
|
||||
!Number.isSafeInteger(byteLength) ||
|
||||
index < 0 ||
|
||||
count < 2 ||
|
||||
count > Math.ceil(MAX_RPC_REASSEMBLED_BYTES / RPC_CHUNK_PAYLOAD_BYTES) ||
|
||||
index >= count ||
|
||||
byteLength <= MAX_RPC_FRAME_BYTES ||
|
||||
byteLength > MAX_RPC_REASSEMBLED_BYTES
|
||||
)
|
||||
throw new Error("invalid rpc chunk metadata");
|
||||
const bytes = decodeBase64(value.data);
|
||||
if (bytes.byteLength > RPC_CHUNK_PAYLOAD_BYTES) throw new Error("rpc chunk payload exceeds the transport limit");
|
||||
|
||||
if (!this.#pending) {
|
||||
if (index !== 0) throw new Error("rpc chunk sequence must start at index 0");
|
||||
this.#pending = { chunkId, count, byteLength, nextIndex: 0, chunks: [], receivedBytes: 0 };
|
||||
}
|
||||
const pending = this.#pending;
|
||||
if (
|
||||
pending.chunkId !== chunkId ||
|
||||
pending.count !== count ||
|
||||
pending.byteLength !== byteLength ||
|
||||
pending.nextIndex !== index
|
||||
)
|
||||
throw new Error("rpc chunk sequence mismatch");
|
||||
pending.chunks.push(bytes);
|
||||
pending.receivedBytes += bytes.byteLength;
|
||||
pending.nextIndex++;
|
||||
if (pending.receivedBytes > pending.byteLength) throw new Error("rpc chunk sequence exceeds declared length");
|
||||
if (pending.nextIndex < pending.count) return undefined;
|
||||
if (pending.receivedBytes !== pending.byteLength) throw new Error("rpc chunk sequence length mismatch");
|
||||
|
||||
this.#pending = undefined;
|
||||
const decoded = new TextDecoder("utf-8", { fatal: true }).decode(Buffer.concat(pending.chunks));
|
||||
const frame: unknown = JSON.parse(decoded);
|
||||
if (!isRecord(frame)) throw new Error("rpc frame must be an object");
|
||||
return frame;
|
||||
}
|
||||
}
|
||||
|
||||
function compactTerminalFrame(
|
||||
frame: object,
|
||||
streamedMessageCount: number,
|
||||
@@ -142,13 +255,27 @@ export function encodeRpcFrame(frame: object, streamedMessageCount = 0, streamed
|
||||
/** Stateful encoder that tracks which messages a client has already received. */
|
||||
export class RpcFrameEncoder {
|
||||
#streamedMessages: unknown[] = [];
|
||||
#protocolVersion: RpcProtocolVersion = 1;
|
||||
#chunkCounter = 0;
|
||||
|
||||
setProtocolVersion(version: number): void {
|
||||
if (version !== 1 && version !== 2) throw new Error(`Unsupported RPC protocol version: ${version}`);
|
||||
this.#protocolVersion = version;
|
||||
}
|
||||
|
||||
encode(frame: object): string {
|
||||
if (isRecord(frame) && frame.type === "agent_start") this.#streamedMessages = [];
|
||||
const encoded = encodeRpcFrame(frame, this.#streamedMessages.length, this.#streamedMessages);
|
||||
const json = JSON.stringify(frame);
|
||||
const encoded =
|
||||
this.#protocolVersion === 2 && serializedFrameBytes(json) > MAX_RPC_FRAME_BYTES
|
||||
? encodeChunkedRpcFrame(frame, `rpc-${++this.#chunkCounter}`)
|
||||
: encodeRpcFrame(frame, this.#streamedMessages.length, this.#streamedMessages);
|
||||
if (!isRecord(frame)) return encoded;
|
||||
if (frame.type === "message_end") {
|
||||
const snapshot = encodedMessageSnapshot(encoded);
|
||||
const snapshot =
|
||||
this.#protocolVersion === 2 && Object.hasOwn(frame, "message")
|
||||
? { message: jsonSnapshot(frame.message) }
|
||||
: encodedMessageSnapshot(encoded);
|
||||
if (snapshot) this.#streamedMessages.push(snapshot.message);
|
||||
} else if (frame.type === "agent_end" && frame.willContinue !== true) this.#streamedMessages = [];
|
||||
return encoded;
|
||||
|
||||
@@ -34,7 +34,7 @@ import type { EventBus } from "../../utils/event-bus";
|
||||
import { initializeExtensions } from "../runtime-init";
|
||||
import { isRpcHostToolResult, isRpcHostToolUpdate, RpcHostToolBridge } from "./host-tools";
|
||||
import { isRpcHostUriResult, RpcHostUriBridge } from "./host-uris";
|
||||
import { RpcFrameEncoder } from "./rpc-frame";
|
||||
import { MAX_RPC_FRAME_BYTES, MAX_RPC_REASSEMBLED_BYTES, RpcFrameEncoder } from "./rpc-frame";
|
||||
import { claimRpcInput } from "./rpc-input";
|
||||
import { RpcSubagentRegistry, readRpcSubagentTranscript } from "./rpc-subagents";
|
||||
import type {
|
||||
@@ -619,9 +619,19 @@ export async function runRpcMode(
|
||||
process.env.PI_NOTIFICATIONS = "off";
|
||||
|
||||
const frameEncoder = new RpcFrameEncoder();
|
||||
process.stdout.write(frameEncoder.encode({ type: "ready" }));
|
||||
process.stdout.write(
|
||||
frameEncoder.encode({
|
||||
type: "ready",
|
||||
protocolVersion: 1,
|
||||
supportedProtocolVersions: [1, 2],
|
||||
maxFrameBytes: MAX_RPC_FRAME_BYTES,
|
||||
maxReassembledFrameBytes: MAX_RPC_REASSEMBLED_BYTES,
|
||||
}),
|
||||
);
|
||||
const output = (obj: RpcResponse | RpcExtensionUIRequest | object) => {
|
||||
process.stdout.write(frameEncoder.encode(obj));
|
||||
if (isRecord(obj) && obj.type === "response" && obj.command === "negotiate_protocol" && obj.success === true)
|
||||
frameEncoder.setProtocolVersion(2);
|
||||
};
|
||||
const emitRpcTitles = shouldEmitRpcTitles();
|
||||
|
||||
@@ -936,6 +946,12 @@ export async function runRpcMode(
|
||||
const id = command.id;
|
||||
|
||||
switch (command.type) {
|
||||
case "negotiate_protocol": {
|
||||
if (command.protocolVersion !== 2)
|
||||
return error(id, "negotiate_protocol", `Unsupported RPC protocol version: ${command.protocolVersion}`);
|
||||
return success(id, "negotiate_protocol", { protocolVersion: 2 });
|
||||
}
|
||||
|
||||
// =================================================================
|
||||
// Prompting
|
||||
// =================================================================
|
||||
|
||||
@@ -25,6 +25,9 @@ import type { TodoPhase } from "../../tools/todo";
|
||||
// ============================================================================
|
||||
|
||||
export type RpcCommand =
|
||||
// Protocol
|
||||
| { id?: string; type: "negotiate_protocol"; protocolVersion: number }
|
||||
|
||||
// Prompting
|
||||
| { id?: string; type: "prompt"; message: string; images?: ImageContent[]; streamingBehavior?: "steer" | "followUp" }
|
||||
| { id?: string; type: "steer"; message: string; images?: ImageContent[] }
|
||||
@@ -132,6 +135,23 @@ export interface RpcPromptResultFrame {
|
||||
agentInvoked: boolean;
|
||||
}
|
||||
|
||||
export interface RpcReadyFrame {
|
||||
type: "ready";
|
||||
protocolVersion: 1;
|
||||
supportedProtocolVersions: [1, 2];
|
||||
maxFrameBytes: number;
|
||||
maxReassembledFrameBytes: number;
|
||||
}
|
||||
|
||||
export interface RpcChunkFrame {
|
||||
type: "rpc_chunk";
|
||||
chunkId: string;
|
||||
index: number;
|
||||
count: number;
|
||||
byteLength: number;
|
||||
data: string;
|
||||
}
|
||||
|
||||
export interface RpcHandoffResult {
|
||||
savedPath?: string;
|
||||
}
|
||||
@@ -168,6 +188,15 @@ export interface RpcSubagentMessagesResult {
|
||||
|
||||
// Success responses with data
|
||||
export type RpcResponse =
|
||||
// Protocol
|
||||
| {
|
||||
id?: string;
|
||||
type: "response";
|
||||
command: "negotiate_protocol";
|
||||
success: true;
|
||||
data: { protocolVersion: 2 };
|
||||
}
|
||||
|
||||
// Prompting (async - events follow)
|
||||
| { id?: string; type: "response"; command: "prompt"; success: true; data?: { agentInvoked: boolean } }
|
||||
| { id?: string; type: "response"; command: "steer"; success: true }
|
||||
|
||||
+51
-6
@@ -14,7 +14,43 @@ if (Bun.env.MOCK_RPC_IGNORE_SIGTERM === "1") {
|
||||
process.on("SIGTERM", () => {});
|
||||
}
|
||||
|
||||
process.stdout.write(`${JSON.stringify({ type: "ready" })}\n`);
|
||||
const supportsProtocolV2 = Bun.env.MOCK_RPC_V2 === "1";
|
||||
let protocolV2Enabled = false;
|
||||
process.stdout.write(
|
||||
`${JSON.stringify(
|
||||
supportsProtocolV2
|
||||
? {
|
||||
type: "ready",
|
||||
protocolVersion: 1,
|
||||
supportedProtocolVersions: [1, 2],
|
||||
maxFrameBytes: 1024 * 1024,
|
||||
maxReassembledFrameBytes: 64 * 1024 * 1024,
|
||||
}
|
||||
: { type: "ready" },
|
||||
)}\n`,
|
||||
);
|
||||
|
||||
function writeFrame(frame: Record<string, unknown>): void {
|
||||
const logical = Buffer.from(JSON.stringify(frame), "utf8");
|
||||
if (!protocolV2Enabled || logical.byteLength <= 1024 * 1024) {
|
||||
process.stdout.write(`${logical.toString("utf8")}\n`);
|
||||
return;
|
||||
}
|
||||
const chunkBytes = 256 * 1024;
|
||||
const count = Math.ceil(logical.byteLength / chunkBytes);
|
||||
for (let index = 0; index < count; index++) {
|
||||
process.stdout.write(
|
||||
`${JSON.stringify({
|
||||
type: "rpc_chunk",
|
||||
chunkId: "mock-rpc-v2",
|
||||
index,
|
||||
count,
|
||||
byteLength: logical.byteLength,
|
||||
data: logical.subarray(index * chunkBytes, (index + 1) * chunkBytes).toString("base64"),
|
||||
})}\n`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Bun's `console` is an AsyncIterable over stdin lines.
|
||||
for await (const raw of console) {
|
||||
@@ -32,15 +68,24 @@ for await (const raw of console) {
|
||||
}
|
||||
if (Bun.env.MOCK_RPC_IGNORE_COMMANDS === "1") continue;
|
||||
const id = typeof frame.id === "string" ? frame.id : undefined;
|
||||
process.stdout.write(
|
||||
`${JSON.stringify({
|
||||
if (frame.type === "negotiate_protocol" && frame.protocolVersion === 2) {
|
||||
writeFrame({
|
||||
id,
|
||||
type: "response",
|
||||
command: frame.type,
|
||||
success: true,
|
||||
data: {},
|
||||
})}\n`,
|
||||
);
|
||||
data: { protocolVersion: 2 },
|
||||
});
|
||||
protocolV2Enabled = true;
|
||||
continue;
|
||||
}
|
||||
writeFrame({
|
||||
id,
|
||||
type: "response",
|
||||
command: frame.type,
|
||||
success: true,
|
||||
data: supportsProtocolV2 ? { payload: "😀".repeat(400_000) } : {},
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
// ignore parse errors — the test harness sends well-formed frames.
|
||||
|
||||
@@ -15,6 +15,17 @@ function isProcessAlive(pid: number): boolean {
|
||||
}
|
||||
|
||||
describe("RpcClient lifecycle (issue #4079 B)", () => {
|
||||
test("auto-negotiates protocol v2 and reassembles an oversized response", async () => {
|
||||
using client = new RpcClient({
|
||||
cliPath: MOCK_AGENT,
|
||||
env: { MOCK_RPC_V2: "1" },
|
||||
});
|
||||
|
||||
await client.start();
|
||||
const state = (await client.getState()) as unknown as { payload: string };
|
||||
expect(state.payload).toBe("😀".repeat(400_000));
|
||||
}, 20_000);
|
||||
|
||||
test("start() succeeds a second time after stop() on the same instance", async () => {
|
||||
using client = new RpcClient({
|
||||
cliPath: MOCK_AGENT,
|
||||
|
||||
@@ -1,5 +1,11 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { encodeRpcFrame, MAX_RPC_FRAME_BYTES, RpcFrameEncoder } from "../src/modes/rpc/rpc-frame";
|
||||
import {
|
||||
encodeRpcFrame,
|
||||
MAX_RPC_FRAME_BYTES,
|
||||
MAX_RPC_REASSEMBLED_BYTES,
|
||||
RpcFrameDecoder,
|
||||
RpcFrameEncoder,
|
||||
} from "../src/modes/rpc/rpc-frame";
|
||||
|
||||
function decode(frame: string): Record<string, unknown> {
|
||||
return JSON.parse(frame) as Record<string, unknown>;
|
||||
@@ -178,4 +184,70 @@ describe("RPC frame encoding", () => {
|
||||
expect(decoded.success).toBe(false);
|
||||
expect(decoded.id).toContain("chars elided for RPC frame");
|
||||
});
|
||||
|
||||
it("losslessly chunks oversized protocol v2 responses into bounded JSONL frames", () => {
|
||||
const frame = {
|
||||
id: "request-v2",
|
||||
type: "response",
|
||||
command: "get_messages",
|
||||
success: true,
|
||||
data: { messages: [{ role: "assistant", content: "😀".repeat(400_000) }] },
|
||||
};
|
||||
const encoder = new RpcFrameEncoder();
|
||||
encoder.setProtocolVersion(2);
|
||||
const encoded = encoder.encode(frame);
|
||||
const lines = encoded.trimEnd().split("\n");
|
||||
const decoder = new RpcFrameDecoder();
|
||||
let decoded: object | undefined;
|
||||
|
||||
expect(lines.length).toBeGreaterThan(1);
|
||||
for (const line of lines) {
|
||||
expect(Buffer.byteLength(`${line}\n`, "utf8")).toBeLessThanOrEqual(MAX_RPC_FRAME_BYTES);
|
||||
decoded = decoder.push(JSON.parse(line));
|
||||
}
|
||||
expect(decoded).toEqual(frame);
|
||||
});
|
||||
|
||||
it("rejects protocol v2 logical frames above the advertised reassembly ceiling", () => {
|
||||
const encoder = new RpcFrameEncoder();
|
||||
encoder.setProtocolVersion(2);
|
||||
const encoded = encoder.encode({
|
||||
id: "request-too-large",
|
||||
type: "response",
|
||||
command: "get_messages",
|
||||
success: true,
|
||||
data: { transcript: "x".repeat(MAX_RPC_REASSEMBLED_BYTES) },
|
||||
});
|
||||
|
||||
expect(decode(encoded)).toEqual({
|
||||
id: "request-too-large",
|
||||
type: "response",
|
||||
command: "get_messages",
|
||||
success: false,
|
||||
error: "RPC response exceeded the transport limit",
|
||||
});
|
||||
});
|
||||
|
||||
it("rejects interrupted protocol v2 chunk sequences", () => {
|
||||
const decoder = new RpcFrameDecoder();
|
||||
decoder.push({
|
||||
type: "rpc_chunk",
|
||||
chunkId: "chunk-1",
|
||||
index: 0,
|
||||
count: 2,
|
||||
byteLength: MAX_RPC_FRAME_BYTES + 1,
|
||||
data: "ew==",
|
||||
});
|
||||
|
||||
expect(() =>
|
||||
decoder.push({
|
||||
type: "rpc_chunk",
|
||||
chunkId: "chunk-2",
|
||||
index: 1,
|
||||
count: 2,
|
||||
byteLength: MAX_RPC_FRAME_BYTES + 1,
|
||||
data: "fQ==",
|
||||
}),
|
||||
).toThrow("rpc chunk sequence mismatch");
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user