Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions packages/client/src/broadcastTransport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,19 @@ describe('broadcast transport ↔ leader relay', () => {
expect(engine.commits).toEqual([handle]);
});

it('wipes an upload chunk it never gets onto a port', async () => {
const bus = new FakeBus();
const ports = new FakeCourierNetwork();
relayOn(bus, new FakeEngineTransport(), ports.courier('leader'));
// No port to move the plaintext out over, so this tab stays its owner.
const follower = followerOn(bus, 'follower-1', unavailableCourier);
const plaintext = Uint8Array.of(4, 3, 2, 1);

await expect(follower.pushChunk(1n, plaintext.buffer as ArrayBuffer)).rejects.toThrow();

expect([...plaintext]).toEqual([0, 0, 0, 0]);
});

it('keeps upload plaintext and command arguments off the channel', async () => {
const bus = new FakeBus();
const ports = new FakeCourierNetwork();
Expand Down
10 changes: 7 additions & 3 deletions packages/client/src/broadcastTransport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -211,9 +211,13 @@ export class BroadcastTransport extends CorrelatedTransport {
build: (requestId: number) => PortRequest,
transfer?: Transferable[]
): Promise<T> {
return this.request<T, MessagePortLike>(this.ensurePort(), (requestId, port) => {
port.postMessage(build(requestId), transfer);
});
return this.request<T, MessagePortLike>(
this.ensurePort(),
(requestId, port) => {
port.postMessage(build(requestId), transfer);
},
transfer
);
}

private ensurePort(): Promise<MessagePortLike> {
Expand Down
162 changes: 162 additions & 0 deletions packages/client/src/correlatedTransport.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,172 @@
import { describe, expect, it } from 'vitest';

import {
CorrelatedTransport,
EngineRequestError,
engineErrorCode,
isRecoverableEngineError,
} from './correlatedTransport.js';
import type { SnapshotDescriptor, WriteHandle } from './worker/protocol.js';

function unsupported(): never {
throw new Error('outside this probe');
}

/**
* A concrete transport wiring only the request skeleton: `pushChunk` carries a
* transfer, `open`/`breakDown` drive the gate and the terminal latch, and the
* rest of the engine surface is out of scope here.
*/
class ProbeTransport extends CorrelatedTransport {
private resolveGate!: () => void;
private rejectGate!: (error: Error) => void;
private readonly gate = new Promise<void>((resolve, reject) => {
this.resolveGate = resolve;
this.rejectGate = reject;
});

constructor(private readonly onSend: (id: number) => void = () => undefined) {
super();
this.gate.catch(() => undefined);
}

pushChunk(_handle: WriteHandle, chunk: ArrayBuffer): Promise<void> {
return this.dispatch(this.gate, (id) => this.onSend(id), [chunk]);
}

open(): void {
this.resolveGate();
}

shut(error: Error): void {
this.rejectGate(error);
}

breakDown(error: Error): void {
this.fail(error);
}

answer(id: number): void {
this.settle(id, true);
}

start(): Promise<void> {
return unsupported();
}
command(): Promise<void> {
return unsupported();
}
beginWrite(): Promise<WriteHandle> {
return unsupported();
}
commitWrite(): Promise<bigint> {
return unsupported();
}
abortWrite(): Promise<void> {
return unsupported();
}
snapshot(): Promise<SnapshotDescriptor> {
return unsupported();
}
siweChallenge(): Promise<string> {
return unsupported();
}
download(): Promise<ArrayBuffer> {
return unsupported();
}
openContentStream(): Promise<WriteHandle> {
return unsupported();
}
readStream(): Promise<ArrayBuffer> {
return unsupported();
}
closeStream(): Promise<void> {
return unsupported();
}
close(): void {
unsupported();
}
}

const plaintext = (): Uint8Array => Uint8Array.of(1, 2, 3, 4);

describe('CorrelatedTransport chunk ownership', () => {
it('leaves the chunk alone once the send has taken it', async () => {
const sent: number[] = [];
const probe = new ProbeTransport((id) => sent.push(id));
probe.open();
const chunk = plaintext();

const pushed = probe.pushChunk(1n, chunk.buffer as ArrayBuffer);
await Promise.resolve();
probe.answer(sent[0]);

// Wiping a sent chunk would zero the bytes the receiver is about to seal.
await expect(pushed).resolves.toBeUndefined();
expect(chunk).toEqual(plaintext());
});

it('wipes the chunk of a request refused by an already-terminal transport', async () => {
const probe = new ProbeTransport();
probe.open();
probe.breakDown(new Error('engine transport closed'));
const chunk = plaintext();

await expect(probe.pushChunk(1n, chunk.buffer as ArrayBuffer)).rejects.toThrow('closed');
expect(chunk).toEqual(new Uint8Array(4));
});

it('wipes the chunk of a request the readiness gate refuses', async () => {
const probe = new ProbeTransport();
const chunk = plaintext();

const pushed = probe.pushChunk(1n, chunk.buffer as ArrayBuffer);
probe.shut(new Error('leader changed; retry'));

await expect(pushed).rejects.toThrow('leader changed; retry');
expect(chunk).toEqual(new Uint8Array(4));
});

it('wipes the chunk of a request the transport outlives its gate to refuse', async () => {
const probe = new ProbeTransport();
const chunk = plaintext();

const pushed = probe.pushChunk(1n, chunk.buffer as ArrayBuffer);
probe.breakDown(new Error('engine transport closed'));
probe.open();

await expect(pushed).rejects.toThrow('closed');
expect(chunk).toEqual(new Uint8Array(4));
});

it('wipes the chunk of a send that throws', async () => {
const probe = new ProbeTransport(() => {
throw new Error('port is dead');
});
probe.open();
const chunk = plaintext();

await expect(probe.pushChunk(1n, chunk.buffer as ArrayBuffer)).rejects.toThrow('port is dead');
expect(chunk).toEqual(new Uint8Array(4));
});

it('wipes a chunk minted in another realm, which instanceof does not answer for', async () => {
// A buffer from a worker or a frame is an ArrayBuffer that `instanceof`
// calls false, and a secret that arrived from there needs the same scrub.
const { runInNewContext } = await import('node:vm');
const foreign = runInNewContext(
'const b = new ArrayBuffer(4); new Uint8Array(b).set([9, 9, 9, 9]); ({ b, v: new Uint8Array(b) })'
) as { b: ArrayBuffer; v: Uint8Array };
expect(foreign.b instanceof ArrayBuffer).toBe(false);

const probe = new ProbeTransport();
probe.open();
probe.breakDown(new Error('engine transport closed'));

await expect(probe.pushChunk(1n, foreign.b)).rejects.toThrow('closed');
expect([...foreign.v]).toEqual([0, 0, 0, 0]);
});
});

describe('engineErrorCode', () => {
it('reads the code off an engine failure and nothing else', () => {
Expand Down
57 changes: 52 additions & 5 deletions packages/client/src/correlatedTransport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,36 @@ export function unknownHandle(kind: HandleKind): EngineRequestError {
);
}

/**
* An `ArrayBuffer`'s length, or `null` for anything that is not one. Branding
* by the `byteLength` getter rather than `instanceof`, which answers false for
* a buffer minted in another realm — a worker's, a frame's — leaving a secret
* that reached this list from there unscrubbed.
*/
const byteLengthOf = Object.getOwnPropertyDescriptor(ArrayBuffer.prototype, 'byteLength')?.get;

function bufferLength(item: Transferable): number | null {
try {
return (byteLengthOf?.call(item) as number | undefined) ?? null;
} catch {
return null;
}
}

/**
* Scrubs the buffers a send would have transferred. A request that rejects
* before its send leaves this frame their terminal owner — nothing detaches
* them and no callee can reach them — so the plaintext is cleared here
* (AGENTS.md 7). A transferred buffer reads as empty, so a send that did run
* leaves this a no-op.
*/
function wipeTransfer(transfer: Transferable[] | undefined): void {
for (const item of transfer ?? []) {
const length = bufferLength(item);
if (length !== null && length > 0) new Uint8Array(item as ArrayBuffer).fill(0);
}
Comment thread
FSM1 marked this conversation as resolved.
}

/**
* Delivers one event to every listener, isolating a throwing subscriber so it
* cannot drop the event for the rest.
Expand Down Expand Up @@ -125,16 +155,24 @@ export abstract class CorrelatedTransport implements EngineTransport {
* synchronous `send` failure deletes the pending entry before rejecting so it
* is never stranded. Resolves with the response's result value (`undefined`
* for a plain ack).
*
* `transfer` is what the send would have moved out of this realm; every route
* to a rejection without it scrubs them ([`wipeTransfer`]).
*/
protected request<T, G = void>(
readyGate: Promise<G>,
send: (id: number, gate: G) => void
send: (id: number, gate: G) => void,
transfer?: Transferable[]
): Promise<T> {
if (this.terminalError) return Promise.reject(this.terminalError);
if (this.terminalError) {
wipeTransfer(transfer);
return Promise.reject(this.terminalError);
}
return readyGate.then(
(gate) =>
new Promise<T>((resolve, reject) => {
if (this.terminalError) {
wipeTransfer(transfer);
reject(this.terminalError);
return;
}
Expand All @@ -144,15 +182,24 @@ export abstract class CorrelatedTransport implements EngineTransport {
send(id, gate);
} catch (error) {
this.pending.delete(id);
wipeTransfer(transfer);
reject(error instanceof Error ? error : new Error(String(error)));
}
})
}),
(error: unknown) => {
wipeTransfer(transfer);
throw error;
}
);
}

/** The void-ack variant of [`request`](CorrelatedTransport.request). */
protected dispatch(readyGate: Promise<void>, send: (id: number) => void): Promise<void> {
return this.request<void>(readyGate, send);
protected dispatch(
readyGate: Promise<void>,
send: (id: number) => void,
transfer?: Transferable[]
): Promise<void> {
return this.request<void>(readyGate, send, transfer);
}

/** Correlates a response to its request id, resolving or rejecting it. */
Expand Down
15 changes: 14 additions & 1 deletion packages/client/src/engineClient.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,17 @@ describe('EngineClient leadership + transport swap', () => {
await leader.dispose();
});

it('scrubs the login secret a closed client refuses', async () => {
const { tab } = origin();
const client = tab();
await client.dispose();
const secret = Uint8Array.of(1, 2, 3, 4);

await expect(client.start(secret.buffer as ArrayBuffer)).rejects.toThrow('closed');

expect(secret).toEqual(new Uint8Array(4));
});

it('refuses a write handle minted by a leadership that has been replaced', async () => {
const { tab, workers } = origin();
const secretSource = {
Expand Down Expand Up @@ -214,9 +225,11 @@ describe('EngineClient leadership + transport swap', () => {
expect.objectContaining({ type: 'pushChunk', handle: 1n })
);

await expect(follower.pushChunk(stale, Uint8Array.of(9).buffer)).rejects.toMatchObject({
const refused = Uint8Array.of(9, 9, 9, 9);
await expect(follower.pushChunk(stale, refused.buffer as ArrayBuffer)).rejects.toMatchObject({
code: 'unknownWriteHandle',
});
expect(refused).toEqual(new Uint8Array(4));
await expect(follower.commitWrite(stale)).rejects.toMatchObject({
code: 'unknownWriteHandle',
});
Expand Down
13 changes: 9 additions & 4 deletions packages/client/src/engineClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -161,13 +161,14 @@ export class EngineClient implements EngineTransport {
// --- EngineTransport ---

start(secret: ArrayBuffer): Promise<void> {
if (this.role === 'closed') return Promise.reject(new Error('engine client closed'));
// This seam is the secret's terminal owner (security rule 7). On the leader
// path the worker becomes the terminal owner — `LocalTransport.start`
// transfers the buffer in (neutered), never copied. On the follower path the
// keyless transport gets no secret: we scrub the buffer we decided not to use
// right here, rather than in a callee that would be zeroing someone else's.
// keyless transport gets no secret, and a closed client no transport at all:
// we scrub the buffer we decided not to use right here, rather than in a
// callee that would be zeroing someone else's.
if (this.role !== 'leader') new Uint8Array(secret).fill(0);
if (this.role === 'closed') return Promise.reject(new Error('engine client closed'));
return this.current.start(secret).then(() => {
this.started = true;
});
Expand All @@ -187,7 +188,11 @@ export class EngineClient implements EngineTransport {

pushChunk(handle: WriteHandle, chunk: ArrayBuffer): Promise<void> {
const inner = this.writes.resolve(handle);
if (inner === undefined) return Promise.reject(unknownHandle('write'));
if (inner === undefined) {
// Refused before any transfer: this seam is the chunk's terminal owner.
new Uint8Array(chunk).fill(0);
return Promise.reject(unknownHandle('write'));
}
return this.current.pushChunk(inner, chunk);
}

Expand Down
13 changes: 13 additions & 0 deletions packages/client/src/transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -237,6 +237,19 @@ describe('LocalTransport', () => {
expect(posted.transfer).toEqual([secret]);
});

it.each([
['the secret', (t: LocalTransport, buffer: ArrayBuffer) => t.start(buffer)],
['an upload chunk', (t: LocalTransport, buffer: ArrayBuffer) => t.pushChunk(1n, buffer)],
])('wipes %s a torn-down worker never took', async (_case, send) => {
const transport = new LocalTransport(new FakeWorker());
transport.close();
const plaintext = Uint8Array.of(1, 2, 3, 4);

await expect(send(transport, plaintext.buffer as ArrayBuffer)).rejects.toThrow('closed');

expect(plaintext).toEqual(new Uint8Array(4));
});

it('transfers the chunk buffer on pushChunk and resolves the write handle and op id', async () => {
const worker = new FakeWorker();
const transport = new LocalTransport(worker);
Expand Down
Loading