From 31036b285026c18bf754d7d3f34db51e874240c4 Mon Sep 17 00:00:00 2001 From: Felix Hanau Date: Mon, 2 Mar 2026 16:41:47 -0500 Subject: [PATCH] Add Socket::proxyTo() API This makes proxying sockets more convenient, should be especially useful for the connect handler. --- src/workerd/api/sockets.c++ | 17 +++++++++++++++++ src/workerd/api/sockets.h | 5 +++++ .../api/tests/connect-handler-test-proxy.js | 8 ++++---- .../generated-snapshot/experimental/index.d.ts | 1 + types/generated-snapshot/experimental/index.ts | 1 + types/generated-snapshot/index.d.ts | 1 + types/generated-snapshot/index.ts | 1 + 7 files changed, 30 insertions(+), 4 deletions(-) diff --git a/src/workerd/api/sockets.c++ b/src/workerd/api/sockets.c++ index 4dfd17227cc..57fafb2a45b 100644 --- a/src/workerd/api/sockets.c++ +++ b/src/workerd/api/sockets.c++ @@ -519,6 +519,23 @@ jsg::Promise Socket::close(jsg::Lock& js) { }); } +void Socket::proxyTo(jsg::Lock& js, jsg::Ref sock, jsg::Optional options) { + jsg::Optional optionsCopy = kj::none; + KJ_IF_SOME(o, options) { + optionsCopy = PipeToOptions{ + .preventAbort = o.preventAbort, + .preventCancel = o.preventCancel, + .preventClose = o.preventClose, + .signal = kj::none, + }; + KJ_IF_SOME(s, o.signal) { + KJ_ASSERT_NONNULL(optionsCopy).signal = s.addRef(); + } + } + sock->readable.pipeTo(js, writable, kj::mv(options).orDefault({})); + readable.pipeTo(js, sock->writable, kj::mv(optionsCopy).orDefault({})); +} + jsg::Ref Socket::startTls(jsg::Lock& js, jsg::Optional tlsOptions) { JSG_REQUIRE( secureTransport != SecureTransportKind::ON, TypeError, "Cannot startTls on a TLS socket."); diff --git a/src/workerd/api/sockets.h b/src/workerd/api/sockets.h index 3745d91dfe8..381cf61857e 100644 --- a/src/workerd/api/sockets.h +++ b/src/workerd/api/sockets.h @@ -131,6 +131,10 @@ class Socket: public jsg::Object { // closing. jsg::Promise close(jsg::Lock& js); + // Proxies to the other socket. Equivalent to: + // a.readable.pipeTo(b.writable); b.readable.pipeTo(a.writable); + void proxyTo(jsg::Lock& js, jsg::Ref sock, jsg::Optional options); + // Flushes write buffers then performs a TLS handshake on the current Socket connection. // The current `Socket` instance is closed and its readable/writable instances are also closed. // All new operations should be performed on the new `Socket` instance. @@ -176,6 +180,7 @@ class Socket: public jsg::Object { JSG_READONLY_PROTOTYPE_PROPERTY(secureTransport, getSecureTransport); JSG_METHOD(close); JSG_METHOD(startTls); + JSG_METHOD(proxyTo); JSG_TS_OVERRIDE({ get secureTransport(): 'on' | 'off' | 'starttls'; diff --git a/src/workerd/api/tests/connect-handler-test-proxy.js b/src/workerd/api/tests/connect-handler-test-proxy.js index a32bc9a2376..722439ed9ff 100644 --- a/src/workerd/api/tests/connect-handler-test-proxy.js +++ b/src/workerd/api/tests/connect-handler-test-proxy.js @@ -8,10 +8,10 @@ export class ConnectProxy extends WorkerEntrypoint { async connect(socket) { // proxy for ConnectEndpoint instance on port 8083. let upstream = connect('localhost:8083'); - await Promise.all([ - socket.readable.pipeTo(upstream.writable), - upstream.readable.pipeTo(socket.writable), - ]); + socket.proxyTo(upstream); + // proxyTo() can't be awaited – wait briefly so that we can be sure the data has been sent by + // the time we return so that the calling worker can read it right away. + await scheduler.wait(10); } } diff --git a/types/generated-snapshot/experimental/index.d.ts b/types/generated-snapshot/experimental/index.d.ts index cc15b634e36..4d257040b86 100755 --- a/types/generated-snapshot/experimental/index.d.ts +++ b/types/generated-snapshot/experimental/index.d.ts @@ -3910,6 +3910,7 @@ interface Socket { get secureTransport(): "on" | "off" | "starttls"; close(): Promise; startTls(options?: TlsOptions): Socket; + proxyTo(sock: Socket, options?: StreamPipeOptions): void; } interface SocketOptions { secureTransport?: string; diff --git a/types/generated-snapshot/experimental/index.ts b/types/generated-snapshot/experimental/index.ts index 3eba9e12c37..d422af53647 100755 --- a/types/generated-snapshot/experimental/index.ts +++ b/types/generated-snapshot/experimental/index.ts @@ -3920,6 +3920,7 @@ export interface Socket { get secureTransport(): "on" | "off" | "starttls"; close(): Promise; startTls(options?: TlsOptions): Socket; + proxyTo(sock: Socket, options?: StreamPipeOptions): void; } export interface SocketOptions { secureTransport?: string; diff --git a/types/generated-snapshot/index.d.ts b/types/generated-snapshot/index.d.ts index 444fb637745..0855f725f23 100755 --- a/types/generated-snapshot/index.d.ts +++ b/types/generated-snapshot/index.d.ts @@ -3822,6 +3822,7 @@ interface Socket { get secureTransport(): "on" | "off" | "starttls"; close(): Promise; startTls(options?: TlsOptions): Socket; + proxyTo(sock: Socket, options?: StreamPipeOptions): void; } interface SocketOptions { secureTransport?: string; diff --git a/types/generated-snapshot/index.ts b/types/generated-snapshot/index.ts index 06fa6ce5033..e484330b46b 100755 --- a/types/generated-snapshot/index.ts +++ b/types/generated-snapshot/index.ts @@ -3832,6 +3832,7 @@ export interface Socket { get secureTransport(): "on" | "off" | "starttls"; close(): Promise; startTls(options?: TlsOptions): Socket; + proxyTo(sock: Socket, options?: StreamPipeOptions): void; } export interface SocketOptions { secureTransport?: string;