From fafe713b2e5bcfc9a4448ea557fae7110713e8d8 Mon Sep 17 00:00:00 2001 From: Samuel Macleod Date: Fri, 14 Aug 2026 13:22:26 +0100 Subject: [PATCH 1/2] Add current-process debug port access --- src/workerd/io/io-channels.h | 5 + src/workerd/server/server-test.c++ | 11 +- src/workerd/server/server.c++ | 226 +++++++++--------- src/workerd/server/server.h | 2 + .../server/workerd-debug-port-client.c++ | 7 + .../server/workerd-debug-port-client.h | 17 +- 6 files changed, 148 insertions(+), 120 deletions(-) diff --git a/src/workerd/io/io-channels.h b/src/workerd/io/io-channels.h index 6fec6da5d6c..a4dec88fc5e 100644 --- a/src/workerd/io/io-channels.h +++ b/src/workerd/io/io-channels.h @@ -452,6 +452,11 @@ class IoChannelFactory: public virtual kj::Refcounted { JSG_FAIL_REQUIRE(Error, "WorkerdDebugPort bindings are not supported by this runtime."); } + // Get direct access to the current workerd process's debug port interface. + virtual rpc::WorkerdDebugPort::Client getWorkerdDebugPort() { + JSG_FAIL_REQUIRE(Error, "WorkerdDebugPort bindings are not supported by this runtime."); + } + // Converts a token created with {SubrequestChannel,ActorClassChannel}::getToken() back into a // live channel. Default implementations throw. virtual kj::Own subrequestChannelFromToken( diff --git a/src/workerd/server/server-test.c++ b/src/workerd/server/server-test.c++ index 78ff5c69635..16043661052 100644 --- a/src/workerd/server/server-test.c++ +++ b/src/workerd/server/server-test.c++ @@ -6995,9 +6995,9 @@ KJ_TEST("Server: debug port RPC calls") { } } -KJ_TEST("Server: workerdDebugPort binding loopback test") { - // This test verifies that a worker can use the workerdDebugPort binding to connect - // back to the same workerd instance's debug port and access other services. +KJ_TEST("Server: workerdDebugPort binding current process test") { + // This test verifies that a worker can use the workerdDebugPort binding to access other services + // in the same process without opening a network connection. TestServer test(R"(( services = [ ( name = "target-service", @@ -7029,8 +7029,7 @@ KJ_TEST("Server: workerdDebugPort binding loopback test") { esModule = `export default { ` async fetch(request, env, ctx) { - ` // Connect to the debug port - ` const client = await env.debugPort.connect("debug-addr"); + ` const client = env.debugPort.current(); ` ` // Test 1: Access the default entrypoint ` const defaultFetcher = client.getEntrypoint("target-service"); @@ -7066,8 +7065,6 @@ KJ_TEST("Server: workerdDebugPort binding loopback test") { ] ))"_kj); - // Enable the debug port on a known address - test.server.enableDebugPort(kj::str("debug-addr")); test.server.allowExperimental(); test.start(); diff --git a/src/workerd/server/server.c++ b/src/workerd/server/server.c++ index 594be83e177..2a3562e2468 100644 --- a/src/workerd/server/server.c++ +++ b/src/workerd/server/server.c++ @@ -3387,6 +3387,7 @@ class Server::WorkerService final: public Service, kj::Array> streamingTails; kj::Array> workerLoaders; kj::Maybe workerdDebugPortNetwork; + kj::Maybe workerdDebugPortServer; }; using LinkCallback = kj::Function; @@ -4461,6 +4462,14 @@ class Server::WorkerService final: public Service, "workerdDebugPort binding is not enabled for this worker"); } + rpc::WorkerdDebugPort::Client getWorkerdDebugPort() override { + auto& channels = + KJ_REQUIRE_NONNULL(ioChannels.tryGet(), "link() has not been called"); + return KJ_REQUIRE_NONNULL( + channels.workerdDebugPortServer, "workerdDebugPort binding is not enabled for this worker") + .makeWorkerdDebugPortClient(); + } + kj::Own subrequestChannelFromToken( ChannelTokenUsage usage, kj::ArrayPtr token) override { return channelTokenHandler.decodeSubrequestChannelToken(usage, token); @@ -5998,6 +6007,7 @@ kj::Promise> Server::makeWorkerImpl(kj::StringPtr if (def.hasWorkerdDebugPortBinding) { result.workerdDebugPortNetwork = network; + result.workerdDebugPortServer = *this; } return result; @@ -6633,133 +6643,135 @@ kj::Promise Server::listenTcp( // ======================================================================================= // Debug port for exposing all services via RPC -class Server::DebugPortListener { +class Server::WorkerdDebugPortImpl final: public rpc::WorkerdDebugPort::Server { public: - DebugPortListener(Server& owner, - kj::Own listener, - capnp::HttpOverCapnpFactory& httpOverCapnpFactory) - : owner(owner), - listener(kj::mv(listener)), + WorkerdDebugPortImpl(Server& srv, capnp::HttpOverCapnpFactory& httpOverCapnpFactory) + : srv(srv), httpOverCapnpFactory(httpOverCapnpFactory) {} - kj::Promise run() { - capnp::TwoPartyServer server(kj::heap(&owner, httpOverCapnpFactory)); - co_return co_await server.listen(*listener); - } + kj::Promise getEntrypoint(GetEntrypointContext context) override { + auto params = context.getParams(); + auto serviceName = params.getService(); + auto propsReader = params.getProps(); - private: - Server& owner; - kj::Own listener; - capnp::HttpOverCapnpFactory& httpOverCapnpFactory; + // Look up the service. + auto& serviceEntry = KJ_ASSERT_NONNULL(srv.services.find(serviceName), + kj::str("jsg.Error: Worker \"", serviceName, "\" not found")); + auto service = serviceEntry->service(); - class WorkerdDebugPortImpl final: public rpc::WorkerdDebugPort::Server { - public: - WorkerdDebugPortImpl( - workerd::server::Server* srvPtr, capnp::HttpOverCapnpFactory& httpOverCapnpFactory) - : srv(*srvPtr), - httpOverCapnpFactory(httpOverCapnpFactory) {} - - kj::Promise getEntrypoint(GetEntrypointContext context) override { - auto params = context.getParams(); - auto serviceName = params.getService(); - auto propsReader = params.getProps(); - - // Look up the service. - auto& serviceEntry = KJ_ASSERT_NONNULL(srv.services.find(serviceName), - kj::str("jsg.Error: Worker \"", serviceName, "\" not found")); - auto service = serviceEntry->service(); - - // Convert props from Frankenvalue if provided - Frankenvalue props; - if (params.hasProps()) { - props = Frankenvalue::fromCapnp(propsReader); - } + // Convert props from Frankenvalue if provided + Frankenvalue props; + if (params.hasProps()) { + props = Frankenvalue::fromCapnp(propsReader); + } - kj::Own targetService; + kj::Own targetService; - // Try to cast to WorkerService to support entrypoints and props - KJ_IF_SOME(workerService, kj::tryDowncast(*service)) { - // This is a WorkerService, use getEntrypoint which supports both entrypoints and props - kj::Maybe maybeEntrypoint; - if (params.hasEntrypoint()) { - maybeEntrypoint = params.getEntrypoint(); - } + // Try to cast to WorkerService to support entrypoints and props + KJ_IF_SOME(workerService, kj::tryDowncast(*service)) { + // This is a WorkerService, use getEntrypoint which supports both entrypoints and props + kj::Maybe maybeEntrypoint; + if (params.hasEntrypoint()) { + maybeEntrypoint = params.getEntrypoint(); + } - targetService = - KJ_ASSERT_NONNULL(workerService.getEntrypoint(maybeEntrypoint, kj::mv(props)), - kj::str("jsg.Error: Worker does not export an entrypoint named \"", - maybeEntrypoint.orDefault("(default)"), "\"")); - } else { - // Not a WorkerService - KJ_ASSERT(!params.hasEntrypoint(), "jsg.Error: Worker does not support named entrypoints"); + targetService = KJ_ASSERT_NONNULL(workerService.getEntrypoint(maybeEntrypoint, kj::mv(props)), + kj::str("jsg.Error: Worker does not export an entrypoint named \"", + maybeEntrypoint.orDefault("(default)"), "\"")); + } else { + // Not a WorkerService + KJ_ASSERT(!params.hasEntrypoint(), "jsg.Error: Worker does not support named entrypoints"); - // Try to apply props if the service supports it - if (params.hasProps()) { - targetService = service->forProps(kj::mv(props), Persistent::NO); - } else { - // No props, just use the service as-is - targetService = kj::addRef(*service); - } + // Try to apply props if the service supports it + if (params.hasProps()) { + targetService = service->forProps(kj::mv(props), Persistent::NO); + } else { + // No props, just use the service as-is + targetService = kj::addRef(*service); } - - // Return a WorkerdBootstrap that wraps this service using the generic implementation. - context.initResults(capnp::MessageSize{4, 1}) - .setEntrypoint( - kj::heap(kj::mv(targetService), httpOverCapnpFactory)); - return kj::READY_NOW; } - kj::Promise getActor(GetActorContext context) override { - auto params = context.getParams(); - auto serviceName = params.getService(); - auto entrypointName = params.getEntrypoint(); - auto actorIdStr = params.getActorId(); + // Return a WorkerdBootstrap that wraps this service using the generic implementation. + context.initResults(capnp::MessageSize{4, 1}) + .setEntrypoint(kj::heap(kj::mv(targetService), httpOverCapnpFactory)); + return kj::READY_NOW; + } - // Look up the service - auto& serviceEntry = KJ_ASSERT_NONNULL(srv.services.find(serviceName), - kj::str("jsg.Error: Worker \"", serviceName, "\" not found")); - auto service = serviceEntry->service(); + kj::Promise getActor(GetActorContext context) override { + auto params = context.getParams(); + auto serviceName = params.getService(); + auto entrypointName = params.getEntrypoint(); + auto actorIdStr = params.getActorId(); + + // Look up the service + auto& serviceEntry = KJ_ASSERT_NONNULL(srv.services.find(serviceName), + kj::str("jsg.Error: Worker \"", serviceName, "\" not found")); + auto service = serviceEntry->service(); + + // Try to cast to WorkerService + auto& workerService = KJ_REQUIRE_NONNULL(kj::tryDowncast(*service), + "jsg.Error: Worker does not support Durable Objects"); + + // Look up the actor namespace + auto& actorNamespace = KJ_ASSERT_NONNULL(workerService.getActorNamespace(entrypointName), + kj::str("jsg.Error: Worker does not export a Durable Object class named \"", entrypointName, + "\"")); + + // Create an actor ID - use the namespace config to determine if it's durable or ephemeral + Worker::Actor::Id actorId; + KJ_SWITCH_ONEOF(actorNamespace.getConfig()) { + KJ_CASE_ONEOF(c, Durable) { + // Durable Object ID (hex-encoded SHA256 hash) + auto decoded = kj::decodeHex(actorIdStr); + KJ_REQUIRE(decoded.size() == SHA256_DIGEST_LENGTH, + "Invalid Durable Object ID: expected 64 hex characters (32 bytes)", decoded.size()); + kj::Own id = + kj::heap(decoded.begin(), kj::none); + actorId = kj::mv(id); + } + KJ_CASE_ONEOF(c, Ephemeral) { + // Ephemeral actor ID (plain string) + actorId = kj::str(actorIdStr); + } + } - // Try to cast to WorkerService - auto& workerService = KJ_REQUIRE_NONNULL(kj::tryDowncast(*service), - "jsg.Error: Worker does not support Durable Objects"); + // Wrap the actor channel using the generic WorkerdBootstrap implementation. + context.initResults(capnp::MessageSize{4, 1}) + .setActor(kj::heap( + actorNamespace.getActorChannel(kj::mv(actorId)), httpOverCapnpFactory)); + return kj::READY_NOW; + } - // Look up the actor namespace - auto& actorNamespace = KJ_ASSERT_NONNULL(workerService.getActorNamespace(entrypointName), - kj::str("jsg.Error: Worker does not export a Durable Object class named \"", - entrypointName, "\"")); + private: + Server& srv; + capnp::HttpOverCapnpFactory& httpOverCapnpFactory; +}; - // Create an actor ID - use the namespace config to determine if it's durable or ephemeral - Worker::Actor::Id actorId; - KJ_SWITCH_ONEOF(actorNamespace.getConfig()) { - KJ_CASE_ONEOF(c, Durable) { - // Durable Object ID (hex-encoded SHA256 hash) - auto decoded = kj::decodeHex(actorIdStr); - KJ_REQUIRE(decoded.size() == SHA256_DIGEST_LENGTH, - "Invalid Durable Object ID: expected 64 hex characters (32 bytes)", decoded.size()); - kj::Own id = - kj::heap(decoded.begin(), kj::none); - actorId = kj::mv(id); - } - KJ_CASE_ONEOF(c, Ephemeral) { - // Ephemeral actor ID (plain string) - actorId = kj::str(actorIdStr); - } - } +class Server::DebugPortListener { + public: + DebugPortListener(Server& owner, + kj::Own listener, + capnp::HttpOverCapnpFactory& httpOverCapnpFactory) + : owner(owner), + listener(kj::mv(listener)), + httpOverCapnpFactory(httpOverCapnpFactory) {} - // Wrap the actor channel using the generic WorkerdBootstrap implementation. - context.initResults(capnp::MessageSize{4, 1}) - .setActor(kj::heap( - actorNamespace.getActorChannel(kj::mv(actorId)), httpOverCapnpFactory)); - return kj::READY_NOW; - } + kj::Promise run() { + capnp::TwoPartyServer server(owner.makeWorkerdDebugPortClient()); + co_return co_await server.listen(*listener); + } - private: - workerd::server::Server& srv; - capnp::HttpOverCapnpFactory& httpOverCapnpFactory; - }; + private: + Server& owner; + kj::Own listener; + capnp::HttpOverCapnpFactory& httpOverCapnpFactory; }; +rpc::WorkerdDebugPort::Client Server::makeWorkerdDebugPortClient() { + return rpc::WorkerdDebugPort::Client( + kj::heap(*this, globalContext->httpOverCapnpFactory)); +} + kj::Promise Server::listenDebugPort(kj::Own listener) { DebugPortListener obj(*this, kj::mv(listener), globalContext->httpOverCapnpFactory); co_return co_await obj.run(); diff --git a/src/workerd/server/server.h b/src/workerd/server/server.h index 689824b2746..c6e2b847fad 100644 --- a/src/workerd/server/server.h +++ b/src/workerd/server/server.h @@ -305,6 +305,7 @@ class Server final: private kj::TaskSet::ErrorHandler, private ChannelTokenHandl kj::Own listener, kj::Own service, kj::StringPtr addrStr); kj::Promise listenDebugPort(kj::Own listener); + rpc::WorkerdDebugPort::Client makeWorkerdDebugPortClient(); class InvalidConfigService; class InvalidConfigActorClass; @@ -318,6 +319,7 @@ class Server final: private kj::TaskSet::ErrorHandler, private ChannelTokenHandl class HttpListener; class TcpListener; class DebugPortListener; + class WorkerdDebugPortImpl; struct ErrorReporter; struct ConfigErrorReporter; diff --git a/src/workerd/server/workerd-debug-port-client.c++ b/src/workerd/server/workerd-debug-port-client.c++ index a8fee836b6a..a708d4623ed 100644 --- a/src/workerd/server/workerd-debug-port-client.c++ +++ b/src/workerd/server/workerd-debug-port-client.c++ @@ -130,4 +130,11 @@ jsg::Ref WorkerdDebugPortConnector::connect( return js.alloc(context.addObject(kj::mv(state))); } +jsg::Ref WorkerdDebugPortConnector::current(jsg::Lock& js) { + auto& context = IoContext::current(); + auto state = + kj::refcounted(context.getIoChannelFactory().getWorkerdDebugPort()); + return js.alloc(context.addObject(kj::mv(state))); +} + } // namespace workerd::server diff --git a/src/workerd/server/workerd-debug-port-client.h b/src/workerd/server/workerd-debug-port-client.h index 664ff5899ec..94a631e026b 100644 --- a/src/workerd/server/workerd-debug-port-client.h +++ b/src/workerd/server/workerd-debug-port-client.h @@ -17,12 +17,13 @@ class Fetcher; namespace workerd::server { -// Holds the I/O state for a debug port connection: the TCP stream, capnp RPC client, -// and debug port capability. Refcounted to support deferred proxying - response bodies -// and WebSockets are proxied through the capnp connection, so it must stay alive until -// they're fully consumed. See WorkerdBootstrapSubrequestChannel::startRequest(). +// Holds the I/O state for a debug port client. Refcounted to support deferred proxying - response +// bodies and WebSockets may use the capability after the originating request has completed. class DebugPortConnectionState: public kj::Refcounted { public: + explicit DebugPortConnectionState(rpc::WorkerdDebugPort::Client debugPort) + : debugPort(kj::mv(debugPort)) {} + DebugPortConnectionState(kj::Own connection, kj::Own rpcClient, rpc::WorkerdDebugPort::Client debugPort) @@ -34,8 +35,8 @@ class DebugPortConnectionState: public kj::Refcounted { return kj::addRef(*this); } - kj::Own connection; - kj::Own rpcClient; + kj::Maybe> connection; + kj::Maybe> rpcClient; rpc::WorkerdDebugPort::Client debugPort; }; @@ -105,8 +106,12 @@ class WorkerdDebugPortConnector: public jsg::Object { // @returns A WorkerdDebugPortClient that lazily connects on first use jsg::Ref connect(jsg::Lock& js, kj::String address); + // Access the current workerd process without opening a network connection. + jsg::Ref current(jsg::Lock& js); + JSG_RESOURCE_TYPE(WorkerdDebugPortConnector) { JSG_METHOD(connect); + JSG_METHOD(current); } }; From 919e625dbe3b7bd7f512eb52cae80c2aba378e17 Mon Sep 17 00:00:00 2001 From: Samuel Macleod Date: Fri, 14 Aug 2026 14:20:51 +0100 Subject: [PATCH 2/2] Fix debug port server type resolution --- src/workerd/server/server.c++ | 15 ++++++--------- 1 file changed, 6 insertions(+), 9 deletions(-) diff --git a/src/workerd/server/server.c++ b/src/workerd/server/server.c++ index 2a3562e2468..1355fffe2c1 100644 --- a/src/workerd/server/server.c++ +++ b/src/workerd/server/server.c++ @@ -6645,7 +6645,8 @@ kj::Promise Server::listenTcp( class Server::WorkerdDebugPortImpl final: public rpc::WorkerdDebugPort::Server { public: - WorkerdDebugPortImpl(Server& srv, capnp::HttpOverCapnpFactory& httpOverCapnpFactory) + WorkerdDebugPortImpl( + workerd::server::Server& srv, capnp::HttpOverCapnpFactory& httpOverCapnpFactory) : srv(srv), httpOverCapnpFactory(httpOverCapnpFactory) {} @@ -6743,18 +6744,15 @@ class Server::WorkerdDebugPortImpl final: public rpc::WorkerdDebugPort::Server { } private: - Server& srv; + workerd::server::Server& srv; capnp::HttpOverCapnpFactory& httpOverCapnpFactory; }; class Server::DebugPortListener { public: - DebugPortListener(Server& owner, - kj::Own listener, - capnp::HttpOverCapnpFactory& httpOverCapnpFactory) + DebugPortListener(Server& owner, kj::Own listener) : owner(owner), - listener(kj::mv(listener)), - httpOverCapnpFactory(httpOverCapnpFactory) {} + listener(kj::mv(listener)) {} kj::Promise run() { capnp::TwoPartyServer server(owner.makeWorkerdDebugPortClient()); @@ -6764,7 +6762,6 @@ class Server::DebugPortListener { private: Server& owner; kj::Own listener; - capnp::HttpOverCapnpFactory& httpOverCapnpFactory; }; rpc::WorkerdDebugPort::Client Server::makeWorkerdDebugPortClient() { @@ -6773,7 +6770,7 @@ rpc::WorkerdDebugPort::Client Server::makeWorkerdDebugPortClient() { } kj::Promise Server::listenDebugPort(kj::Own listener) { - DebugPortListener obj(*this, kj::mv(listener), globalContext->httpOverCapnpFactory); + DebugPortListener obj(*this, kj::mv(listener)); co_return co_await obj.run(); }