fix(collab): clip oversized arrays alongside strings in shrinkForReplication
Per #3740 review: many short strings (e.g. a tool result whose content array holds thousands of small text blocks) could sum past MAX_REPLICATED_PAYLOAD_BYTES without any individual field crossing the per-string floor, so the helper exited the truncation loop and shipped an oversized frame — the relay close/reconnect loop the helper was meant to prevent. Replace the string-only truncation pass with a single walker that head-truncates strings AND head-clips arrays in one descent, driven by a concrete SHRINK_PASSES schedule that tightens both axes together. The final pass clamps every string to 64 B and every array to one element, so any payload converges. Add direct unit tests for shrinkForReplication covering: identity for small values, single-giant-string clamp, many-short-strings array clamp (no field above floor), and discriminator preservation on a fully-shrunk payload.
This commit is contained in:
@@ -4,19 +4,22 @@
|
||||
* The host wraps every {@link CollabFrame} in an AES-GCM envelope and ships it
|
||||
* through the relay's WebSocket. WebSocket servers enforce a per-frame
|
||||
* `maxPayloadLength` (Bun's default is 16 MB; many proxies cap lower). A
|
||||
* single oversized SessionEntry — typically a `read`/`bash`/`search` tool
|
||||
* result that captured a multi-megabyte blob — would otherwise ship as its own
|
||||
* oversized chunk and trip that limit, killing the host's WebSocket with
|
||||
* `1006 Received too big message`. `CollabSocket` treats 1006 as transient and
|
||||
* reconnects, the next guest hello triggers the same oversized send, and the
|
||||
* loop never breaks (issue #3739).
|
||||
* single oversized payload — typically a `read`/`bash`/`search` tool result
|
||||
* captured as one multi-megabyte string, or a tool result whose `content`
|
||||
* array holds thousands of small blocks — would otherwise ship as its own
|
||||
* oversized frame and trip that limit, killing the host's WebSocket with
|
||||
* `1006 Received too big message`. `CollabSocket` treats 1006 as transient
|
||||
* and reconnects, the next guest hello triggers the same oversized send, and
|
||||
* the loop never breaks (issue #3739).
|
||||
*
|
||||
* This helper bounds any JSON-serializable payload below
|
||||
* {@link MAX_REPLICATED_PAYLOAD_BYTES}. Already-small payloads pass through
|
||||
* untouched; oversized ones are returned as a deep-cloned shadow where long
|
||||
* strings are head-truncated with an `[…N bytes elided for collab session]`
|
||||
* marker. Guests still see the structural mirror; tool outputs degrade
|
||||
* gracefully instead of looping the whole session.
|
||||
* strings are head-truncated AND long arrays are head-clipped, with
|
||||
* `[…N chars elided for collab session]` / `[…N items elided for collab
|
||||
* session]` markers. Both axes are needed: string truncation alone leaves
|
||||
* the cap unenforced for a payload built of many short strings, where no
|
||||
* field exceeds the per-string floor.
|
||||
*/
|
||||
|
||||
/**
|
||||
@@ -28,36 +31,62 @@
|
||||
export const MAX_REPLICATED_PAYLOAD_BYTES = 1 * 1024 * 1024;
|
||||
|
||||
/**
|
||||
* Starting per-string head-truncation cap. The shrinker halves this until the
|
||||
* shrunk payload fits {@link MAX_REPLICATED_PAYLOAD_BYTES}; the floor
|
||||
* (`MIN_STRING_CAP_BYTES`) bounds the worst case so a payload with very many
|
||||
* long strings still converges in a handful of passes.
|
||||
* Progressive shrink passes. Each pass tightens both the per-string cap and
|
||||
* the per-array head limit; the loop stops at the first pass whose output
|
||||
* fits {@link MAX_REPLICATED_PAYLOAD_BYTES}. The schedule is concrete (not
|
||||
* recomputed) so the failure modes the helper guards against are visible:
|
||||
*
|
||||
* - One giant string → the first pass already truncates it under 64 KB.
|
||||
* - Array of many small blocks (e.g. a tool result with thousands of
|
||||
* `{type:"text", text:"..."}` content items) → later passes head-clip the
|
||||
* array to a small sample with a `[…N items elided]` summary element.
|
||||
*
|
||||
* The final pass clamps every string to 64 B and every array to one element,
|
||||
* so even pathological mixes converge.
|
||||
*/
|
||||
const INITIAL_STRING_CAP_BYTES = 64 * 1024;
|
||||
interface ShrinkPass {
|
||||
stringCap: number;
|
||||
arrayLimit: number;
|
||||
}
|
||||
|
||||
const MIN_STRING_CAP_BYTES = 256;
|
||||
const SHRINK_PASSES: readonly ShrinkPass[] = [
|
||||
{ stringCap: 64 * 1024, arrayLimit: 256 },
|
||||
{ stringCap: 16 * 1024, arrayLimit: 128 },
|
||||
{ stringCap: 4 * 1024, arrayLimit: 64 },
|
||||
{ stringCap: 1 * 1024, arrayLimit: 32 },
|
||||
{ stringCap: 256, arrayLimit: 16 },
|
||||
{ stringCap: 256, arrayLimit: 4 },
|
||||
{ stringCap: 64, arrayLimit: 1 },
|
||||
];
|
||||
|
||||
const STRING_ELISION_RESERVE = 80;
|
||||
|
||||
/**
|
||||
* Recursively walk `value`, head-truncating any string longer than `cap`.
|
||||
* Returns a deep-cloned copy when truncation occurs and the original
|
||||
* reference for already-small subtrees, so unchanged subgraphs avoid the
|
||||
* structural clone cost.
|
||||
* Recursively walk `value`, head-truncating any string longer than
|
||||
* `stringCap` and head-clipping any array longer than `arrayLimit`. Returns
|
||||
* a freshly built deep clone — every object/array is rebuilt so the
|
||||
* recursive output can be safely serialized in isolation. Cheap to call on
|
||||
* small values: short strings, numbers, and booleans pass through without
|
||||
* allocation.
|
||||
*/
|
||||
function truncateStrings(value: unknown, cap: number): unknown {
|
||||
function shrinkWalk(value: unknown, stringCap: number, arrayLimit: number): unknown {
|
||||
if (typeof value === "string") {
|
||||
if (value.length <= cap) return value;
|
||||
const headLen = Math.max(0, cap - 80);
|
||||
if (value.length <= stringCap) return value;
|
||||
const headLen = Math.max(0, stringCap - STRING_ELISION_RESERVE);
|
||||
return `${value.slice(0, headLen)}\n…[${value.length - headLen} chars elided for collab session]`;
|
||||
}
|
||||
if (Array.isArray(value)) {
|
||||
const out: unknown[] = new Array(value.length);
|
||||
for (let i = 0; i < value.length; i++) out[i] = truncateStrings(value[i], cap);
|
||||
const keep = Math.min(value.length, arrayLimit);
|
||||
const elided = value.length - keep;
|
||||
const out: unknown[] = new Array(elided > 0 ? keep + 1 : keep);
|
||||
for (let i = 0; i < keep; i++) out[i] = shrinkWalk(value[i], stringCap, arrayLimit);
|
||||
if (elided > 0) out[keep] = `…[${elided} items elided for collab session]`;
|
||||
return out;
|
||||
}
|
||||
if (value && typeof value === "object") {
|
||||
const src = value as Record<string, unknown>;
|
||||
const out: Record<string, unknown> = {};
|
||||
for (const k in src) out[k] = truncateStrings(src[k], cap);
|
||||
for (const k in src) out[k] = shrinkWalk(src[k], stringCap, arrayLimit);
|
||||
return out;
|
||||
}
|
||||
return value;
|
||||
@@ -65,19 +94,18 @@ function truncateStrings(value: unknown, cap: number): unknown {
|
||||
|
||||
/**
|
||||
* Return `value` unchanged when its JSON serialization already fits
|
||||
* {@link MAX_REPLICATED_PAYLOAD_BYTES}; otherwise return a deep-cloned shadow
|
||||
* with long strings head-truncated until the payload fits. The cap is halved
|
||||
* across passes so a value with one giant string and one with many medium
|
||||
* strings both converge.
|
||||
* {@link MAX_REPLICATED_PAYLOAD_BYTES}; otherwise return a deep-cloned
|
||||
* shadow shrunk along both string and array axes until the payload fits.
|
||||
* The function is generic over `T` because the wire shape is preserved:
|
||||
* only string leaves and array tails change; discriminator fields, ids, and
|
||||
* other small metadata pass through untouched.
|
||||
*/
|
||||
export function shrinkForReplication<T>(value: T): T {
|
||||
if (JSON.stringify(value).length <= MAX_REPLICATED_PAYLOAD_BYTES) return value;
|
||||
let cap = INITIAL_STRING_CAP_BYTES;
|
||||
let shrunk = value as unknown;
|
||||
while (cap >= MIN_STRING_CAP_BYTES) {
|
||||
shrunk = truncateStrings(value, cap);
|
||||
let shrunk: unknown = value;
|
||||
for (const pass of SHRINK_PASSES) {
|
||||
shrunk = shrinkWalk(value, pass.stringCap, pass.arrayLimit);
|
||||
if (JSON.stringify(shrunk).length <= MAX_REPLICATED_PAYLOAD_BYTES) return shrunk as T;
|
||||
cap = Math.floor(cap / 2);
|
||||
}
|
||||
return shrunk as T;
|
||||
}
|
||||
|
||||
@@ -30,6 +30,10 @@ import {
|
||||
unpackEnvelope,
|
||||
} from "@oh-my-pi/pi-coding-agent/collab/protocol";
|
||||
import { CollabSocket } from "@oh-my-pi/pi-coding-agent/collab/relay-client";
|
||||
import {
|
||||
MAX_REPLICATED_PAYLOAD_BYTES,
|
||||
shrinkForReplication,
|
||||
} from "@oh-my-pi/pi-coding-agent/collab/replication-shrink";
|
||||
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
|
||||
import type { SessionEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
|
||||
|
||||
@@ -271,3 +275,60 @@ describe("collab replication shrinking (#3739)", () => {
|
||||
expect(shrunkContent).toContain("chars elided for collab session");
|
||||
});
|
||||
});
|
||||
|
||||
describe("shrinkForReplication (#3740 review)", () => {
|
||||
it("passes already-small values through by reference", () => {
|
||||
const small = { type: "message", id: "x", text: "hi" };
|
||||
expect(shrinkForReplication(small)).toBe(small);
|
||||
});
|
||||
|
||||
it("clamps a single giant string under the cap with an elision marker", () => {
|
||||
const giant = { content: "x".repeat(5 * 1024 * 1024) };
|
||||
const shrunk = shrinkForReplication(giant);
|
||||
const size = JSON.stringify(shrunk).length;
|
||||
expect(size).toBeLessThanOrEqual(MAX_REPLICATED_PAYLOAD_BYTES);
|
||||
expect(shrunk).not.toBe(giant);
|
||||
expect(shrunk.content).toContain("chars elided for collab session");
|
||||
});
|
||||
|
||||
it("clamps a payload built of many short strings (no individual oversized) under the cap", () => {
|
||||
// Realistic shape: a tool result content array with thousands of small
|
||||
// text blocks. ~3 MB total; no individual string crosses the 64 B
|
||||
// floor of the final shrink pass, so the helper MUST clip the array,
|
||||
// not just the strings.
|
||||
const content = Array.from({ length: 100_000 }, (_, i) => ({
|
||||
type: "text",
|
||||
text: `block-${i}`,
|
||||
}));
|
||||
const payload = { role: "toolResult", content };
|
||||
const original = JSON.stringify(payload).length;
|
||||
expect(original).toBeGreaterThan(MAX_REPLICATED_PAYLOAD_BYTES);
|
||||
const shrunk = shrinkForReplication(payload);
|
||||
const shrunkSize = JSON.stringify(shrunk).length;
|
||||
expect(shrunkSize).toBeLessThanOrEqual(MAX_REPLICATED_PAYLOAD_BYTES);
|
||||
// Array was clipped with a summary marker reporting the dropped count.
|
||||
const shrunkContent = shrunk.content;
|
||||
expect(Array.isArray(shrunkContent)).toBe(true);
|
||||
const marker = shrunkContent.find(
|
||||
(item: unknown) => typeof item === "string" && item.includes("items elided for collab session"),
|
||||
);
|
||||
expect(marker).toBeDefined();
|
||||
});
|
||||
|
||||
it("preserves the wire discriminator on a fully-shrunk payload", () => {
|
||||
// Even the worst-case final pass keeps the discriminator key/value
|
||||
// pairs intact (only string leaves and array tails are touched), so
|
||||
// guests still see the entry/event as the correct kind.
|
||||
const evt = {
|
||||
type: "tool_execution_end",
|
||||
toolCallId: "call-1",
|
||||
toolName: "read",
|
||||
result: { text: "x".repeat(8 * 1024 * 1024) },
|
||||
};
|
||||
const shrunk = shrinkForReplication(evt);
|
||||
expect(shrunk.type).toBe("tool_execution_end");
|
||||
expect(shrunk.toolCallId).toBe("call-1");
|
||||
expect(shrunk.toolName).toBe("read");
|
||||
expect(JSON.stringify(shrunk).length).toBeLessThanOrEqual(MAX_REPLICATED_PAYLOAD_BYTES);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user