From 5720a2736be1a191a3524d02e78dacc9a623f839 Mon Sep 17 00:00:00 2001 From: Henry Jang Date: Thu, 16 Jul 2026 13:38:48 +0900 Subject: [PATCH 1/7] feat(miniflare): tunnel raw TCP connect() through the remote-bindings proxy client Remote bindings were wired locally as a proxy worker that only relayed fetch/JSRPC, so `binding.connect(address)` against a raw-TCP binding (e.g. a VPC network) failed with "Incoming CONNECT on a worker not supported". Add an inbound `connect` handler to the proxy client worker: it recovers the target address from `socket.opened.localAddress` (the verbatim `connect()` authority, per workerd#6059), opens a dedicated WebSocket to the remote proxy server, and relays bytes in both directions via a new `pipeSocketOverWebSocket` helper (with `truncateCloseReason` for the 123-byte close-reason limit). The relay tears both directions down together so a parked read or a missing close event can't leak the socket. The inbound connect handler is experimental-gated in workerd, so opt the proxy client worker into the `experimental` compat flag only for raw-TCP bindings via a new `rawTcp` option on `remoteProxyClientWorker`; VPC networks pass it. The existing HTTP/JSRPC proxy call sites are unchanged. Co-Authored-By: Claude Fable 5 --- .../miniflare/src/plugins/shared/constants.ts | 17 +- .../src/plugins/vpc-networks/index.ts | 7 +- .../src/plugins/vpc-services/index.ts | 7 +- .../workers/shared/remote-bindings-utils.ts | 173 ++++++++++++++++++ .../shared/remote-proxy-client.worker.ts | 42 +++++ 5 files changed, 243 insertions(+), 3 deletions(-) diff --git a/packages/miniflare/src/plugins/shared/constants.ts b/packages/miniflare/src/plugins/shared/constants.ts index 02c5c2ef60..1f40e4eb59 100644 --- a/packages/miniflare/src/plugins/shared/constants.ts +++ b/packages/miniflare/src/plugins/shared/constants.ts @@ -90,9 +90,24 @@ export function objectEntryWorker( // into a per-binding service. The only static, non-props-able binding is the // loopback service (used to surface diagnostics back to the Miniflare host, // e.g. a Cloudflare Access block detected on the remote proxy response). -export function remoteProxyClientWorker(script?: () => string) { +// +// `options.rawTcp` opts the service into raw TCP `connect()` tunnelling (see the +// `experimental` compatibility flag below). That is a property of the service +// itself, not of any individual binding, so it stays valid under the shared, +// props-based service model. +export function remoteProxyClientWorker( + script?: () => string, + options?: { rawTcp?: boolean } +) { return { compatibilityDate: "2025-01-01", + // Raw TCP bindings (e.g. VPC networks) tunnel `binding.connect()` traffic + // through this worker's inbound `connect` handler, which requires the + // `experimental` compatibility flag (workerd#6059). Other bindings only + // proxy HTTP/JSRPC and must not opt in. + ...(options?.rawTcp + ? { compatibilityFlags: ["experimental"] } + : {}), modules: [ { name: "index.worker.js", diff --git a/packages/miniflare/src/plugins/vpc-networks/index.ts b/packages/miniflare/src/plugins/vpc-networks/index.ts index bc363f2bba..ff100f6e50 100644 --- a/packages/miniflare/src/plugins/vpc-networks/index.ts +++ b/packages/miniflare/src/plugins/vpc-networks/index.ts @@ -69,7 +69,12 @@ export const VPC_NETWORKS_PLUGIN: Plugin = { return [ { name: VPC_NETWORKS_REMOTE_SERVICE_NAME, - worker: remoteProxyClientWorker(), + // VPC networks expose raw TCP via `binding.connect()`, tunnelled + // through the proxy client's inbound `connect` handler. The shared + // `vpc-networks:remote` service is dedicated to VPC networks, so + // opting it into raw TCP leaves every other binding's service + // untouched. + worker: remoteProxyClientWorker(undefined, { rawTcp: true }), }, ]; }, diff --git a/packages/miniflare/src/plugins/vpc-services/index.ts b/packages/miniflare/src/plugins/vpc-services/index.ts index 597155d4f6..d70cef08b1 100644 --- a/packages/miniflare/src/plugins/vpc-services/index.ts +++ b/packages/miniflare/src/plugins/vpc-services/index.ts @@ -60,7 +60,12 @@ export const VPC_SERVICES_PLUGIN: Plugin = { return [ { name: VPC_SERVICES_REMOTE_SERVICE_NAME, - worker: remoteProxyClientWorker(), + // VPC services also expose raw TCP via `binding.connect()`, tunnelled + // through the proxy client's inbound `connect` handler. The shared + // `vpc-services:remote` service is dedicated to VPC services, so + // opting it into raw TCP leaves every other binding's service + // untouched. + worker: remoteProxyClientWorker(undefined, { rawTcp: true }), }, ]; }, diff --git a/packages/miniflare/src/workers/shared/remote-bindings-utils.ts b/packages/miniflare/src/workers/shared/remote-bindings-utils.ts index fcf19a22d8..55acdd9eab 100644 --- a/packages/miniflare/src/workers/shared/remote-bindings-utils.ts +++ b/packages/miniflare/src/workers/shared/remote-bindings-utils.ts @@ -174,6 +174,179 @@ export function makeFetch( }; } +/** + * Clamp a WebSocket close reason to the protocol's 123-byte limit. + * + * `WebSocket.close(code, reason)` requires the reason to be at most 123 UTF-8 + * bytes. Slicing by JavaScript string length (UTF-16 code units) can both + * overshoot the byte budget and split a multi-byte character, so truncate on a + * UTF-8 byte boundary instead. + */ +function truncateCloseReason(reason: string): string { + const bytes = new TextEncoder().encode(reason); + if (bytes.length <= 123) { + return reason; + } + // Back off to a UTF-8 character boundary: step left over any trailing + // continuation bytes (10xxxxxx) so we never cut a multi-byte sequence. + let end = 123; + while (end > 0 && ((bytes[end] ?? 0) & 0b1100_0000) === 0b1000_0000) { + end--; + } + return new TextDecoder().decode(bytes.subarray(0, end)); +} + +/** + * Relay raw TCP bytes between a `connect()` socket and a WebSocket, in both + * directions, propagating close and error state. + * + * Used to tunnel `binding.connect()` traffic through the remote-bindings proxy: + * the local proxy client pipes the caller's inbound socket over a WebSocket to + * the remote proxy server, which pipes it into the real binding's socket. + * + * The returned promise resolves once both directions have closed cleanly, and + * rejects if either side errors. Rejecting lets the caller's `connect()` handler + * surface the failure on its socket. + * + * Both directions are torn down together: when one side finishes — a WebSocket + * close/error, or the socket reaching EOF / erroring — the opposite direction is + * actively cancelled. The socket read is woken with `reader.cancel()` and the + * WebSocket wait is settled directly, so neither a parked `reader.read()` nor a + * never-arriving close event can leak the relay (and the underlying socket) for + * the lifetime of the process. + * + * Note: WebSockets have no backpressure signal, so `socket.readable` is drained + * as fast as it produces. Incoming WebSocket messages are written in order via a + * serialised write chain. There is no TCP half-close: when either direction ends + * the tunnel is torn down in full — the WebSocket is closed and the socket's + * writable side is closed — rather than closing a single direction. + */ +export async function pipeSocketOverWebSocket( + socket: Socket, + ws: WebSocket +): Promise { + const writer = socket.writable.getWriter(); + const reader = socket.readable.getReader(); + + let wsClosed = false; + function closeWebSocket(code: number, reason?: string) { + if (wsClosed) { + return; + } + wsClosed = true; + try { + ws.close( + code, + reason === undefined ? undefined : truncateCloseReason(reason) + ); + } catch { + // Already closing/closed. + } + } + + // ws -> socket writes are serialised through this chain to preserve byte + // order. The writable side is closed exactly once (guarded by `writerClosed`) + // so both teardown paths below can invoke it idempotently. + let writeChain = Promise.resolve(); + let writerClosed = false; + function closeWriter(): Promise { + if (writerClosed) { + return Promise.resolve(); + } + writerClosed = true; + return writeChain.then(() => writer.close()); + } + + // `fromWebSocket` settles from the ws close/error events, or is settled by the + // socket -> ws direction once it has closed the ws itself (in which case no + // inbound close event will arrive to settle it). + let resolveFromWs!: () => void; + let rejectFromWs!: (reason: unknown) => void; + const fromWebSocket = new Promise((resolve, reject) => { + resolveFromWs = resolve; + rejectFromWs = reject; + }); + + // WebSocket -> socket. Message events aren't awaited by the runtime, so writes + // are serialised through the promise chain to preserve byte order. + ws.addEventListener("message", (event) => { + const chunk = + typeof event.data === "string" + ? new TextEncoder().encode(event.data) + : new Uint8Array(event.data); + writeChain = writeChain + .then(() => writer.write(chunk)) + .catch((error) => { + // A socket write failed: tear the tunnel down rather than leaving the + // rejection unhandled. Close the ws (1011), cancel the opposite read + // direction, and reject so the caller sees the failure. + closeWebSocket( + 1011, + (error as Error)?.message ?? "socket write failed" + ); + reader.cancel().catch(() => {}); + rejectFromWs(error); + }); + }); + ws.addEventListener("close", (event) => { + wsClosed = true; + // Actively cancel the opposite direction so a parked `reader.read()` can't + // keep the socket -> ws relay (and this whole pipe) alive forever. + reader.cancel().catch(() => {}); + // Close code 1011 signals the remote end errored. + if (event.code === 1011) { + rejectFromWs( + new Error(event.reason || "Remote tunnel closed with an error") + ); + return; + } + // Flush any queued writes, then close the writable side so the caller sees + // EOF. + closeWriter().then(resolveFromWs, rejectFromWs); + }); + ws.addEventListener("error", () => { + wsClosed = true; + reader.cancel().catch(() => {}); + rejectFromWs(new Error("Tunnel WebSocket errored")); + }); + + // socket -> WebSocket. On EOF close cleanly (1000); on read error close with + // 1011 and reject so the caller sees the failure. + const toWebSocket = (async () => { + try { + for (;;) { + const { value, done } = await reader.read(); + if (done) { + break; + } + // Re-check after the parked read: the ws may have closed while we were + // waiting. If so, terminate quietly instead of erroring on `send`. + if (wsClosed) { + break; + } + ws.send( + value.buffer.slice( + value.byteOffset, + value.byteOffset + value.byteLength + ) + ); + } + closeWebSocket(1000); + } catch (error) { + closeWebSocket(1011, (error as Error).message); + throw error; + } finally { + reader.releaseLock(); + // We closed (or observed the close of) the ws on this side, so no inbound + // close event will arrive to settle `fromWebSocket`. Close the writable + // side and settle it here; idempotent with the ws `close` handler above. + closeWriter().then(resolveFromWs, rejectFromWs); + } + })(); + + await Promise.all([toWebSocket, fromWebSocket]); +} + /** * Create a remote proxy stub that proxies to a remote binding via capnweb. * diff --git a/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts b/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts index a4090913d8..24712cc854 100644 --- a/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts +++ b/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts @@ -3,6 +3,7 @@ import { SharedBindings } from "./constants"; import { makeFetch, makeRemoteProxyStub, + pipeSocketOverWebSocket, throwRemoteRequired, } from "./remote-bindings-utils"; import type { @@ -25,6 +26,47 @@ export default class Client extends WorkerEntrypoint< )(request); } + // Handles `binding.connect(address)` for raw TCP bindings (e.g. VPC networks). + // Only reachable when the worker is configured with the `experimental` + // compatibility flag, which enables inbound `connect` handlers (workerd#6059). + async connect(socket: Socket): Promise { + const { remoteProxyConnectionString, binding, cfTraceId } = + this.ctx.props; + if (!remoteProxyConnectionString) { + throwRemoteRequired(binding); + } + + // The address passed to `binding.connect("host:port")` arrives verbatim as + // the inbound socket's `localAddress` on the service-binding path. See + // https://github.com/cloudflare/workerd/pull/6059. + const { localAddress } = await socket.opened; + if (!localAddress) { + throw new Error( + `Binding ${binding} received a connection without a target address` + ); + } + + const headers = new Headers({ + Upgrade: "websocket", + "MF-Binding": binding, + "MF-Connect-Address": localAddress, + }); + if (cfTraceId) { + headers.set("cf-trace-id", cfTraceId); + } + + const response = await fetch(remoteProxyConnectionString, { headers }); + const ws = response.webSocket; + if (!ws) { + throw new Error( + `Binding ${binding} failed to open a tunnel to ${localAddress} (status ${response.status})` + ); + } + ws.accept(); + + await pipeSocketOverWebSocket(socket, ws); + } + constructor( ctx: ExecutionContext, env: RemoteBindingEnv From f658ae81b0c80ff6011449a5e11867051a317cec Mon Sep 17 00:00:00 2001 From: Henry Jang Date: Thu, 16 Jul 2026 13:38:56 +0900 Subject: [PATCH 2/7] feat(wrangler): add a raw TCP connect tunnel endpoint to ProxyServerWorker The remote-bindings proxy server now detects the WebSocket connect upgrade (carrying the target in `MF-Connect-Address`), opens the real binding socket via `fetcher.connect(address)`, and relays bytes bidirectionally over the WebSocket. This is the edge half of the Miniflare `connect()` tunnel. The relay helper is a byte-for-byte mirror of the one in miniflare's remote-bindings-utils; this template is bundled standalone for the edge, so the logic is duplicated here rather than imported. Co-Authored-By: Claude Fable 5 --- .../remoteBindings/ProxyServerWorker.ts | 197 +++++++++++++++++- 1 file changed, 196 insertions(+), 1 deletion(-) diff --git a/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts b/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts index 0ca9c5b8a6..4a0fed465e 100644 --- a/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts +++ b/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts @@ -113,10 +113,205 @@ function isJSRPCBinding(request: Request): boolean { return request.headers.has("Upgrade") && url.searchParams.has("MF-Binding"); } +/** + * A raw TCP tunnel request from the local proxy client's `connect` handler: a + * WebSocket upgrade carrying the target address in `MF-Connect-Address`. We open + * the binding's socket (routing through the VPC tunnel) and relay bytes between + * it and the WebSocket in both directions. + */ +function isConnectBinding(request: Request): boolean { + return ( + request.headers.get("Upgrade") === "websocket" && + request.headers.has("MF-Connect-Address") + ); +} + +function handleConnect(request: Request, env: Env): Response { + const address = request.headers.get("MF-Connect-Address")!; + const fetcher = getExposedFetcher(request, env) as Fetcher; + + const { 0: client, 1: server } = new WebSocketPair(); + server.accept(); + + const socket = fetcher.connect(address); + // Relay runs for the lifetime of the tunnel; failures are surfaced by closing + // the WebSocket with code 1011 (see pipeSocketOverWebSocket). + pipeSocketOverWebSocket(socket, server).catch(() => {}); + + return new Response(null, { status: 101, webSocket: client }); +} + +/** + * Clamp a WebSocket close reason to the protocol's 123-byte limit. + * + * `WebSocket.close(code, reason)` requires the reason to be at most 123 UTF-8 + * bytes. Slicing by JavaScript string length (UTF-16 code units) can both + * overshoot the byte budget and split a multi-byte character, so truncate on a + * UTF-8 byte boundary instead. + */ +function truncateCloseReason(reason: string): string { + const bytes = new TextEncoder().encode(reason); + if (bytes.length <= 123) { + return reason; + } + // Back off to a UTF-8 character boundary: step left over any trailing + // continuation bytes (10xxxxxx) so we never cut a multi-byte sequence. + let end = 123; + while (end > 0 && ((bytes[end] ?? 0) & 0b1100_0000) === 0b1000_0000) { + end--; + } + return new TextDecoder().decode(bytes.subarray(0, end)); +} + +/** + * Relay raw TCP bytes between a `connect()` socket and a WebSocket in both + * directions, propagating close and error state. Mirrors the helper in the + * Miniflare proxy client (packages/miniflare/src/workers/shared/remote-bindings-utils.ts); + * this file is bundled standalone for the edge, so the logic is duplicated. Keep + * the two implementations byte-for-byte identical (comments aside). + * + * Both directions are torn down together: when one side finishes — a WebSocket + * close/error, or the socket reaching EOF / erroring — the opposite direction is + * actively cancelled (`reader.cancel()` for the socket read, settling the + * promise for the WebSocket wait), so neither a parked read nor a never-arriving + * close event can leak the relay or the underlying socket. + * + * Note: WebSockets have no backpressure signal, and there is no TCP half-close — + * when either direction ends the tunnel is torn down in full (the WebSocket is + * closed and the socket's writable side is closed). + */ +async function pipeSocketOverWebSocket( + socket: Socket, + ws: WebSocket +): Promise { + const writer = socket.writable.getWriter(); + const reader = socket.readable.getReader(); + + let wsClosed = false; + function closeWebSocket(code: number, reason?: string) { + if (wsClosed) { + return; + } + wsClosed = true; + try { + ws.close( + code, + reason === undefined ? undefined : truncateCloseReason(reason) + ); + } catch { + // Already closing/closed. + } + } + + // ws -> socket writes are serialised through this chain to preserve byte + // order. The writable side is closed exactly once (guarded by `writerClosed`) + // so both teardown paths below can invoke it idempotently. + let writeChain = Promise.resolve(); + let writerClosed = false; + function closeWriter(): Promise { + if (writerClosed) { + return Promise.resolve(); + } + writerClosed = true; + return writeChain.then(() => writer.close()); + } + + // `fromWebSocket` settles from the ws close/error events, or is settled by the + // socket -> ws direction once it has closed the ws itself (in which case no + // inbound close event will arrive to settle it). + let resolveFromWs!: () => void; + let rejectFromWs!: (reason: unknown) => void; + const fromWebSocket = new Promise((resolve, reject) => { + resolveFromWs = resolve; + rejectFromWs = reject; + }); + + // WebSocket -> socket. Message events aren't awaited by the runtime, so writes + // are serialised through the promise chain to preserve byte order. + ws.addEventListener("message", (event) => { + const chunk = + typeof event.data === "string" + ? new TextEncoder().encode(event.data) + : new Uint8Array(event.data); + writeChain = writeChain + .then(() => writer.write(chunk)) + .catch((error) => { + // A socket write failed: tear the tunnel down rather than leaving the + // rejection unhandled. Close the ws (1011), cancel the opposite read + // direction, and reject so the caller sees the failure. + closeWebSocket( + 1011, + (error as Error)?.message ?? "socket write failed" + ); + reader.cancel().catch(() => {}); + rejectFromWs(error); + }); + }); + ws.addEventListener("close", (event) => { + wsClosed = true; + // Actively cancel the opposite direction so a parked `reader.read()` can't + // keep the socket -> ws relay (and this whole pipe) alive forever. + reader.cancel().catch(() => {}); + // Close code 1011 signals the remote end errored. + if (event.code === 1011) { + rejectFromWs( + new Error(event.reason || "Remote tunnel closed with an error") + ); + return; + } + // Flush any queued writes, then close the writable side so the caller sees + // EOF. + closeWriter().then(resolveFromWs, rejectFromWs); + }); + ws.addEventListener("error", () => { + wsClosed = true; + reader.cancel().catch(() => {}); + rejectFromWs(new Error("Tunnel WebSocket errored")); + }); + + // socket -> WebSocket. On EOF close cleanly (1000); on read error close with + // 1011 and reject so the caller sees the failure. + const toWebSocket = (async () => { + try { + for (;;) { + const { value, done } = await reader.read(); + if (done) { + break; + } + // Re-check after the parked read: the ws may have closed while we were + // waiting. If so, terminate quietly instead of erroring on `send`. + if (wsClosed) { + break; + } + ws.send( + value.buffer.slice( + value.byteOffset, + value.byteOffset + value.byteLength + ) + ); + } + closeWebSocket(1000); + } catch (error) { + closeWebSocket(1011, (error as Error).message); + throw error; + } finally { + reader.releaseLock(); + // We closed (or observed the close of) the ws on this side, so no inbound + // close event will arrive to settle `fromWebSocket`. Close the writable + // side and settle it here; idempotent with the ws `close` handler above. + closeWriter().then(resolveFromWs, rejectFromWs); + } + })(); + + await Promise.all([toWebSocket, fromWebSocket]); +} + export default { async fetch(request, env) { try { - if (isJSRPCBinding(request)) { + if (isConnectBinding(request)) { + return handleConnect(request, env); + } else if (isJSRPCBinding(request)) { return await newWorkersRpcResponse( request, getExposedJSRPCBinding(request, env) From 13813d25d3ebb447f841b8e122d2a67d0e6b31ae Mon Sep 17 00:00:00 2001 From: Henry Jang Date: Thu, 16 Jul 2026 13:39:05 +0900 Subject: [PATCH 3/7] test(miniflare): cover the raw TCP connect() tunnel Add 13 tests for the `connect()` tunnel: end-to-end byte relay and address propagation plus a close/error matrix run against the real bundled ProxyServerWorker (in a second "edge" Miniflare instance), unit coverage of `pipeSocketOverWebSocket` teardown/error paths, and opt-in assertions that the `experimental` flag is set only when `rawTcp` is requested. Add a `@relay-under-test` vitest alias so the worker-side relay helper can be unit-tested without pulling worker-typed source into the node-side tsconfig (which excludes `src/workers/**`). Co-Authored-By: Claude Fable 5 --- .../shared/remote-bindings-connect.spec.ts | 741 ++++++++++++++++++ packages/miniflare/vitest.config.mts | 9 + 2 files changed, 750 insertions(+) create mode 100644 packages/miniflare/test/plugins/shared/remote-bindings-connect.spec.ts diff --git a/packages/miniflare/test/plugins/shared/remote-bindings-connect.spec.ts b/packages/miniflare/test/plugins/shared/remote-bindings-connect.spec.ts new file mode 100644 index 0000000000..44fbe11dfa --- /dev/null +++ b/packages/miniflare/test/plugins/shared/remote-bindings-connect.spec.ts @@ -0,0 +1,741 @@ +import path from "node:path"; +// The relay helper under unit test. It is the same implementation bundled into +// miniflare's dist and mirrored (byte-for-byte, comments aside) into the edge +// ProxyServerWorker. It is resolved via a vitest alias (see vitest.config.mts) +// rather than a real path so that its worker-typed source isn't pulled into the +// node-side tsconfig, which excludes `src/workers/**`. tsc has no matching path +// mapping, hence the expected error below. +// @ts-expect-error resolved at runtime by a vitest alias; see vitest.config.mts +import { pipeSocketOverWebSocket } from "@relay-under-test"; +import { build } from "esbuild"; +import { + Miniflare, + remoteProxyClientWorker, + VPC_SERVICES_PLUGIN, +} from "miniflare"; +import { beforeAll, describe, test } from "vitest"; +import { useDispose } from "../../test-shared"; +import type { RemoteProxyConnectionString } from "miniflare"; + +// The raw-TCP `connect()` tunnel spans two workers: +// - the Miniflare-side remote-proxy client (built into miniflare's dist), whose +// inbound `connect` handler tunnels bytes over a WebSocket, and +// - the edge-side ProxyServerWorker (from the @cloudflare/remote-bindings +// package), which relays them into the real binding via +// `env[binding].connect(address)`. +// +// To exercise the real ProxyServerWorker we bundle it (it imports capnweb) and +// run it in a second "edge" Miniflare instance. A `vpc-target` worker with a +// `connect` handler stands in for a private VPC service: it reflects the address +// it observed (proving verbatim propagation) and echoes bytes back. +const PROXY_SERVER_WORKER_PATH = path.resolve( + __dirname, + "../../../../remote-bindings/templates/remoteBindings/ProxyServerWorker.ts" +); + +const COMPAT_DATE = "2025-01-01"; + +// Shared bundle of the real edge ProxyServerWorker, built once for the file. +let proxyServerBundle: string; +beforeAll(async () => { + const result = await build({ + entryPoints: [PROXY_SERVER_WORKER_PATH], + bundle: true, + format: "esm", + target: "esnext", + platform: "browser", + // Workers runtime built-ins are provided by workerd, not bundled. + external: ["cloudflare:*"], + write: false, + }); + proxyServerBundle = result.outputFiles[0].text; +}); + +// Reflects `socket.opened.localAddress` (the verbatim `connect()` address), then +// echoes received bytes, closing once it sees a newline sentinel. +const VPC_TARGET_SCRIPT = /* javascript */ ` + export default { + async connect(socket) { + const { localAddress } = await socket.opened; + const writer = socket.writable.getWriter(); + const enc = new TextEncoder(); + const dec = new TextDecoder(); + await writer.write(enc.encode("ADDR:" + localAddress + "|")); + const reader = socket.readable.getReader(); + let buf = ""; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + buf += dec.decode(value, { stream: true }); + await writer.write(value); + if (buf.includes("\\n")) break; + } + await writer.close(); + }, + }; +`; + +// Instrumented VPC target used by the close/error matrix. It selects its +// behaviour from the verbatim `connect()` address and records the lifecycle of +// each connection in module globals, exposed over `/state`. Because the relay's +// internal teardown isn't directly observable, `connectDone` (incremented in the +// handler's `finally`) is our proxy signal that the relay actually terminated +// this side of the tunnel — if a direction leaks, the handler parks forever on +// `reader.read()` and `connectDone` never advances. +const INSTRUMENTED_TARGET_SCRIPT = /* javascript */ ` + let connectStarted = 0; + let connectDone = 0; + let lastOutcome = ""; + export default { + async fetch(request) { + const url = new URL(request.url); + if (url.pathname === "/state") { + return Response.json({ connectStarted, connectDone, lastOutcome }); + } + return new Response("vpc-target"); + }, + async connect(socket) { + connectStarted++; + let outcome = "opening"; + try { + const { localAddress } = await socket.opened; + const enc = new TextEncoder(); + const writer = socket.writable.getWriter(); + if (localAddress.startsWith("error")) { + // (a) Handler errors immediately: aborting the writable errors the + // socket so the far (user) end sees a failure rather than a clean EOF. + outcome = "aborted"; + await writer.abort(new Error("vpc target refused connection")); + return; + } + if (localAddress.startsWith("eof")) { + // (b) Send some data then EOF (close our writable). Then read until + // our readable reaches EOF, which only happens once the relay has + // torn down the opposite direction. + await writer.write(enc.encode("SERVER-DATA\\n")); + await writer.close(); + const reader = socket.readable.getReader(); + outcome = "sent-eof-waiting-readable-eof"; + for (;;) { + const { done } = await reader.read(); + if (done) break; + } + outcome = "eof-both-sides"; + return; + } + // (c) Echo until our readable reaches EOF, i.e. until the user closes + // first and the relay propagates that close to us. + outcome = "echo-waiting-user-close"; + const reader = socket.readable.getReader(); + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + await writer.write(value); + } + await writer.close().catch(() => {}); + outcome = "user-closed"; + } catch (e) { + outcome = "threw:" + (e && e.message ? e.message : String(e)); + } finally { + connectDone++; + lastOutcome = outcome; + } + }, + }; +`; + +// User worker: opens the tunnel via the VPC binding, sends a payload and reads +// the reflected address + echo back. +const USER_SCRIPT = /* javascript */ ` + export default { + async fetch(request, env) { + const socket = env.VPC.connect("db.internal:3306"); + await socket.opened; + const writer = socket.writable.getWriter(); + await writer.write(new TextEncoder().encode("PING\\n")); + writer.releaseLock(); + const reader = socket.readable.getReader(); + const dec = new TextDecoder(); + let out = ""; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + out += dec.decode(value, { stream: true }); + if (out.includes("PING\\n")) break; + } + return Response.json({ received: out }); + }, + }; +`; + +// Matrix user worker. Behaviour is chosen from the `?scenario=` query so a +// single worker can drive every close/error case. +const MATRIX_USER_SCRIPT = /* javascript */ ` + export default { + async fetch(request, env) { + const scenario = new URL(request.url).searchParams.get("scenario"); + const dec = new TextDecoder(); + + if (scenario === "error") { + // (a) Expect the tunnel to surface the target's failure. + const socket = env.VPC.connect("error.internal:1"); + try { + await socket.opened; + const reader = socket.readable.getReader(); + for (;;) { + const { done } = await reader.read(); + if (done) break; + } + // A clean EOF (no error) also counts as a resolved read loop. + return Response.json({ errored: false }); + } catch (e) { + return Response.json({ + errored: true, + message: e && e.message ? e.message : String(e), + }); + } + } + + if (scenario === "eof") { + // (b) Target sends data then closes; we read to EOF and return without + // closing our own send side. The relay must still tear the tunnel down. + const socket = env.VPC.connect("eof.internal:3306"); + await socket.opened; + const reader = socket.readable.getReader(); + let out = ""; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + out += dec.decode(value, { stream: true }); + } + return Response.json({ received: out, sawEof: true }); + } + + if (scenario === "userclose") { + // (c) We write, confirm the echo, then close first. + const socket = env.VPC.connect("echo.internal:3306"); + await socket.opened; + const writer = socket.writable.getWriter(); + await writer.write(new TextEncoder().encode("HELLO\\n")); + const reader = socket.readable.getReader(); + let out = ""; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + out += dec.decode(value, { stream: true }); + if (out.includes("HELLO\\n")) break; + } + reader.releaseLock(); + writer.releaseLock(); + await socket.close(); + return Response.json({ received: out, closed: true }); + } + + return new Response("unknown scenario", { status: 400 }); + }, + }; +`; + +function makeEdge(targetScript: string, directSockets = false): Miniflare { + return new Miniflare({ + workers: [ + { + name: "proxy-server", + compatibilityDate: COMPAT_DATE, + compatibilityFlags: ["experimental"], + modules: [ + { + type: "ESModule", + path: "ProxyServerWorker.js", + contents: proxyServerBundle, + }, + ], + serviceBindings: { VPC: "vpc-target" }, + }, + { + name: "vpc-target", + compatibilityDate: COMPAT_DATE, + compatibilityFlags: ["experimental"], + modules: true, + script: targetScript, + // Enable direct access so the test can read the target's `/state`. + ...(directSockets + ? { unsafeDirectSockets: [{ host: "127.0.0.1" }] } + : {}), + }, + ], + }); +} + +function makeLocal(userScript: string, edgeUrl: URL): Miniflare { + return new Miniflare({ + modules: true, + compatibilityDate: COMPAT_DATE, + script: userScript, + vpcNetworks: { + VPC: { + network_id: "test-network", + remoteProxyConnectionString: + edgeUrl as unknown as RemoteProxyConnectionString, + }, + }, + }); +} + +// Same as `makeLocal`, but wires the user worker's `VPC` binding through a VPC +// *service* binding rather than a VPC network. VPC services also support +// `connect()`, so the raw-TCP tunnel must work identically. +function makeLocalService(userScript: string, edgeUrl: URL): Miniflare { + return new Miniflare({ + modules: true, + compatibilityDate: COMPAT_DATE, + script: userScript, + vpcServices: { + VPC: { + service_id: "test-service", + remoteProxyConnectionString: + edgeUrl as unknown as RemoteProxyConnectionString, + }, + }, + }); +} + +type TargetState = { + connectStarted: number; + connectDone: number; + lastOutcome: string; +}; + +async function readTargetState(edge: Miniflare): Promise { + const url = await edge.unsafeGetDirectURL("vpc-target"); + const res = await fetch(new URL("/state", url)); + return (await res.json()) as TargetState; +} + +// Poll the instrumented target until its connect handler has finished (the relay +// tore this side down), or give up after ~2s and return the last state seen. +async function waitForConnectDone( + edge: Miniflare, + expected = 1 +): Promise { + let state = await readTargetState(edge); + for (let i = 0; i < 100 && state.connectDone < expected; i++) { + await new Promise((resolve) => setTimeout(resolve, 20)); + state = await readTargetState(edge); + } + return state; +} + +describe("remote-bindings proxy client: raw TCP connect tunnel", () => { + test("tunnels bytes end-to-end and forwards the connect address verbatim", async ({ + expect, + }) => { + // "Edge" instance running the real ProxyServerWorker, with a VPC binding + // pointing at a worker that has a `connect` handler. + const edge = makeEdge(VPC_TARGET_SCRIPT); + useDispose(edge); + const edgeUrl = await edge.ready; + + // "Local" instance with the VPC networks binding, routed at the edge above. + const local = makeLocal(USER_SCRIPT, edgeUrl); + useDispose(local); + + const response = await local.dispatchFetch("http://localhost/"); + const { received } = (await response.json()) as { received: string }; + + // The address the user passed to `binding.connect(...)` must arrive + // verbatim at the far end of the tunnel. + expect(received).toContain("ADDR:db.internal:3306|"); + // And bytes must round-trip through the tunnel. + expect(received).toContain("PING\n"); + }); + + // Same round-trip, but the local binding is a VPC *service* rather than a VPC + // network. Both expose `connect()`, so the raw-TCP tunnel must behave + // identically through the shared `vpc-services:remote` proxy client. + test("tunnels bytes end-to-end through a VPC service binding", async ({ + expect, + }) => { + const edge = makeEdge(VPC_TARGET_SCRIPT); + useDispose(edge); + const edgeUrl = await edge.ready; + + const local = makeLocalService(USER_SCRIPT, edgeUrl); + useDispose(local); + + const response = await local.dispatchFetch("http://localhost/"); + const { received } = (await response.json()) as { received: string }; + + expect(received).toContain("ADDR:db.internal:3306|"); + expect(received).toContain("PING\n"); + }); +}); + +describe("remote-bindings proxy client: connect tunnel close/error matrix", () => { + // (a) The target's `connect` handler errors immediately (aborts its socket). + // + // NOTE ON A HARNESS LIMITATION: Miniflare's service-binding socket layer + // collapses a mid-tunnel error into a clean EOF, so the user's `opened`/`read` + // does NOT reject here — it resolves and then sees EOF. The relay's actual + // reject / 1011 behaviour (which is what surfaces a real VPC failure on a live + // socket) is instead covered deterministically by the + // "pipeSocketOverWebSocket: teardown & error propagation" unit tests below. + // What we *can* assert end-to-end is that a target-side error still terminates + // the tunnel cleanly, the user request completes (no hang), and the target's + // handler runs to completion. + test("a target connect failure terminates the tunnel without hanging", async ({ + expect, + }) => { + const edge = makeEdge(INSTRUMENTED_TARGET_SCRIPT, true); + useDispose(edge); + const edgeUrl = await edge.ready; + + const local = makeLocal(MATRIX_USER_SCRIPT, edgeUrl); + useDispose(local); + + const response = await local.dispatchFetch( + "http://localhost/?scenario=error" + ); + // The user request completes rather than hanging. + const result = (await response.json()) as { + errored: boolean; + message?: string; + }; + expect(result).toBeTruthy(); + + // The target handler ran and terminated (its `finally` fired). + const state = await waitForConnectDone(edge); + expect(state.connectStarted).toBe(1); + expect(state.connectDone).toBe(1); + expect(state.lastOutcome).toBe("aborted"); + }); + + // (b) The target sends data then closes (EOF). The user reads to EOF and + // returns WITHOUT closing its own send side. This is the HIGH-1 regression: + // without cross-direction cancellation the relay would park on the opposite + // read forever and the target handler would never finish. + test("tears the relay down when the target half-closes with the user still open", async ({ + expect, + }) => { + const edge = makeEdge(INSTRUMENTED_TARGET_SCRIPT, true); + useDispose(edge); + const edgeUrl = await edge.ready; + + const local = makeLocal(MATRIX_USER_SCRIPT, edgeUrl); + useDispose(local); + + const response = await local.dispatchFetch( + "http://localhost/?scenario=eof" + ); + const result = (await response.json()) as { + received: string; + sawEof: boolean; + }; + + // The user saw the target's data followed by EOF. + expect(result.received).toContain("SERVER-DATA"); + expect(result.sawEof).toBe(true); + + // The relay must have closed the target's readable side, letting its + // handler complete. If HIGH-1 regressed, `connectDone` stays 0 here. + const state = await waitForConnectDone(edge); + expect(state.connectDone).toBe(1); + expect(state.lastOutcome).toBe("eof-both-sides"); + }); + + // (c) The user closes first. The relay must propagate that close so the + // target's `connect` handler observes EOF and its pipe is cleaned up. + test("propagates a user-initiated close to the target pipe", async ({ + expect, + }) => { + const edge = makeEdge(INSTRUMENTED_TARGET_SCRIPT, true); + useDispose(edge); + const edgeUrl = await edge.ready; + + const local = makeLocal(MATRIX_USER_SCRIPT, edgeUrl); + useDispose(local); + + const response = await local.dispatchFetch( + "http://localhost/?scenario=userclose" + ); + const result = (await response.json()) as { + received: string; + closed: boolean; + }; + + // The echo round-tripped before the user closed. + expect(result.received).toContain("HELLO"); + expect(result.closed).toBe(true); + + // The user-initiated close must reach the target as EOF. + const state = await waitForConnectDone(edge); + expect(state.connectDone).toBe(1); + expect(state.lastOutcome).toBe("user-closed"); + }); +}); + +// --- Focused unit tests for the relay helper ------------------------------ +// +// The e2e matrix above proves the tunnel works end-to-end, but Miniflare's +// service-binding socket layer collapses a mid-tunnel *error* into a clean EOF +// (see the (a) test's note), so it can't exercise the relay's reject / 1011 +// paths, nor prove that a parked read is actively cancelled. These unit tests +// drive `pipeSocketOverWebSocket` directly with controllable fakes to cover +// exactly those paths deterministically. + +type FakeSocket = { + socket: { + readable: ReadableStream; + writable: WritableStream; + }; + writes: Uint8Array[]; + writableClosed: () => boolean; + cancelledWith: () => unknown; + pushRead: (bytes: Uint8Array) => void; + endRead: () => void; + errorRead: (err: unknown) => void; +}; + +function makeFakeSocket(writeError?: unknown): FakeSocket { + const writes: Uint8Array[] = []; + let writableClosed = false; + const writable = new WritableStream({ + write(chunk) { + if (writeError) { + throw writeError; + } + writes.push(chunk); + }, + close() { + writableClosed = true; + }, + abort() { + writableClosed = true; + }, + }); + + let controller: ReadableStreamDefaultController; + let cancelled: unknown = undefined; + let didCancel = false; + const readable = new ReadableStream({ + start(c) { + controller = c; + }, + cancel(reason) { + didCancel = true; + cancelled = reason; + }, + }); + + return { + socket: { readable, writable }, + writes, + writableClosed: () => writableClosed, + cancelledWith: () => (didCancel ? { reason: cancelled } : undefined), + pushRead: (bytes) => controller.enqueue(bytes), + endRead: () => controller.close(), + errorRead: (err) => controller.error(err), + }; +} + +type FakeWebSocket = { + ws: { + addEventListener: (type: string, fn: (event: unknown) => void) => void; + send: (data: ArrayBuffer) => void; + close: (code?: number, reason?: string) => void; + }; + sent: ArrayBuffer[]; + closedWith: () => { code?: number; reason?: string } | undefined; + emitMessage: (data: ArrayBuffer | string) => void; + emitClose: (code: number, reason?: string) => void; + emitError: () => void; +}; + +function makeFakeWebSocket(): FakeWebSocket { + const listeners: Record void)[]> = { + message: [], + close: [], + error: [], + }; + const sent: ArrayBuffer[] = []; + let closed: { code?: number; reason?: string } | undefined; + const ws = { + addEventListener(type: string, fn: (event: unknown) => void) { + (listeners[type] ??= []).push(fn); + }, + send(data: ArrayBuffer) { + sent.push(data); + }, + close(code?: number, reason?: string) { + closed ??= { code, reason }; + }, + }; + + return { + ws, + sent, + closedWith: () => closed, + emitMessage: (data) => listeners.message.forEach((fn) => fn({ data })), + emitClose: (code, reason) => + listeners.close.forEach((fn) => fn({ code, reason })), + emitError: () => listeners.error.forEach((fn) => fn({})), + }; +} + +// Let a parked read (and any microtasks) settle before asserting. +const flush = () => new Promise((resolve) => setTimeout(resolve, 0)); + +describe("pipeSocketOverWebSocket: teardown & error propagation", () => { + test("[HIGH-1] a clean WS close cancels the parked socket read and resolves", async ({ + expect, + }) => { + const socket = makeFakeSocket(); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + // The socket->ws direction is parked on `reader.read()` (no data). A clean + // remote close must cancel it so the pipe resolves instead of hanging. + await flush(); + ws.emitClose(1000, ""); + + await expect(pipe).resolves.toBeUndefined(); + expect(socket.cancelledWith()).not.toBeUndefined(); + expect(socket.writableClosed()).toBe(true); + }); + + test("[HIGH-1] socket EOF closes the WS and settles the WS-wait direction", async ({ + expect, + }) => { + const socket = makeFakeSocket(); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + socket.pushRead(new Uint8Array([1, 2, 3])); + socket.endRead(); + + // The pipe must resolve even though no inbound WS `close` event arrives — + // the socket->ws direction settles the WS-wait side itself. + await expect(pipe).resolves.toBeUndefined(); + expect(ws.sent.length).toBe(1); + expect(ws.closedWith()).toEqual({ code: 1000, reason: undefined }); + }); + + test("a remote 1011 close rejects and cancels the opposite direction", async ({ + expect, + }) => { + const socket = makeFakeSocket(); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + await flush(); + ws.emitClose(1011, "remote boom"); + + await expect(pipe).rejects.toThrow("remote boom"); + expect(socket.cancelledWith()).not.toBeUndefined(); + }); + + test("[LOW-A] a WS error event rejects, cancels the opposite direction, and closes the writer", async ({ + expect, + }) => { + const socket = makeFakeSocket(); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + // The socket->ws direction is parked on `reader.read()`. A WS `error` event + // must reject the pipe, cancel that parked read, and (via the woken + // `toWebSocket` finally) close the writable side. + await flush(); + ws.emitError(); + + await expect(pipe).rejects.toThrow("Tunnel WebSocket errored"); + expect(socket.cancelledWith()).not.toBeUndefined(); + // The writer is closed asynchronously by the woken `toWebSocket` finally. + await flush(); + expect(socket.writableClosed()).toBe(true); + }); + + test("a socket read error closes the WS with 1011 and rejects", async ({ + expect, + }) => { + const socket = makeFakeSocket(); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + socket.errorRead(new Error("read failed")); + + await expect(pipe).rejects.toThrow("read failed"); + expect(ws.closedWith()?.code).toBe(1011); + }); + + test("[LOW-1] a socket write failure tears down with 1011 and rejects", async ({ + expect, + }) => { + const socket = makeFakeSocket(new Error("write failed")); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + ws.emitMessage(new Uint8Array([9]).buffer); + + await expect(pipe).rejects.toThrow("write failed"); + expect(ws.closedWith()?.code).toBe(1011); + expect(socket.cancelledWith()).not.toBeUndefined(); + }); + + test("[LOW-3] the 1011 close reason is truncated to <=123 UTF-8 bytes on a char boundary", async ({ + expect, + }) => { + const socket = makeFakeSocket(); + const ws = makeFakeWebSocket(); + const pipe = pipeSocketOverWebSocket(socket.socket, ws.ws); + + // 200 x "é" = 400 UTF-8 bytes, well over the 123-byte cap and multi-byte so + // a naive UTF-16 slice could split a character. + socket.errorRead(new Error("é".repeat(200))); + + await expect(pipe).rejects.toThrow(); + const reason = ws.closedWith()?.reason ?? ""; + expect(new TextEncoder().encode(reason).length).toBeLessThanOrEqual(123); + // No replacement character => no multi-byte sequence was split. + expect(reason).not.toContain("�"); + }); +}); + +describe("remoteProxyClientWorker: raw TCP opt-in", () => { + test("does not enable the experimental flag by default", ({ expect }) => { + const worker = remoteProxyClientWorker(); + expect( + (worker as { compatibilityFlags?: string[] }).compatibilityFlags + ).toBeUndefined(); + }); + + test("enables the experimental flag when rawTcp is set", ({ expect }) => { + const worker = remoteProxyClientWorker(undefined, { + rawTcp: true, + }); + expect( + (worker as { compatibilityFlags?: string[] }).compatibilityFlags + ).toEqual(["experimental"]); + }); +}); + +describe("VPC_SERVICES plugin: raw TCP opt-in", () => { + test("opts its remote-proxy client service into the experimental flag", async ({ + expect, + }) => { + const services = await VPC_SERVICES_PLUGIN.getServices({ + options: { + vpcServices: { + VPC: { service_id: "test-service" }, + }, + }, + // The plugin's `getServices` only reads `options`; the remaining fields + // of the services context are irrelevant to the raw-TCP opt-in. + } as unknown as Parameters[0]); + + const serviceList = services as Array<{ + worker?: { compatibilityFlags?: string[] }; + }>; + expect(serviceList).toHaveLength(1); + expect(serviceList[0].worker?.compatibilityFlags).toEqual(["experimental"]); + }); +}); diff --git a/packages/miniflare/vitest.config.mts b/packages/miniflare/vitest.config.mts index ffa4d5d92a..15b10f6ec5 100644 --- a/packages/miniflare/vitest.config.mts +++ b/packages/miniflare/vitest.config.mts @@ -18,6 +18,15 @@ export default defineConfig({ resolve: { alias: { miniflare: path.resolve(__dirname, "dist/src/index.js"), + // Exposes the worker-side raw-TCP relay helper to focused unit tests + // without importing it by a real path (which would drag worker-typed + // source into the node-side tsconfig, whose `exclude` covers + // `src/workers/**`). The spec imports this id with an `@ts-expect-error` + // since tsc has no matching path mapping. + "@relay-under-test": path.resolve( + __dirname, + "src/workers/shared/remote-bindings-utils.ts" + ), "miniflare:shared": path.resolve( __dirname, "src/workers/shared/index.ts" From 919708fd3522b3a9bd2c554ed650c3a394227b0e Mon Sep 17 00:00:00 2001 From: Henry Jang Date: Thu, 16 Jul 2026 13:39:11 +0900 Subject: [PATCH 4/7] chore: add changeset for remote VPC Network connect() support Co-Authored-By: Claude Fable 5 --- .changeset/vpc-networks-connect-tunnel.md | 10 ++++++++++ 1 file changed, 10 insertions(+) create mode 100644 .changeset/vpc-networks-connect-tunnel.md diff --git a/.changeset/vpc-networks-connect-tunnel.md b/.changeset/vpc-networks-connect-tunnel.md new file mode 100644 index 0000000000..fcb39cbba2 --- /dev/null +++ b/.changeset/vpc-networks-connect-tunnel.md @@ -0,0 +1,10 @@ +--- +"miniflare": minor +"wrangler": minor +--- + +Support `connect()` on remote VPC Network and VPC Service bindings in local development + +Remote VPC Network and VPC Service bindings previously only supported HTTP and JSRPC, so calling `binding.connect(address)` against a private TCP service (for example a database) failed in local dev with `Incoming CONNECT on a worker not supported`. Raw TCP connections through remote VPC Network and VPC Service bindings now work in local development. + +This feature is experimental. Existing HTTP and JSRPC usage of remote VPC Network and VPC Service bindings is unaffected, and no new configuration is required. From 9af6eb18671ae687127a238de3a1e26582ded655 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 24 Jul 2026 11:37:47 +0000 Subject: [PATCH 5/7] style: fix oxfmt formatting in miniflare remote-bindings files Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01BuKCFc94BFDpbwag1UfPbB --- packages/miniflare/src/plugins/shared/constants.ts | 4 +--- .../src/workers/shared/remote-proxy-client.worker.ts | 3 +-- 2 files changed, 2 insertions(+), 5 deletions(-) diff --git a/packages/miniflare/src/plugins/shared/constants.ts b/packages/miniflare/src/plugins/shared/constants.ts index 1f40e4eb59..6d64fc61f8 100644 --- a/packages/miniflare/src/plugins/shared/constants.ts +++ b/packages/miniflare/src/plugins/shared/constants.ts @@ -105,9 +105,7 @@ export function remoteProxyClientWorker( // through this worker's inbound `connect` handler, which requires the // `experimental` compatibility flag (workerd#6059). Other bindings only // proxy HTTP/JSRPC and must not opt in. - ...(options?.rawTcp - ? { compatibilityFlags: ["experimental"] } - : {}), + ...(options?.rawTcp ? { compatibilityFlags: ["experimental"] } : {}), modules: [ { name: "index.worker.js", diff --git a/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts b/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts index 24712cc854..0de0e3a089 100644 --- a/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts +++ b/packages/miniflare/src/workers/shared/remote-proxy-client.worker.ts @@ -30,8 +30,7 @@ export default class Client extends WorkerEntrypoint< // Only reachable when the worker is configured with the `experimental` // compatibility flag, which enables inbound `connect` handlers (workerd#6059). async connect(socket: Socket): Promise { - const { remoteProxyConnectionString, binding, cfTraceId } = - this.ctx.props; + const { remoteProxyConnectionString, binding, cfTraceId } = this.ctx.props; if (!remoteProxyConnectionString) { throwRemoteRequired(binding); } From bb6480e918aefebf3cee039c6faa99ee54d4379c Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 24 Jul 2026 11:42:03 +0000 Subject: [PATCH 6/7] refactor: replace non-null assertion with a 400 guard in handleConnect oxlint (typescript-eslint/no-non-null-assertion) forbids `!`; return a 400 for a connect upgrade missing the MF-Connect-Address header instead. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01BuKCFc94BFDpbwag1UfPbB --- .../templates/remoteBindings/ProxyServerWorker.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts b/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts index 4a0fed465e..bfd51efed8 100644 --- a/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts +++ b/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts @@ -127,7 +127,10 @@ function isConnectBinding(request: Request): boolean { } function handleConnect(request: Request, env: Env): Response { - const address = request.headers.get("MF-Connect-Address")!; + const address = request.headers.get("MF-Connect-Address"); + if (address === null) { + return new Response("Missing MF-Connect-Address header", { status: 400 }); + } const fetcher = getExposedFetcher(request, env) as Fetcher; const { 0: client, 1: server } = new WebSocketPair(); From 2cc4927667c3d8bc120dc26d65f55764a2738888 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 24 Jul 2026 12:58:54 +0000 Subject: [PATCH 7/7] fix: null-guard the read-error close reason in the connect relay Mirror the write-path pattern in both relay copies so a non-object rejection can't throw while building the 1011 close reason and mask the original error. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01BuKCFc94BFDpbwag1UfPbB --- packages/miniflare/src/workers/shared/remote-bindings-utils.ts | 2 +- .../templates/remoteBindings/ProxyServerWorker.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/miniflare/src/workers/shared/remote-bindings-utils.ts b/packages/miniflare/src/workers/shared/remote-bindings-utils.ts index 55acdd9eab..848f275c43 100644 --- a/packages/miniflare/src/workers/shared/remote-bindings-utils.ts +++ b/packages/miniflare/src/workers/shared/remote-bindings-utils.ts @@ -333,7 +333,7 @@ export async function pipeSocketOverWebSocket( } closeWebSocket(1000); } catch (error) { - closeWebSocket(1011, (error as Error).message); + closeWebSocket(1011, (error as Error)?.message ?? "socket read failed"); throw error; } finally { reader.releaseLock(); diff --git a/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts b/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts index bfd51efed8..055bfeda91 100644 --- a/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts +++ b/packages/remote-bindings/templates/remoteBindings/ProxyServerWorker.ts @@ -295,7 +295,7 @@ async function pipeSocketOverWebSocket( } closeWebSocket(1000); } catch (error) { - closeWebSocket(1011, (error as Error).message); + closeWebSocket(1011, (error as Error)?.message ?? "socket read failed"); throw error; } finally { reader.releaseLock();