From ca3da33615370fb66c240a354df40ae39a399822 Mon Sep 17 00:00:00 2001 From: Jonas Templestein <242550+jonastemplestein@users.noreply.github.com> Date: Mon, 17 Aug 2026 15:20:22 +0100 Subject: [PATCH] Add deferred WebSocket upgrade materialization On Cloudflare Workers, a tunneled upgrade Response's socket now arrives by default as an opaque { readable, writable, init } byte-stream pair (DeferredWebSocketUpgrade) instead of an eagerly materialized WebSocketPair end: at a session endpoint with internal RPC hops beyond it, eager materialization was only ever correct when that endpoint itself served the upgrade, and materialization is a one-way door (once the pump attaches, the tunnel streams are consumed), so the receiver must build the forwardable form up front. The new materializeUpgrade() export rebuilds the real upgrade Response from the pair at the hop that actually serves the 101; pair.init carries the provider's upgrade headers (e.g. a negotiated Sec-WebSocket-Protocol) so they survive to that response. Other runtimes keep delivering a usable TunneledWebSocket by default (no internal hops to cross), and the new deferUpgradeMaterialization session option overrides the default in either direction -- false on Workers restores eager materialization for an endpoint that serves the socket itself. This exists to thread a tunneled socket through boundaries that can serialize byte streams but not sockets -- specifically native Cloudflare Workers RPC between isolates, whose serializer refuses a live WebSocket (DataCloneError). Without it, an upgrade Response received from a capnweb session materializes one hop too early and dies on the next internal Workers-RPC hop. This composes with the capnweb-in-a-Worker, native-RPC-to-the-DO topology upstream recommends (capnweb issue #36). The pair is byte-oriented because workerd RPC proxies streams as byte pipes only (a value-chunk stream fails with "This ReadableStream did not return bytes"), so the tunnel's text/binary/close frames travel in a length-prefixed framing spoken only by the deferring session and materializeUpgrade(). Since the two speakers are independently deployed workers, the header layout is frozen and evolution is append-only (an unknown frame type tears the tunnel down loudly). The capnweb wire format is unchanged; deferral is purely a receive-side choice, threaded to the deserializer through the Importer interface the same way upstream threads RpcLimits (as an optional method, so only the RPC session implements it). Ownership follows the ordinary semantics of streams received over RPC: the pair's inner ends stay payload-owned for the tunnel's whole life (first use locks them but takes no ownership), and the outer streams use highWaterMark 0 so nothing pulls -- or locks -- before a real read; an untouched deferred Response releases the tunnel on payload disposal, like an unclaimed TunneledWebSocket. materializeUpgrade() consumes the pair, locking both streams synchronously so a second call fails fast. The decoder tolerates arbitrary re-chunking, assembles split frames in linear time, caps the frame length on both the encode and decode side, and fails closed -- aborting the sender's socket -- on malformed frames and on truncation at end-of-stream. Flow control is inherited end to end: stream acks fire as the pair is read, so an unconsumed pair keeps the provider throttled to the flow-control window (asserted by test), and a materialized endpoint behaves like a non-deferred one. Reviewed with two staged multi-agent adversarial passes (eleven lenses total; every bug claim independently reproduced, several fixes empirically probed down to a raw-socket RFC 6455 handshake against the materialized 101) plus upstream-maintainer design evidence. The review removed a parallel writable mode on TunneledWebSocket in favor of the existing WritableStreamStubHook machinery, collapsed four stream wrappers into two, made the deferred branch a one-line call in serialize.ts, and reverted map.ts to byte-identical with upstream. Tested on Node (framing pinned byte-for-byte, split/coalesced/large/ empty frames, header carriage, mixed upgrade and plain traffic on one deferring session, window-bounded backpressure of an unread pair, double-materialize fail-fast, malformed and oversized and truncated frames failing closed, abort/cancel teardown, ignored-params release, RPC-session death yielding error + close 1006, and the full session battery over a deferred-then-materialized socket) and on workerd (in-isolate defer/materialize with a real WebSocketPair; the pair crossing a native service-binding hop in a call param AND in a return payload -- the relay-to-DO shape -- with init and close propagation asserted through both boundaries; and a REAL HTTP upgrade served from a fetch handler via materializeUpgrade, subprotocol echoed on the actual 101). Co-Authored-By: Claude Fable 5 --- .changeset/defer-upgrade-materialization.md | 7 + README.md | 6 + __tests__/test-server-workerd.js | 92 +++- __tests__/websocket-tunnel.test.ts | 514 +++++++++++++++++++- __tests__/workerd.test.ts | 134 ++++- src/index.ts | 6 +- src/rpc.ts | 44 ++ src/serialize.ts | 20 +- src/websocket-streams.ts | 335 +++++++++++++ 9 files changed, 1151 insertions(+), 7 deletions(-) create mode 100644 .changeset/defer-upgrade-materialization.md diff --git a/.changeset/defer-upgrade-materialization.md b/.changeset/defer-upgrade-materialization.md new file mode 100644 index 00000000..45bd371b --- /dev/null +++ b/.changeset/defer-upgrade-materialization.md @@ -0,0 +1,7 @@ +--- +"@iterate-com/capnweb": minor +--- + +Tunneled WebSocket upgrade Responses now arrive in forwardable form by default on Cloudflare Workers: `Response.webSocket` deserializes to an opaque `DeferredWebSocketUpgrade` — a `{ readable, writable, init }` byte-stream pair — instead of an eagerly materialized `WebSocketPair` end. The pair can be carried across hops that serialize byte streams but not sockets (in particular native Workers RPC between isolates, whose serializer refuses a live WebSocket), and the new `materializeUpgrade(pair, init?)` export rebuilds a real upgrade `Response` at the hop that actually serves it; `pair.init` carries the provider's upgrade headers (e.g. a negotiated `Sec-WebSocket-Protocol`) so they survive to the served 101. On other runtimes the previous behavior (a usable `TunneledWebSocket`) remains the default. The new `deferUpgradeMaterialization` session option overrides the default in either direction — set it to `false` on Workers to restore eager materialization when the session endpoint itself serves the upgrade. Receive-side only: the wire format is unchanged. + +Behavior change on Workers only: code that received a tunneled upgrade over a capnweb session on workerd and used `response.webSocket` as a native socket at the session endpoint must either call `materializeUpgrade(response.webSocket)` or set `deferUpgradeMaterialization: false` on the session. diff --git a/README.md b/README.md index 53b8961c..4207222b 100644 --- a/README.md +++ b/README.md @@ -7,6 +7,12 @@ > Delta vs upstream: > - **WebSocket-over-RPC** — a `Response` with a Workers-style `webSocket` > upgrade can be passed over RPC (tunneled as a stream pair). +> - **Deferred upgrade materialization** — on Cloudflare Workers a tunneled +> upgrade arrives by default as an opaque byte-stream pair that can cross +> hops which serialize byte streams but not sockets (e.g. native Workers +> RPC between isolates); `materializeUpgrade()` rebuilds the real socket at +> the hop that serves it. The `deferUpgradeMaterialization` session option +> overrides the default in either direction. > - **`onCall` session option** — server-side per-call hook for observability > (used by Iterate OS ITX tracing); propagates through promise pipelining. > - Small `Provider` type tweak for better go-to-definition through stubs. diff --git a/__tests__/test-server-workerd.js b/__tests__/test-server-workerd.js index 6a40e4af..38027b2a 100644 --- a/__tests__/test-server-workerd.js +++ b/__tests__/test-server-workerd.js @@ -10,7 +10,7 @@ // build step for it. Instead, we're getting by configuring the worker in vitest.config.ts by // just specifying the raw JS modules. -import { newWorkersRpcResponse } from "../dist/index-workers.js"; +import { newWorkersRpcResponse, newWebSocketRpcSession, materializeUpgrade } from "../dist/index-workers.js"; import { RpcTarget, DurableObject } from "cloudflare:workers"; // TODO(cleanup): At present we clone the implementation of Counter and TestTarget because @@ -81,8 +81,41 @@ export class TestTarget extends RpcTarget { } } +// See openDeferredEchoRemote below. +let deferredEchoRetainer = null; + export default { async fetch(req, env, ctx) { + // The design doc's "fetch-lane exit", end to end: defer a tunneled upgrade from an + // in-isolate capnweb session, materialize it, and serve it as a REAL HTTP upgrade from + // this fetch handler -- headers carried on pair.init and all. + if (new URL(req.url).pathname === "/deferred-upgrade") { + let origin = new WebSocketPair(); + origin[1].accept(); + origin[1].addEventListener("message", event => origin[1].send(event.data)); + + class Target extends RpcTarget { + openEcho() { + return new Response(null, { + status: 101, + webSocket: origin[0], + headers: { "sec-websocket-protocol": "itx-v1", "x-provider-custom": "survives" }, + }); + } + } + + // (Holding a capnweb session in this request-scoped context triggers workerd + // hang-detector warnings after the test finishes; production holds sessions in a + // long-lived relay WebSocket context. Accepted as test-only noise.) + let pair = new WebSocketPair(); + pair[0].accept(); + pair[1].accept(); + let api = newWebSocketRpcSession(pair[0]); // deferred by default on Workers + newWebSocketRpcSession(pair[1], new Target()); + let response = await api.openEcho(); + return materializeUpgrade(response.webSocket); + } + return newWorkersRpcResponse(req, new TestTarget(env), { onSendError(err) { return err; } }); @@ -90,5 +123,62 @@ export default { async greet(name, env, ctx) { return `Hello, ${name}!`; + }, + + // The relay shape: THIS isolate holds a deferring capnweb session and returns the raw pair + // in a native RPC return payload; the caller materializes after the call settles. The session + // and the delivering Response must be retained past the call (module state) -- disposing them + // would release the payload-owned tunnel -- which is exactly the retention contract a real + // relay must follow. + async openDeferredEchoRemote(unused, env, ctx) { + let origin = new WebSocketPair(); + origin[1].accept(); + origin[1].addEventListener("message", event => origin[1].send(event.data)); + let closeEvent = null; + origin[1].addEventListener("close", + event => { closeEvent = { code: event.code, reason: event.reason }; }); + + class Target extends RpcTarget { + openEcho() { + return new Response(null, { + status: 101, + webSocket: origin[0], + headers: { "sec-websocket-protocol": "itx-v1" }, + }); + } + } + + // (Session held in an RPC-call context: same accepted hang-detector noise as the + // /deferred-upgrade route above.) + let pair = new WebSocketPair(); + pair[0].accept(); + pair[1].accept(); + let api = newWebSocketRpcSession(pair[0]); // deferred by default on Workers + newWebSocketRpcSession(pair[1], new Target()); + let response = await api.openEcho(); + deferredEchoRetainer = { api, response, getClose: () => closeEvent }; + return response.webSocket; + }, + + async deferredEchoCloseEvent(unused, env, ctx) { + return deferredEchoRetainer ? deferredEchoRetainer.getClose() : null; + }, + + // Native workers-RPC leg of deferred upgrade materialization: receives a tunneled socket as a + // raw { readable, writable } pair (which workerd RPC serializes; a live WebSocket it refuses), + // materializes the socket in THIS isolate, and runs one echo round trip through it. + async materializeEcho(webSocket, env, ctx) { + let response = materializeUpgrade(webSocket); + let socket = response.webSocket; + socket.accept(); + let reply = new Promise(resolve => { + socket.addEventListener("message", event => resolve(event.data), { once: true }); + }); + socket.send("ping across isolates"); + try { + return await reply; + } finally { + socket.close(1000, ""); + } } } diff --git a/__tests__/websocket-tunnel.test.ts b/__tests__/websocket-tunnel.test.ts index 2f92c1c5..47e1831d 100644 --- a/__tests__/websocket-tunnel.test.ts +++ b/__tests__/websocket-tunnel.test.ts @@ -5,7 +5,8 @@ import { expect, it, describe, beforeAll, afterAll } from "vitest"; import type { AddressInfo } from "node:net"; import { WebSocket as NodeWebSocket, WebSocketServer } from "ws"; -import { newWebSocketRpcSession, RpcTarget } from "../src/index.js"; +import { newWebSocketRpcSession, materializeUpgrade, RpcSession, RpcTarget, + type RpcTransport } from "../src/index.js"; import { TestTarget } from "./test-util.js"; import { registerSessionTestBattery } from "./session-battery.js"; @@ -300,3 +301,514 @@ describe("Cap'n Web over a WebSocket obtained via fetch() over Cap'n Web", () => }; }); }); + +// The deferUpgradeMaterialization session option delivers an upgrade Response's tunneled socket +// as an opaque framed-byte { readable, writable, init } pair, for infrastructure that needs to +// carry the tunnel across a further hop (one that can serialize byte streams but not sockets) +// before materializeUpgrade() turns it back into a real socket at the final hop. +describe("deferred upgrade materialization", () => { + let echoServer = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + echoServer.on("connection", socket => { + socket.on("message", (data, isBinary) => socket.send(data, { binary: isBinary })); + }); + + // A server that greets in response to any message ("knock"), plus a log of its server-side + // sockets so tests can observe closure from the provider's side. (The greeting is + // knock-triggered rather than sent on connect because a `ws` socket's early messages can + // arrive before the tunnel attaches its listeners; see webSocketToStreams.) + let greetingSockets: any[] = []; + let greetingServer = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + greetingServer.on("connection", socket => { + greetingSockets.push(socket); + socket.on("message", () => socket.send("welcome!")); + }); + + class DeferredTestApi extends RpcTarget { + constructor(private echoPort: number, private greetingPort: number) { super(); } + + async openEcho(): Promise { + return responseWithWebSocket(await openWebSocket(this.echoPort)); + } + + async openGreeting(): Promise { + return responseWithWebSocket(await openWebSocket(this.greetingPort)); + } + + // An upgrade Response carrying a negotiated subprotocol header, to test that headers + // survive deferral and re-materialization. + async openEchoWithProtocol(): Promise { + let response = new Response(null, { headers: { "sec-websocket-protocol": "itx" } }); + Object.defineProperty(response, "webSocket", + { value: await openWebSocket(this.echoPort) }); + return response; + } + + async openPlain(): Promise { + return new Response("plain traffic", { status: 418 }); + } + } + + // Identity adapter that re-chunks a byte stream into pieces of at most `size` bytes, + // simulating a transport that splits frames at arbitrary boundaries (as a real workerd RPC + // byte pipe may). + function rechunkBytes(readable: ReadableStream, size: number): ReadableStream { + let reader = readable.getReader(); + let buffer = new Uint8Array(0); + return new ReadableStream({ + async pull(controller) { + while (buffer.length === 0) { + let { done, value } = await reader.read(); + if (done) { + controller.close(); + return; + } + buffer = value; + } + controller.enqueue(buffer.subarray(0, Math.min(size, buffer.length))); + buffer = buffer.length > size ? buffer.subarray(size) : new Uint8Array(0); + }, + cancel(reason) { + return reader.cancel(reason); + }, + }); + } + + function textFrame(text: string): Uint8Array { + let payload = new TextEncoder().encode(text); + let frame = new Uint8Array(5 + payload.length); + frame[0] = 0; + new DataView(frame.buffer).setUint32(1, payload.length); + frame.set(payload, 5); + return frame; + } + + let rpcServer = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + let clientSocket: NodeWebSocket; + let api: any; + + beforeAll(async () => { + let [echoPort, greetingPort, rpcPort] = await Promise.all([ + listening(echoServer), listening(greetingServer), listening(rpcServer)]); + rpcServer.on("connection", socket => { + newWebSocketRpcSession(socket as any, new DeferredTestApi(echoPort, greetingPort)); + }); + clientSocket = new NodeWebSocket(`ws://127.0.0.1:${rpcPort}`); + api = newWebSocketRpcSession(clientSocket as any, undefined, + { deferUpgradeMaterialization: true }); + }); + + afterAll(async () => { + clientSocket.close(); + await Promise.all([ + closeServer(rpcServer), closeServer(echoServer), closeServer(greetingServer)]); + }); + + it("delivers the tunnel as an opaque byte-stream pair", async () => { + let response = await api.openGreeting(); + let webSocket = response.webSocket; + expect(webSocket.readable).toBeInstanceOf(ReadableStream); + expect(webSocket.writable).toBeInstanceOf(WritableStream); + + // The pair carries the tunnel's frames as opaque bytes -- so it can cross byte-oriented + // boundaries like native workerd RPC. Speak the internal framing by hand (1 byte type, + // 4 bytes big-endian length, payload; type 0 = text) to knock, and decode the reply. This + // deliberately pins the frame format byte-for-byte: the encoder and decoder ship together, + // so a drifted format would still round-trip green in every other test here and only fail + // against an already-deployed peer -- in the worst case as nothing but a hung tunnel. + let writer = webSocket.writable.getWriter(); + let reader = webSocket.readable.getReader(); + let knock = new TextEncoder().encode("knock"); + let frame = new Uint8Array(5 + knock.length); + frame[0] = 0; + new DataView(frame.buffer).setUint32(1, knock.length); + frame.set(knock, 5); + await writer.write(frame); + + let { value } = await reader.read(); + expect(value).toBeInstanceOf(Uint8Array); + expect(value![0]).toBe(0); // a text frame + expect(new DataView(value!.buffer, value!.byteOffset).getUint32(1)).toBe(value!.length - 5); + expect(new TextDecoder().decode(value!.subarray(5))).toBe("welcome!"); + + // Closing the writable closes the provider's socket. + let serverSocket = greetingSockets[greetingSockets.length - 1]; + let closed = new Promise(resolve => serverSocket.once("close", resolve)); + await writer.close(); + await closed; + }); + + it("materializes a real upgrade Response with materializeUpgrade()", async () => { + let response = await api.openEcho(); + let materialized = materializeUpgrade(response.webSocket); + let socket: any = (materialized as any).webSocket; + expect(socket).toBeTruthy(); + + let echoed = nextEvent(socket, "message"); + socket.send("hello"); + expect((await echoed).data).toBe("hello"); + + let binaryEchoed = nextEvent(socket, "message"); + socket.send(new Uint8Array([4, 5, 6])); + let bytes = (await binaryEchoed).data; + expect(typeof bytes).not.toBe("string"); + expect(Array.from(bytes)).toEqual([4, 5, 6]); + + let closeEvent = nextEvent(socket, "close"); + socket.close(1000, "done"); + expect(await closeEvent).toMatchObject({ code: 1000, reason: "done" }); + expect(socket.readyState).toBe(3); // CLOSED + }); + + it("survives arbitrary re-chunking of the byte pair", async () => { + // A transport carrying the pair may split frames at any byte boundary; FrameDecoder must + // reassemble them. Split into 3-byte chunks (every header is split), and also push a large + // frame through moderate chunks. + let webSocket = (await api.openEcho()).webSocket; + let materialized = materializeUpgrade({ + readable: rechunkBytes(webSocket.readable, 3), + writable: webSocket.writable, + }); + let socket: any = (materialized as any).webSocket; + + let echoed = nextEvent(socket, "message"); + socket.send("split me"); + expect((await echoed).data).toBe("split me"); + + // An empty text frame is a legal frame (zero-length payload). + let emptyEchoed = nextEvent(socket, "message"); + socket.send(""); + expect((await emptyEchoed).data).toBe(""); + socket.close(1000, ""); + }); + + it("survives re-chunking of a large binary frame", async () => { + let webSocket = (await api.openEcho()).webSocket; + let materialized = materializeUpgrade({ + readable: rechunkBytes(webSocket.readable, 64 * 1024), + writable: webSocket.writable, + }); + let socket: any = (materialized as any).webSocket; + + let big = new Uint8Array(1024 * 1024); + for (let i = 0; i < big.length; i += 4096) big[i] = i & 0xff; + let echoed = nextEvent(socket, "message"); + socket.send(big); + let bytes = (await echoed).data; + expect(bytes.length).toBe(big.length); + expect(Array.from(bytes.subarray(0, 16))).toEqual(Array.from(big.subarray(0, 16))); + expect(bytes[8192]).toBe(big[8192]); + socket.close(1000, ""); + }); + + it("decodes several frames coalesced into one chunk", async () => { + // The inverse of re-chunking: one byte chunk carrying two whole frames. + let webSocket = (await api.openGreeting()).webSocket; + let writer = webSocket.writable.getWriter(); + let reader = webSocket.readable.getReader(); + + let knock = textFrame("knock"); + let both = new Uint8Array(knock.length * 2); + both.set(knock); + both.set(knock, knock.length); + await writer.write(both); + + // Two knocks produce two greetings. + expect((await reader.read()).value![0]).toBe(0); + expect((await reader.read()).value![0]).toBe(0); + await writer.close(); + }); + + it("materializes headers carried on the pair (e.g. a negotiated subprotocol)", async () => { + let response = await api.openEchoWithProtocol(); + expect(response.headers.get("sec-websocket-protocol")).toBe("itx"); + + // The pair itself carries the init, so infrastructure that forwards the whole pair object + // across a hop preserves the headers without extra plumbing. + let webSocket = response.webSocket; + expect(webSocket.init).toBeTruthy(); + let materialized = materializeUpgrade(webSocket); + expect(materialized.headers.get("sec-websocket-protocol")).toBe("itx"); + (materialized as any).webSocket.close(1000, ""); + }); + + it("refuses to materialize the same pair twice", async () => { + let webSocket = (await api.openEcho()).webSocket; + let materialized = materializeUpgrade(webSocket); + let socket: any = (materialized as any).webSocket; + + // The second call must fail fast (the pair is consumed), leaving the first materialization + // fully working. + expect(() => materializeUpgrade(webSocket)).toThrow(TypeError); + + let echoed = nextEvent(socket, "message"); + socket.send("still mine"); + expect((await echoed).data).toBe("still mine"); + socket.close(1000, ""); + }); + + it("tears the tunnel down on a malformed frame instead of leaving it half-open", async () => { + let webSocket = (await api.openGreeting()).webSocket; + let serverSocket = greetingSockets[greetingSockets.length - 1]; + let closed = new Promise(resolve => serverSocket.once("close", resolve)); + + // Frame type 7 does not exist. The write must reject AND the provider's socket must be + // torn down (fail closed), not linger half-open. + let writer = webSocket.writable.getWriter(); + await expect(writer.write(new Uint8Array([7, 0, 0, 0, 0]))).rejects.toThrow(/frame type/); + await closed; + + // Likewise a header declaring an absurd length: fail closed immediately rather than + // buffering toward 4 GiB. + let webSocket2 = (await api.openGreeting()).webSocket; + let serverSocket2 = greetingSockets[greetingSockets.length - 1]; + let closed2 = new Promise(resolve => serverSocket2.once("close", resolve)); + let writer2 = webSocket2.writable.getWriter(); + await expect(writer2.write(new Uint8Array([0, 255, 255, 255, 255]))) + .rejects.toThrow(/maximum/); + await closed2; + }); + + it("treats a truncated frame at end-of-stream as an error, not a clean close", async () => { + // Read direction: a byte stream that ends mid-frame means the transport lost data. The + // materialized socket must fail (1006), not report a normal closure. + let partial = textFrame("lost message").subarray(0, 7); + let materialized = materializeUpgrade({ + readable: new ReadableStream({ + start(controller) { + controller.enqueue(partial); + controller.close(); + }, + }), + writable: new WritableStream(), + }); + let socket: any = (materialized as any).webSocket; + let error = nextEvent(socket, "error"); + let close = nextEvent(socket, "close"); + socket.accept(); + expect((await error).error.message).toMatch(/truncated/); + expect(await close).toMatchObject({ code: 1006 }); + + // Write direction: closing the pair's writable with a partial frame buffered must reject + // and tear the provider's socket down. + let webSocket = (await api.openGreeting()).webSocket; + let serverSocket = greetingSockets[greetingSockets.length - 1]; + let closed = new Promise(resolve => serverSocket.once("close", resolve)); + let writer = webSocket.writable.getWriter(); + await writer.write(partial); + await expect(writer.close()).rejects.toThrow(/truncated/); + await closed; + }); + + it("closes the provider socket when the pair is aborted or canceled", async () => { + // These are the teardown verbs infrastructure (or a workerd RPC proxy whose far end died) + // drives on the raw pair. Aborting the writable tears the provider down eagerly. + { + let webSocket = (await api.openGreeting()).webSocket; + let serverSocket = greetingSockets[greetingSockets.length - 1]; + let closed = new Promise(resolve => serverSocket.once("close", resolve)); + await webSocket.writable.abort(new Error("infra teardown")); + await closed; + } + // Canceling the readable propagates lazily, like any received stream's cancel: the sender + // notices when it next has a message to deliver, then tears down. + { + let webSocket = (await api.openGreeting()).webSocket; + let serverSocket = greetingSockets[greetingSockets.length - 1]; + let closed = new Promise(resolve => serverSocket.once("close", resolve)); + await webSocket.readable.cancel(); + let writer = webSocket.writable.getWriter(); + await writer.write(textFrame("knock")); // provokes a greeting nobody can receive + await closed; + } + }); + + it("releases the tunnel when a deferred Response received in params is ignored", async () => { + // The deferred mirror of the base feature's ignore() test: a server session with + // deferUpgradeMaterialization receives an upgrade Response in params, never touches the + // pair, and returns -- payload disposal must release the tunnel and close the socket. + class DeferredIgnorer extends RpcTarget { + async ignore(response: Response): Promise {} + } + let server = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + server.on("connection", s => { + newWebSocketRpcSession(s as any, new DeferredIgnorer(), + { deferUpgradeMaterialization: true }); + }); + let socket = new NodeWebSocket(`ws://127.0.0.1:${await listening(server)}`); + try { + let ignorer: any = newWebSocketRpcSession(socket as any); + let provider = await openWebSocket((echoServer.address() as AddressInfo).port); + let closed = nextEvent(provider, "close"); + await ignorer.ignore(responseWithWebSocket(provider)); + await closed; + } finally { + socket.close(); + await closeServer(server); + } + }); + + it("leaves non-upgrade Responses untouched (mixed traffic)", async () => { + // A deferring session (the relay) carries plain HTTP through live capabilities too; only + // webSocket-bearing Responses are affected by the option. + let response = await api.openPlain(); + expect(response.status).toBe(418); + expect(await response.text()).toBe("plain traffic"); + expect((response as any).webSocket ?? null).toBeNull(); + }); + + it("throttles the provider to the flow-control window while the pair is unread", async () => { + // The bounded-memory property the relay leans on: capnweb stream acks fire as the pair is + // READ, so an unconsumed deferred pair keeps the provider-side sender throttled to the + // flow-control window instead of letting a firehose pool at the deferring receiver. + // + // In-memory transport pair, counting the bytes the provider side emits. Deterministic: with + // zero reads there are zero acks, so the window never advances no matter how long we wait. + let providerSent = 0; + function pipeTransports(): [RpcTransport, RpcTransport] { + function makeEndpoint(counted: boolean) { + let queue: string[] = []; + let waiting: ((m: string) => void) | undefined; + return { + deliver(message: string) { + if (waiting) { let w = waiting; waiting = undefined; w(message); } + else queue.push(message); + }, + transport: undefined as unknown as RpcTransport, + make(peer: () => { deliver(m: string): void }): RpcTransport { + return { + send(message: string) { + if (counted) providerSent += message.length; + peer().deliver(message); + }, + receive() { + if (queue.length > 0) return Promise.resolve(queue.shift()!); + return new Promise(resolve => { waiting = resolve; }); + }, + }; + }, + }; + } + let a = makeEndpoint(true), b = makeEndpoint(false); + return [a.make(() => b), b.make(() => a)]; + } + + // A minimal provider-side socket the test can flood on demand. + let listeners = new Map void)[]>(); + let providerSocket = { + send(data: unknown) {}, + close(code?: number, reason?: string) { + for (let l of listeners.get("close") ?? []) l({ code: code ?? 1005, reason: reason ?? "" }); + }, + addEventListener(type: string, listener: (event: any) => void) { + let list = listeners.get(type) ?? []; + list.push(listener); + listeners.set(type, list); + }, + }; + + class FloodApi extends RpcTarget { + async openFlood(): Promise { + return responseWithWebSocket(providerSocket); + } + } + + let [providerTransport, receiverTransport] = pipeTransports(); + new RpcSession(providerTransport, new FloodApi()); + let receiver = new RpcSession(receiverTransport, undefined, + { deferUpgradeMaterialization: true }); + let response: any = await (receiver.getRemoteMain() as any).openFlood(); + let webSocket = response.webSocket; + + let settle = () => new Promise(resolve => setTimeout(resolve, 50)); + let baseline = providerSent; + let message = "x".repeat(32 * 1024); + for (let i = 0; i < 1024; i++) { + for (let l of listeners.get("message") ?? []) l({ data: message }); + } + await settle(); + let unreadSent = providerSent - baseline; + // 32 MiB was offered; only about one flow-control window's worth may cross. + expect(unreadSent).toBeGreaterThan(0); + expect(unreadSent).toBeLessThan(4 * 1024 * 1024); + + // Reading the pair acks chunks and lets more flow. + let reader = webSocket.readable.getReader(); + for (let i = 0; i < 8; i++) await reader.read(); + await settle(); + expect(providerSent - baseline).toBeGreaterThan(unreadSent); + expect(providerSent - baseline).toBeLessThan(16 * 1024 * 1024); + }); + + it("fails closed when the RPC session dies under a materialized socket", async () => { + // A dedicated session, so terminating the transport doesn't break the shared one. + let transport = new NodeWebSocket( + `ws://127.0.0.1:${(rpcServer.address() as AddressInfo).port}`); + let session: any = newWebSocketRpcSession(transport as any, undefined, + { deferUpgradeMaterialization: true }); + let materialized = materializeUpgrade((await session.openEcho()).webSocket); + let socket: any = (materialized as any).webSocket; + + let echoed = nextEvent(socket, "message"); + socket.send("alive"); + expect((await echoed).data).toBe("alive"); + + // Kill the RPC transport abnormally: the socket must fail (error and/or close 1006), not + // hang and not report a clean close. + let close = nextEvent(socket, "close"); + transport.terminate(); + expect(await close).toMatchObject({ code: 1006 }); + expect(socket.readyState).toBe(3); // CLOSED + }); +}); + +// Prove that deferring and re-materializing changes nothing about the socket's behavior: run the +// full session battery over a socket that took the long way around -- tunneled over RPC, +// delivered as the opaque byte pair, then rebuilt by materializeUpgrade() -- mirroring the +// "Cap'n Web over a WebSocket obtained via fetch() over Cap'n Web" battery above. +describe("Cap'n Web over a deferred-then-materialized WebSocket", () => { + let innerServer = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + innerServer.on("connection", socket => { + newWebSocketRpcSession(socket as any, new TestTarget()); + }); + + class Gateway extends RpcTarget { + async fetch(request: Request): Promise { + if (request.headers.get("Upgrade")?.toLowerCase() !== "websocket") { + return new Response("Expected a WebSocket upgrade.", { status: 426 }); + } + let port = (innerServer.address() as AddressInfo).port; + return responseWithWebSocket(await openWebSocket(port)); + } + } + + let gatewayServer = new WebSocketServer({ host: "127.0.0.1", port: 0 }); + gatewayServer.on("connection", socket => { + newWebSocketRpcSession(socket as any, new Gateway()); + }); + + afterAll(async () => { + await Promise.all([closeServer(gatewayServer), closeServer(innerServer)]); + }); + + registerSessionTestBattery(async () => { + let gatewayPort = await listening(gatewayServer); + await listening(innerServer); + + let gatewaySocket = new NodeWebSocket(`ws://127.0.0.1:${gatewayPort}`); + let gateway: any = newWebSocketRpcSession(gatewaySocket as any, undefined, + { deferUpgradeMaterialization: true }); + let response = await gateway.fetch(new Request("https://inner.example/rpc", { + headers: { Upgrade: "websocket" }, + })); + + let materialized = materializeUpgrade(response.webSocket); + let stub = newWebSocketRpcSession((materialized as any).webSocket); + return { + stub, + async [Symbol.asyncDispose]() { + stub[Symbol.dispose](); + gatewaySocket.close(); + }, + }; + }); +}); diff --git a/__tests__/workerd.test.ts b/__tests__/workerd.test.ts index aa685926..ed6fd2cd 100644 --- a/__tests__/workerd.test.ts +++ b/__tests__/workerd.test.ts @@ -5,7 +5,7 @@ /// import { expect, it, describe } from "vitest"; import { RpcStub as NativeRpcStub, RpcTarget as NativeRpcTarget, env, DurableObject } from "cloudflare:workers"; -import { newHttpBatchRpcSession, newWebSocketRpcSession, RpcStub, RpcTarget } from "../src/index-workers.js"; +import { newHttpBatchRpcSession, newWebSocketRpcSession, materializeUpgrade, RpcStub, RpcTarget } from "../src/index-workers.js"; import { v, wrapServerTarget, type ServiceValidator } from "../packages/capnweb-validate/src/internal/core.js"; import { Counter, TestTarget } from "./test-util.js"; @@ -256,7 +256,10 @@ describe("workerd RPC server", () => { let pair = new WebSocketPair(); pair[0].accept(); pair[1].accept(); - let api: any = newWebSocketRpcSession(pair[0]); + // On Workers, upgrades arrive deferred by default; this session opts out because it serves + // the socket right here at the session endpoint. + let api: any = newWebSocketRpcSession(pair[0], undefined, + { deferUpgradeMaterialization: false }); newWebSocketRpcSession(pair[1], new WebSocketResponseTarget()); // On workerd, the received Response holds a native WebSocket, suitable for completing a real @@ -276,6 +279,133 @@ describe("workerd RPC server", () => { socket!.close(); }); + // Opens an in-isolate deferring session against an echoing WebSocketPair target, returning + // the deferred Response plus a promise for the close event observed at the origin end -- so + // tests can assert closure makes it all the way back. + function openDeferredEcho(): { + response: Promise, + originClose: Promise<{ code: number, reason: string }>, + } { + let resolveClose: (event: { code: number, reason: string }) => void; + let originClose = new Promise<{ code: number, reason: string }>(resolve => { + resolveClose = resolve; + }); + + class EchoUpgradeTarget extends RpcTarget { + openEcho() { + let pair = new WebSocketPair(); + pair[1].accept(); + pair[1].addEventListener("message", event => pair[1].send(event.data)); + pair[1].addEventListener("close", + event => resolveClose({ code: event.code, reason: event.reason })); + return new Response(null, { status: 101, webSocket: pair[0] }); + } + } + + let pair = new WebSocketPair(); + pair[0].accept(); + pair[1].accept(); + let api: any = newWebSocketRpcSession(pair[0]); // deferred by default on Workers + newWebSocketRpcSession(pair[1], new EchoUpgradeTarget()); + return { response: api.openEcho(), originClose }; + } + + it("defers upgrade materialization by default on Workers", async () => { + // On Workers a tunneled upgrade arrives as the transportable pair rather than a + // materialized socket (which could never cross another RPC hop anyway). + let response = await openDeferredEcho().response; + let webSocket: any = (response as any).webSocket; + expect(webSocket instanceof WebSocket).toBe(false); + expect(webSocket.readable).toBeInstanceOf(ReadableStream); + expect(webSocket.writable).toBeInstanceOf(WritableStream); + + // materializeUpgrade() is the final hop: a real 101 carrying a native socket, suitable for + // completing an actual HTTP upgrade. + let materialized = materializeUpgrade(webSocket); + expect(materialized.status).toBe(101); + let socket = materialized.webSocket!; + socket.accept(); + let message = new Promise(resolve => { + socket.addEventListener("message", event => resolve(event.data), { once: true }); + }); + socket.send("hello through a deferred tunnel"); + expect(await message).toBe("hello through a deferred tunnel"); + socket.close(); + }); + + it("carries a deferred stream pair across a native workers RPC boundary", async () => { + // The reason deferUpgradeMaterialization exists: workerd's RPC serializer refuses a live + // WebSocket but carries ReadableStream/WritableStream natively. Receive a tunnel deferred, + // hand the raw pair across a real service-binding RPC hop, and let the OTHER isolate + // materialize the socket and run an echo round trip through the full path: + // far isolate <-native RPC streams-> this isolate <-capnweb tunnel-> echoing pair end. + let { response, originClose } = openDeferredEcho(); + let webSocket: any = ((await response) as any).webSocket; + let reply = await ((env).testServer).materializeEcho(webSocket); + expect(reply).toBe("ping across isolates"); + + // The far isolate closed its materialized socket with 1000; the close must traverse the + // native hop and the tunnel back to the origin pair end. + expect(await originClose).toMatchObject({ code: 1000 }); + }); + + it("carries a deferred pair in a native RPC return payload (the relay-to-DO shape)", async () => { + // The production topology is the inverse of the test above: the capnweb session lives in + // the FAR isolate (the relay), the pair crosses in the RETURN payload of a native RPC + // call, and this isolate materializes and uses the socket after that call has settled. + let webSocket: any = await ((env).testServer).openDeferredEchoRemote(null); + expect(webSocket.readable).toBeInstanceOf(ReadableStream); + expect(webSocket.writable).toBeInstanceOf(WritableStream); + // init rides the pair object across the native hop as plain data. + expect(webSocket.init).toBeTruthy(); + + let materialized = materializeUpgrade(webSocket); + expect(materialized.status).toBe(101); + expect(materialized.headers.get("sec-websocket-protocol")).toBe("itx-v1"); + let socket = materialized.webSocket!; + socket.accept(); + let message = new Promise(resolve => { + socket.addEventListener("message", event => resolve(event.data), { once: true }); + }); + socket.send("ping on the return path"); + expect(await message).toBe("ping on the return path"); + + // Close from the materialized end; the origin -- two hops away, behind the native + // boundary AND the capnweb tunnel, in the other isolate -- must observe it. + socket.close(1000, "done"); + let closeEvent: any = null; + for (let i = 0; i < 100 && !closeEvent; i++) { + await new Promise(resolve => setTimeout(resolve, 20)); + closeEvent = await ((env).testServer).deferredEchoCloseEvent(null); + } + expect(closeEvent).toMatchObject({ code: 1000, reason: "done" }); + }); + + it("serves a real HTTP upgrade from a fetch handler via materializeUpgrade", async () => { + // The design doc's fetch-lane exit, literally: the aux worker defers a tunneled upgrade + // from an in-isolate capnweb session, materializes it, and returns it from its fetch + // handler. We complete a real WebSocket handshake against it over the service binding and + // check the provider's negotiated subprotocol (carried on pair.init) reached the real 101. + // (Over a service binding the 101 is fabricated in-process: headers pass through verbatim + // and no Sec-WebSocket-Accept is computed. Real-wire handshake sanitization -- reserved + // headers recomputed/dropped, Accept derived from the client key -- is owned by workerd/kj + // and was verified externally with a raw-socket probe during review.) + let resp = await (env).testServer.fetch("http://foo/deferred-upgrade", { + headers: { Upgrade: "websocket", "Sec-WebSocket-Protocol": "itx-v1" }, + }); + expect(resp.status).toBe(101); + expect(resp.headers.get("sec-websocket-protocol")).toBe("itx-v1"); + expect(resp.headers.get("x-provider-custom")).toBe("survives"); + let ws = resp.webSocket!; + ws.accept(); + let message = new Promise(resolve => { + ws.addEventListener("message", event => resolve(event.data), { once: true }); + }); + ws.send("through the fetch lane"); + expect(await message).toBe("through the fetch lane"); + ws.close(1000, ""); + }); + it("can accept WebSocket RPC connections", async () => { let resp = await (env).testServer.fetch("http://foo", {headers: {Upgrade: "websocket"}}); let ws = resp.webSocket; diff --git a/src/index.ts b/src/index.ts index 97d68be0..396ffd2a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -13,6 +13,7 @@ import { newWebSocketRpcSession as newWebSocketRpcSessionImpl, import { newHttpBatchRpcSession as newHttpBatchRpcSessionImpl, newHttpBatchRpcResponse, nodeHttpBatchRpcResponse } from "./batch.js"; import { newMessagePortRpcSession as newMessagePortRpcSessionImpl } from "./messageport.js"; +import { materializeUpgrade, type DeferredWebSocketUpgrade } from "./websocket-streams.js"; import { forceInitMap } from "./map.js"; import { forceInitStreams } from "./streams.js"; @@ -21,9 +22,10 @@ forceInitStreams(); // Re-export public API types. export { serialize, deserialize, newWorkersWebSocketRpcResponse, newHttpBatchRpcResponse, - nodeHttpBatchRpcResponse, WebSocketTransport, DEFAULT_LIMITS, DEFAULT_MAX_DEPTH }; + nodeHttpBatchRpcResponse, WebSocketTransport, DEFAULT_LIMITS, DEFAULT_MAX_DEPTH, + materializeUpgrade }; export type { RpcCallInfo, RpcTransport, RpcTransportWithCustomEncoding, AnyRpcTransport, - RpcSessionOptions, RpcCompatible, EncodingLevel, RpcLimits }; + RpcSessionOptions, RpcCompatible, EncodingLevel, RpcLimits, DeferredWebSocketUpgrade }; // Hack the type system to make RpcStub's types work nicely! /** diff --git a/src/rpc.ts b/src/rpc.ts index 78de3bbd..9f417c71 100644 --- a/src/rpc.ts +++ b/src/rpc.ts @@ -467,6 +467,46 @@ export type RpcSessionOptions = { * spans the full asynchronous call. The hook is propagated through promise pipelining. */ onCall?: RpcCallHandler; + + /** + * Controls what a tunneled WebSocket upgrade Response's `webSocket` property deserializes to + * on this session: `true` delivers the opaque `{ readable, writable, init }` byte-stream + * pair (`DeferredWebSocketUpgrade`), `false` a materialized socket (a native WebSocketPair + * end on Cloudflare Workers, a TunneledWebSocket elsewhere). + * + * Unset, the runtime picks the shape that can actually be used there: **deferred on + * Cloudflare Workers** -- where a socket materialized at the session endpoint could never + * cross another RPC hop, the exact mistake behind + * `DataCloneError: Could not serialize object of type "WebSocket"` -- and **materialized + * everywhere else**, where there are no internal hops and the app wants a socket directly. + * Set `false` on Workers only when the session endpoint itself serves the upgrade and wants + * `response.webSocket` to be a native socket without calling `materializeUpgrade()`. + * + * The deferred pair exists to be carried across hops that serialize byte streams but not + * sockets -- e.g. a Workers RPC boundary between isolates -- and rebuilt into a real upgrade + * Response with `materializeUpgrade()` at the hop that serves it. Forward the whole pair + * object: `init` carries the provider's upgrade headers (e.g. a negotiated + * Sec-WebSocket-Protocol) so the materialized 101 can echo them. The deferred Response + * itself has status 200 (a constructed Response can't be 1xx) -- detect an upgrade by + * `response.webSocket != null`, never by status -- and must not be forwarded across a native + * boundary: its `webSocket` property is a local JS extension and silently disappears, + * leaving an ordinary-looking 200. + * + * The pair's inner ends follow the ordinary ownership semantics of streams received over + * RPC: they belong to the containing payload for the tunnel's whole life -- first use locks + * them but takes no ownership -- and disposing the payload releases the tunnel and closes + * the socket. So a pair received in call *params* cannot be kept past the call (there is no + * analog of TunneledWebSocket's accept()): finish forwarding it before returning, or receive + * the pair as a call *result* instead. The awaited result carries the payload; keep the + * delivering Response undisposed for as long as the tunnel is in use (undisposed results + * live until the session ends). Infrastructure that needs to reclaim per-tunnel state before + * then can wrap the pair's streams in observing pass-throughs before forwarding, or tie + * retention to whatever connection delivered the pair onward. + * + * Purely a receive-side choice: the wire format is unchanged, the sending side needs no + * corresponding option, and either side may run an older version. + */ + deferUpgradeMaterialization?: boolean; }; class RpcSessionImpl implements Importer, Exporter { @@ -716,6 +756,10 @@ class RpcSessionImpl implements Importer, Exporter { return this.limits; } + deferUpgradeMaterialization(): boolean | undefined { + return this.options.deferUpgradeMaterialization; + } + createPipe(readable: ReadableStream, readableHook: StubHook): ImportId { if (this.abortReason) throw this.abortReason; diff --git a/src/serialize.ts b/src/serialize.ts index a869c25b..4de1c54b 100644 --- a/src/serialize.ts +++ b/src/serialize.ts @@ -3,7 +3,8 @@ // https://opensource.org/license/mit import { StubHook, RpcPayload, typeForRpc, RpcStub, RpcPromise, LocatedPromise, RpcTarget, unwrapStubAndPath, streamImpl, PromiseStubHook, PayloadStubHook, type RpcCallHandler } from "./core.js"; -import { webSocketToStreams, makeUpgradeResponse } from "./websocket-streams.js"; +import { webSocketToStreams, makeUpgradeResponse, makeDeferredUpgradeResponse, + deferUpgradesByDefault } from "./websocket-streams.js"; export type ImportId = number; export type ExportId = number; @@ -745,6 +746,14 @@ export interface Importer { // Importer (rather than the Evaluator constructor) so that the per-session options reach the // Evaluator without changing how Evaluators are constructed throughout the codebase. getLimits(): RpcLimits; + + // The session's explicit `deferUpgradeMaterialization` setting, or undefined to use the + // runtime default (deferred on Cloudflare Workers, materialized elsewhere -- see + // RpcSessionOptions). Like getLimits(), surfaced through the Importer so the per-session + // option reaches the Evaluator at every construction site -- but optional, since only a real + // RPC session can meaningfully answer it (an importer that can't receive an upgrade Response + // at all simply doesn't implement it). + deferUpgradeMaterialization?(): boolean | undefined; } class NullImporter implements Importer { @@ -1101,6 +1110,15 @@ export class Evaluator { this.hooks.push(writableHook); delete init.webSocket; + + if (this.importer.deferUpgradeMaterialization?.() ?? deferUpgradesByDefault()) { + // Deliver the tunnel in forwardable form (the default on Workers, where an + // eagerly materialized socket could never cross another RPC hop anyway; a + // serving endpoint rebuilds the socket with materializeUpgrade()). See + // "Deferred materialization" in websocket-streams.ts. + return makeDeferredUpgradeResponse(readable, writableHook, init as ResponseInit); + } + return makeUpgradeResponse(readable, writableHook, init as ResponseInit); } diff --git a/src/websocket-streams.ts b/src/websocket-streams.ts index b7863434..eeac5fff 100644 --- a/src/websocket-streams.ts +++ b/src/websocket-streams.ts @@ -371,6 +371,341 @@ export function makeUpgradeResponse( } } +// ======================================================================================= +// Deferred materialization +// +// A deferring session -- the default on Cloudflare Workers, or any session with the +// `deferUpgradeMaterialization` option set -- delivers an upgrade Response's tunneled socket as +// an opaque `{ readable, writable, init }` byte pair (the DeferredWebSocketUpgrade below) +// instead of materializing a socket, so that infrastructure code can carry the tunnel across +// further hops before materializeUpgrade() turns it back into a real socket at the hop that +// serves it. +// +// The pair is BYTE-oriented (chunks are Uint8Arrays): the boundaries the pair exists to cross +// -- in particular native Cloudflare Workers RPC between isolates -- proxy only byte streams, +// refusing streams of arbitrary values. So the tunnel's frames are wrapped in a +// length-prefixed encoding: 1 byte frame type (0 text, 1 binary, 2 close), 4 bytes big-endian +// payload length, then the payload (UTF-8 text, raw bytes, or `{"code","reason"}` JSON +// respectively). +// +// Although the encoding never appears on a Cap'n Web session's own wire, it IS a persistent +// cross-deployment format: the deferring worker and the worker that materializes are deployed +// independently and can run different versions of this library during a rollout. The header +// layout (1-byte type + 4-byte big-endian length) is therefore frozen, and evolution is +// append-only: new frame types may be added, and a decoder that meets an unknown type fails +// loudly (tearing the tunnel down) rather than desynchronizing. Apps must still treat the pair +// as opaque -- the encoding is a contract between two versions of this library, not an API. +// +// The pair inherits the tunnel's flow control end to end: stream acks fire as the pair is +// read, so an in-transit or unconsumed pair keeps the provider throttled to the flow-control +// window, and a deferring receiver pools at most a window's worth of chunks. Once +// materialized, the endpoint behaves like a non-deferred one: messages are dispatched as they +// arrive, and a slow final consumer is absorbed by the native socket's send buffer, as on a +// direct WebSocket. + +// The runtime default when a session doesn't set `deferUpgradeMaterialization` explicitly: +// defer on Cloudflare Workers, where a socket materialized at the session endpoint could never +// cross another RPC hop anyway (and the serving isolate rebuilds it with materializeUpgrade()); +// materialize elsewhere, where there are no internal RPC hops to worry about and the app wants +// a usable socket directly. +export function deferUpgradesByDefault(): boolean { + return typeof WebSocketPair !== "undefined"; +} + +const FRAME_TEXT = 0, FRAME_BINARY = 1, FRAME_CLOSE = 2; + +// Far above any message a supported transport delivers (Workers caps incoming WebSocket +// messages at 1 MiB; `ws` defaults to 100 MiB), this bounds what a corrupt length header can +// make the decoder buffer. Exceeding it tears the tunnel down. +const MAX_FRAME_LENGTH = 128 * 1024 * 1024; + +const textEncoder = new TextEncoder(); +const textDecoder = new TextDecoder(); + +function encodeFrame(chunk: unknown): Uint8Array { + let type: number, payload: Uint8Array; + if (isCloseRecord(chunk)) { + type = FRAME_CLOSE; + payload = textEncoder.encode( + JSON.stringify({ code: chunk.close.code, reason: chunk.close.reason })); + } else { + let data = toStringOrBytes(chunk); + if (typeof data === "string") { + type = FRAME_TEXT; + payload = textEncoder.encode(data); + } else { + type = FRAME_BINARY; + payload = data; + } + } + if (payload.length > MAX_FRAME_LENGTH) { + // Enforced on both sides: failing here attributes the error at the offending send instead + // of at the far end's decoder. + throw new TypeError(`WebSocket tunnel frame of ${payload.length} bytes exceeds the maximum.`); + } + let frame = new Uint8Array(5 + payload.length); + frame[0] = type; + new DataView(frame.buffer).setUint32(1, payload.length); + frame.set(payload, 5); + return frame; +} + +// Decodes frames from a byte stream that may be re-chunked at arbitrary boundaries in transit. +class FrameDecoder { + // Bytes of the frame in progress, kept as a chunk list with a running total so a frame + // delivered in many pieces is assembled once when complete, not re-copied on every push. + // Invariant: the buffered bytes always begin at a frame boundary. + #chunks: Uint8Array[] = []; + #size = 0; + + // A partial frame is buffered; end-of-stream in this state means truncation, not a clean end. + get hasPartialFrame(): boolean { + return this.#size > 0; + } + + push(bytes: Uint8Array): unknown[] { + if (bytes.length > 0) { + this.#chunks.push(bytes); + this.#size += bytes.length; + } + + let frames: unknown[] = []; + while (this.#size >= 5) { + let header = this.#peek(5); + let length = new DataView(header.buffer, header.byteOffset).getUint32(1); + if (length > MAX_FRAME_LENGTH) { + throw new TypeError(`WebSocket tunnel frame of ${length} bytes exceeds the maximum.`); + } + if (this.#size < 5 + length) break; + + let frame = this.#take(5 + length); + let payload = frame.subarray(5); + switch (frame[0]) { + case FRAME_TEXT: + frames.push(textDecoder.decode(payload)); + break; + case FRAME_BINARY: + // Copy: the payload may be a view into a caller's chunk. + frames.push(payload.slice()); + break; + case FRAME_CLOSE: { + let { code, reason } = JSON.parse(textDecoder.decode(payload)); + frames.push({ close: { + code: typeof code === "number" ? code : 1005, + reason: typeof reason === "string" ? reason : "", + } }); + break; + } + default: + throw new TypeError(`Unknown tunnel frame type: ${frame[0]}`); + } + } + return frames; + } + + // A view whose first n bytes are the first n buffered bytes, without consuming them. May be + // longer than n (the fast path returns the whole first chunk). Requires n <= #size. + #peek(n: number): Uint8Array { + if (this.#chunks[0].length >= n) return this.#chunks[0]; + let out = new Uint8Array(n); + let pos = 0; + for (let chunk of this.#chunks) { + let take = Math.min(n - pos, chunk.length); + out.set(chunk.subarray(0, take), pos); + pos += take; + if (pos === n) break; + } + return out; + } + + // Removes the first n buffered bytes and returns them contiguously. + #take(n: number): Uint8Array { + this.#size -= n; + let first = this.#chunks[0]; + if (first.length === n) { + return this.#chunks.shift()!; + } else if (first.length > n) { + this.#chunks[0] = first.subarray(n); + return first.subarray(0, n); + } + let out = new Uint8Array(n); + let pos = 0; + while (pos < n) { + let chunk = this.#chunks[0]; + let take = Math.min(n - pos, chunk.length); + out.set(chunk.subarray(0, take), pos); + pos += take; + if (take === chunk.length) this.#chunks.shift(); + else this.#chunks[0] = chunk.subarray(take); + } + return out; + } +} + +function coerceBytes(chunk: unknown): Uint8Array { + let bytes = toStringOrBytes(chunk); + if (typeof bytes === "string") { + throw new TypeError("Expected bytes on a deferred tunnel stream, got a string."); + } + return bytes; +} + +// A ReadableStream whose chunks are `getReader()`'s chunks mapped through `transform`, which +// may yield zero or more output chunks per input (a decoder can need more bytes; one byte chunk +// can hold several frames). Built with highWaterMark 0 and a lazy reader thunk so the inner +// stream is locked only when a consumer actually reads: an untouched deferred Response's inner +// streams stay payload-owned and are released on disposal, like an unclaimed TunneledWebSocket. +// (A TransformStream can't do this: pipeThrough locks the inner stream immediately, and the +// default highWaterMark of 1 would pull -- and lock -- at construction time.) `flush` runs at +// end-of-stream; it can throw to turn a truncated stream into an error instead of a clean end. +function transformReadable(getReader: () => ReadableStreamDefaultReader, + transform: (chunk: unknown) => unknown[], flush?: () => void): ReadableStream { + let reader: ReadableStreamDefaultReader | undefined; + return new ReadableStream({ + async pull(controller) { + reader ??= getReader(); + // Loop until we enqueue something (or end): a transform can consume a chunk without + // producing output (a decoder mid-frame), and the stream machinery only re-invokes + // pull() after an enqueue -- returning empty-handed would strand the pending read. + while (true) { + let { done, value } = await reader.read(); + if (done) { + flush?.(); + controller.close(); + return; + } + let chunks = transform(value); + if (chunks.length > 0) { + for (let chunk of chunks) { + controller.enqueue(chunk); + } + return; + } + } + }, + cancel(reason) { + reader ??= getReader(); + return reader.cancel(reason); + }, + }, { highWaterMark: 0 }); +} + +// The WritableStream counterpart: transforms each incoming chunk and forwards the results to +// the writer, acquired lazily on first use. A transform or flush error means the tunnel's +// framing is broken, so fail closed: abort the inner writable (tearing down the sender's +// socket) rather than leaving it half-open. +function transformWritable(getWriter: () => WritableStreamDefaultWriter, + transform: (chunk: unknown) => unknown[], flush?: () => void): WritableStream { + let writer: WritableStreamDefaultWriter | undefined; + return new WritableStream({ + async write(chunk) { + writer ??= getWriter(); + try { + for (let out of transform(chunk)) { + await writer.write(out); + } + } catch (err) { + writer.abort(err).catch(() => {}); + throw err; + } + }, + close() { + writer ??= getWriter(); + try { + flush?.(); + } catch (err) { + writer.abort(err).catch(() => {}); + throw err; + } + return writer.close(); + }, + abort(reason) { + writer ??= getWriter(); + return writer.abort(reason); + }, + }); +} + +// The opaque pair a deferring session delivers as `Response.webSocket`. `init` carries the +// provider's upgrade ResponseInit (notably negotiated headers such as Sec-WebSocket-Protocol); +// it is plain data, so infrastructure that forwards the whole pair object across a hop +// preserves the headers for materializeUpgrade() automatically. +export interface DeferredWebSocketUpgrade { + readable: ReadableStream; + writable: WritableStream; + init?: ResponseInit; +} + +// Builds the Response a deferring session delivers: the tunnel's value streams wrapped into the +// opaque framed-byte pair. The readable encodes the tunnel's chunks into framed bytes; the +// writable decodes framed bytes and forwards the frames to the sender's WritableStream hook. +// Both inner ends stay owned by the containing payload, exactly like a ReadableStream and +// WritableStream received over RPC as plain values, and are only locked when the pair is used. +export function makeDeferredUpgradeResponse( + readable: ReadableStream, writableHook: StubHook, init: ResponseInit): Response { + let decoder = new FrameDecoder(); + let webSocket: DeferredWebSocketUpgrade = { + readable: transformReadable(() => readable.getReader(), chunk => [encodeFrame(chunk)]), + writable: transformWritable( + () => streamImpl.createWritableStreamFromHook(writableHook).getWriter(), + bytes => decoder.push(coerceBytes(bytes)), + () => { + if (decoder.hasPartialFrame) { + throw new Error("WebSocket tunnel closed with a truncated frame."); + } + }), + init, + }; + let response = new Response(null, init); + Object.defineProperty(response, "webSocket", { value: webSocket, configurable: true }); + return response; +} + +// Turns a tunneled socket's deferred byte pair -- as delivered by a deferring session -- back +// into a real upgrade Response. This is the +// counterpart the final hop calls: infrastructure code forwards the pair across boundaries +// that can serialize byte streams but not sockets (e.g. Cloudflare Workers RPC between +// isolates), then materializes the socket exactly once, where the 101 is actually returned to +// the client. +// +// The provider's upgrade headers ride on `webSocket.init`; an explicit `init` argument +// replaces it wholesale (merge the provider's headers in yourself if you want both). Reserved +// handshake headers that ride along (Connection, Upgrade, Sec-WebSocket-Accept/-Key/-Version, +// Content-Length, ...) are recomputed or dropped by the runtime when the Response completes a +// real upgrade, so only non-reserved headers such as Sec-WebSocket-Protocol reach the wire; +// callers on in-process hops (service bindings) see init's headers verbatim. +// +// The pair is consumed: both streams are locked synchronously, so a second +// materializeUpgrade() of the same pair throws instead of corrupting the first. On Workers the +// result is pumped by in-memory listeners, so a live materialized tunnel keeps its isolate (or +// Durable Object) resident for the socket's lifetime. +export function materializeUpgrade( + webSocket: DeferredWebSocketUpgrade, init?: ResponseInit): Response { + let reader = webSocket.readable.getReader(); + let writer: WritableStreamDefaultWriter; + try { + writer = webSocket.writable.getWriter(); + } catch (err) { + reader.releaseLock(); + throw err; + } + + // The inverse wrappers of makeDeferredUpgradeResponse's, with the write direction wrapped in + // a local WritableStream hook so the socket manages it exactly as it manages a hook received + // over RPC. + let decoder = new FrameDecoder(); + let readable = transformReadable(() => reader, + bytes => decoder.push(coerceBytes(bytes)), + () => { + if (decoder.hasPartialFrame) { + throw new Error("WebSocket tunnel ended with a truncated frame."); + } + }); + let writableHook = streamImpl.createWritableStreamHook( + transformWritable(() => writer, chunk => [encodeFrame(chunk)])); + return makeUpgradeResponse(readable, writableHook, init ?? webSocket.init ?? {}); +} + // Forward messages and closure between a native WebSocket (one end of a WebSocketPair) and a // tunneled socket, in both directions. //