Files
oh-my-pi/packages/coding-agent/test/collab/read-only.test.ts
T
roboomp 5b937511c4 fix(collab): chunked welcome so large session snapshots can join
The host used to ship the entire transcript inside a single welcome
frame, so a multi-MB session spent the guest's 30s first-welcome
timeout on the relay transfer itself: ~1.3 MB took ~3s, ~4.2 MB took
~12s, and ~13.6 MB never arrived before the guest gave up with
'timed out waiting for the host's welcome'.

Bump COLLAB_PROTO to 2 and split the welcome:

- welcome carries metadata only (header, state, agents, entryCount,
  readOnly) and lands in well under one second.
- a train of snapshot-chunk frames (SNAPSHOT_CHUNK_BYTES = 512 KB,
  oversize entries ship alone) carries the transcript. Last chunk
  flips final: true; an empty snapshot still emits one final chunk.
- the host queues welcome + chunks synchronously inside #handleHello,
  preserving the host comment's ordering invariant (later broadcast
  frames cannot interleave between them).
- the TUI guest accumulates chunks under a SNAPSHOT_PROGRESS_TIMEOUT_MS
  that resets per chunk; only after final does it write the replica
  jsonl, switchSession, and render. The first-welcome timeout still
  guards arrival of the small welcome.
- the collab-web GuestClient streams entries into the snapshot as
  chunks arrive and flips phase to 'live' on final.

Includes a contract test (in-process relay) asserting the welcome is
metadata-only, the chunk train fans the 1.5 MB synthetic transcript
across multiple frames with only the last marked final, and the
flattened entries match the source snapshot.

Fixes #3144
2026-06-20 18:33:59 +00:00

344 lines
12 KiB
TypeScript

/**
* End-to-end contract: a host started with both link variants marks view-link
* guests read-only in `welcome` and refuses their mutating frames, while
* full-link guests keep prompt/abort/agent-cmd capability. Runs over an
* in-process relay + fake WebSocket transport (no real sockets, no handshake
* or polling latency) that speaks the documented relay forwarding contract,
* with real AES-GCM sealing — only the TUI context and the network transport
* are stubbed. One host/relay boots once and is reused; guest frames ride the
* in-memory transport, so the suite stays fast and time-independent.
*/
import { afterAll, afterEach, beforeAll, describe, expect, it } from "bun:test";
import { importRoomKey } from "@oh-my-pi/pi-coding-agent/collab/crypto";
import { CollabHost } from "@oh-my-pi/pi-coding-agent/collab/host";
import {
COLLAB_PROTO,
type CollabFrame,
parseCollabLink,
rewriteEnvelopePeer,
unpackEnvelope,
} from "@oh-my-pi/pi-coding-agent/collab/protocol";
import { CollabSocket } from "@oh-my-pi/pi-coding-agent/collab/relay-client";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
// ── In-memory transport ────────────────────────────────────────────────────
// FakeWebSocket + InMemoryRelay replace the real Bun.serve relay and loopback
// WebSocket. They mirror the production relay's forwarding contract exactly
// (4-byte peerId envelope routing, peer-joined/peer-left control frames) but
// deliver every frame on a microtask with zero network or timer latency. Real
// CollabSocket / CollabHost run unchanged on top, so sealing, enveloping, the
// hello→welcome handshake, and read-only enforcement are all exercised.
/** Active relay the fake transport routes through; set for the lifetime of this file. */
let activeRelay: InMemoryRelay | null = null;
class FakeWebSocket {
static readonly CONNECTING = 0;
static readonly OPEN = 1;
static readonly CLOSING = 2;
static readonly CLOSED = 3;
binaryType = "blob";
readyState: number = FakeWebSocket.CONNECTING;
readonly role: "host" | "guest";
peerId = 0;
onopen: (() => void) | null = null;
onmessage: ((event: { data: unknown }) => void) | null = null;
onerror: (() => void) | null = null;
onclose: ((event: { code: number; reason: string }) => void) | null = null;
readonly #relay: InMemoryRelay;
constructor(url: string) {
const relay = activeRelay;
if (!relay) throw new Error("FakeWebSocket: no active in-memory relay");
this.#relay = relay;
this.role = new URL(url).searchParams.get("role") === "host" ? "host" : "guest";
queueMicrotask(() => {
if (this.readyState !== FakeWebSocket.CONNECTING) return;
this.readyState = FakeWebSocket.OPEN;
relay.connect(this);
this.onopen?.();
});
}
send(data: Uint8Array): void {
if (this.readyState !== FakeWebSocket.OPEN) return;
// Snapshot: the relay rewrites the peerId in place, and the sender may
// reuse the buffer once send() returns.
const bytes = new Uint8Array(data);
queueMicrotask(() => this.#relay.forward(this, bytes));
}
close(_code?: number): void {
if (this.readyState === FakeWebSocket.CLOSED) return;
this.readyState = FakeWebSocket.CLOSED;
this.#relay.disconnect(this);
queueMicrotask(() => this.onclose?.({ code: 1000, reason: "closed" }));
}
/** Relay → this socket: a binary frame, delivered as ArrayBuffer (binaryType "arraybuffer"). */
deliver(bytes: Uint8Array): void {
if (this.readyState !== FakeWebSocket.OPEN) return;
const copy = new Uint8Array(bytes);
queueMicrotask(() => this.onmessage?.({ data: copy.buffer }));
}
/** Relay → this socket: a JSON control message. */
deliverControl(json: string): void {
if (this.readyState !== FakeWebSocket.OPEN) return;
queueMicrotask(() => this.onmessage?.({ data: json }));
}
}
/** Single-room in-memory relay mirroring the production forwarding contract. */
class InMemoryRelay {
#host: FakeWebSocket | null = null;
readonly #guests = new Map<number, FakeWebSocket>();
#nextPeerId = 1;
connect(ws: FakeWebSocket): void {
if (ws.role === "host") {
this.#host = ws;
return;
}
ws.peerId = this.#nextPeerId++;
this.#guests.set(ws.peerId, ws);
this.#host?.deliverControl(JSON.stringify({ t: "peer-joined", peer: ws.peerId }));
}
forward(from: FakeWebSocket, bytes: Uint8Array): void {
if (from.role === "host") {
const envelope = unpackEnvelope(bytes);
if (!envelope) return;
if (envelope.peerId === 0) {
for (const guest of this.#guests.values()) guest.deliver(bytes);
} else {
this.#guests.get(envelope.peerId)?.deliver(bytes);
}
return;
}
rewriteEnvelopePeer(bytes, from.peerId);
this.#host?.deliver(bytes);
}
disconnect(ws: FakeWebSocket): void {
if (ws.role === "host") {
if (this.#host === ws) this.#host = null;
return;
}
this.#guests.delete(ws.peerId);
this.#host?.deliverControl(JSON.stringify({ t: "peer-left", peer: ws.peerId }));
}
}
interface HostHarness {
ctx: InteractiveModeContext;
prompts: { from?: string }[];
aborts: { count: number };
/** Resolves on the next promptCustomMessage call — no polling. */
nextPrompt(): Promise<{ from?: string }>;
}
/** Minimal InteractiveModeContext double: only the members CollabHost touches. */
function makeHostContext(): HostHarness {
const prompts: { from?: string }[] = [];
const aborts = { count: 0 };
const promptWaiters: ((details: { from?: string }) => void)[] = [];
const ctx = {
settings: { get: () => "" },
sessionManager: {
getSessionId: () => "sess-1",
getCwd: () => "/tmp",
snapshotForReplication: () => ({
header: { type: "session", id: "sess-1", timestamp: new Date().toISOString(), cwd: "/tmp" },
entries: [],
}),
onEntryAppended: undefined,
},
session: {
isStreaming: false,
queuedMessageCount: 0,
sessionName: "test",
model: undefined,
thinkingLevel: undefined,
subscribe: () => () => {},
emitNotice: () => {},
promptCustomMessage: (message: { details?: { from?: string } }) => {
const details = message.details ?? {};
prompts.push(details);
for (const waiter of promptWaiters.splice(0)) waiter(details);
return Promise.resolve();
},
abort: () => {
aborts.count++;
return Promise.resolve();
},
},
eventBus: undefined,
statusLine: {
setCollabStatus: () => {},
invalidate: () => {},
getCachedContextBreakdown: () => ({ usedTokens: 0, contextWindow: 0 }),
},
ui: { requestRender: () => {} },
showStatus: () => {},
collabHost: undefined,
} as unknown as InteractiveModeContext;
const nextPrompt = (): Promise<{ from?: string }> => {
const { promise, resolve } = Promise.withResolvers<{ from?: string }>();
promptWaiters.push(resolve);
return promise;
};
return { ctx, prompts, aborts, nextPrompt };
}
interface TestGuest {
socket: CollabSocket;
nextFrame(): Promise<CollabFrame>;
}
/**
* Frames the test harness skips: the host's debounced broadcasts (state,
* agents, entry, event, bus) and the per-peer snapshot-chunk train that
* follows every welcome. They interleave nondeterministically with the
* directed welcome/error frames these tests actually assert on.
*/
const FILTERED_FRAME_TYPES: Record<string, true> = {
state: true,
agents: true,
entry: true,
event: true,
bus: true,
"snapshot-chunk": true,
};
/**
* Raw guest speaking the wire protocol directly. `writeToken` overrides the link's token (e.g. forged).
* Broadcast frames interleave nondeterministically with directed replies (the post-hello state
* broadcast races the first prompt's error reply), so `nextFrame` drops them and yields only the
* welcome/error frames these tests assert on.
*/
async function joinAsGuest(link: string, name: string, writeTokenOverride?: string): Promise<TestGuest> {
const parsed = parseCollabLink(link);
if ("error" in parsed) throw new Error(parsed.error);
const writeToken =
writeTokenOverride ?? (parsed.writeToken ? Buffer.from(parsed.writeToken).toString("base64url") : undefined);
const key = await importRoomKey(parsed.key);
const socket = new CollabSocket({ wsUrl: parsed.wsUrl, role: "guest", key });
const queue: CollabFrame[] = [];
const waiters: ((frame: CollabFrame) => void)[] = [];
socket.onFrame = frame => {
if (FILTERED_FRAME_TYPES[frame.t]) return;
const waiter = waiters.shift();
if (waiter) waiter(frame);
else queue.push(frame);
};
socket.onOpen = () => socket.send({ t: "hello", proto: COLLAB_PROTO, name, writeToken });
socket.connect();
const nextFrame = (): Promise<CollabFrame> => {
const queued = queue.shift();
if (queued) return Promise.resolve(queued);
const { promise, resolve } = Promise.withResolvers<CollabFrame>();
waiters.push(resolve);
return promise;
};
return { socket, nextFrame };
}
// ── Shared host/relay, booted once ──────────────────────────────────────────
// Booting the relay + host and connecting the host socket is the only heavy
// step; it is identical across all three tests (none mutate host config), so it
// runs once. Per-test guest state is reset in afterEach.
const RealWebSocket = globalThis.WebSocket;
const guestCleanups: (() => void)[] = [];
let harness: HostHarness;
let host: CollabHost;
beforeAll(async () => {
globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket;
activeRelay = new InMemoryRelay();
harness = makeHostContext();
host = new CollabHost(harness.ctx);
// Port is irrelevant: the fake transport routes by the `role` query param.
await host.start("ws://localhost:8787");
});
afterEach(() => {
for (const cleanup of guestCleanups.splice(0).reverse()) cleanup();
harness.prompts.length = 0;
harness.aborts.count = 0;
});
afterAll(async () => {
// Restore the real transport first so the global is clean even if stop() throws;
// the host's socket holds its own FakeWebSocket/relay refs, so teardown still works.
globalThis.WebSocket = RealWebSocket;
activeRelay = null;
await host.stop("test done");
});
describe("collab read-only links", () => {
it("welcomes view-link guests read-only and refuses their mutating frames", async () => {
const { prompts, aborts } = harness;
expect(host.viewLink).not.toBe(host.link);
const guest = await joinAsGuest(host.viewLink, "viewer");
guestCleanups.push(() => guest.socket.close());
const welcome = await guest.nextFrame();
if (welcome.t !== "welcome") throw new Error(`expected welcome, got ${welcome.t}`);
expect(welcome.readOnly).toBe(true);
guest.socket.send({ t: "prompt", text: "do something" });
const promptReply = await guest.nextFrame();
if (promptReply.t !== "error") throw new Error(`expected error, got ${promptReply.t}`);
expect(promptReply.message).toContain("read-only");
expect(prompts).toHaveLength(0);
guest.socket.send({ t: "abort" });
const abortReply = await guest.nextFrame();
expect(abortReply.t).toBe("error");
expect(aborts.count).toBe(0);
guest.socket.send({ t: "agent-cmd", cmd: "kill", agentId: "nope" });
const cmdReply = await guest.nextFrame();
expect(cmdReply.t).toBe("error");
expect(host.participants.find(p => p.name === "viewer")?.readOnly).toBe(true);
});
it("keeps full write capability for guests holding the write token", async () => {
const { prompts, nextPrompt } = harness;
const guest = await joinAsGuest(host.link, "writer");
guestCleanups.push(() => guest.socket.close());
const welcome = await guest.nextFrame();
if (welcome.t !== "welcome") throw new Error(`expected welcome, got ${welcome.t}`);
expect(welcome.readOnly).toBeUndefined();
const prompted = nextPrompt();
guest.socket.send({ t: "prompt", text: "real prompt" });
expect(await prompted).toEqual({ from: "writer" });
expect(prompts).toHaveLength(1);
expect(host.participants.find(p => p.name === "writer")?.readOnly).toBeUndefined();
});
it("treats a forged write token as read-only", async () => {
const { prompts } = harness;
// A viewer knows the room key but not the token; garbage must not escalate.
const forged = Buffer.alloc(16, 0xab).toString("base64url");
const guest = await joinAsGuest(host.viewLink, "forger", forged);
guestCleanups.push(() => guest.socket.close());
const welcome = await guest.nextFrame();
if (welcome.t !== "welcome") throw new Error(`expected welcome, got ${welcome.t}`);
expect(welcome.readOnly).toBe(true);
guest.socket.send({ t: "prompt", text: "escalation attempt" });
const reply = await guest.nextFrame();
expect(reply.t).toBe("error");
expect(prompts).toHaveLength(0);
});
});