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. //