File
Blob: src/workerd/server/workerd-debug-port-client.c++
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #include "workerd-debug-port-client.h" |
| 6 | |
| 7 | #include <workerd/api/http.h> |
| 8 | #include <workerd/io/frankenvalue.h> |
| 9 | #include <workerd/io/io-context.h> |
| 10 | #include <workerd/io/worker-interface.h> |
| 11 | |
| 12 | #include <kj/memory.h> |
| 13 | |
| 14 | namespace workerd::server { |
| 15 | |
| 16 | namespace { |
| 17 | // A SubrequestChannel that makes requests to a remote worker via the debug port. |
| 18 | // |
| 19 | // The connection ref is attached to WorkerInterfaces returned by startRequest(). |
| 20 | // For HTTP fetch, the response body/WebSocket gets this attached (deferred proxying), |
| 21 | // ensuring the connection stays alive as long as the response is in use. |
| 22 | class WorkerdBootstrapSubrequestChannel final: public IoChannelFactory::SubrequestChannel { |
| 23 | public: |
| 24 | WorkerdBootstrapSubrequestChannel(rpc::WorkerdBootstrap::Client bootstrap, |
| 25 | capnp::HttpOverCapnpFactory& httpOverCapnpFactory, |
| 26 | capnp::ByteStreamFactory& byteStreamFactory, |
| 27 | kj::Own<DebugPortConnectionState> connectionState) |
| 28 | : bootstrap(kj::mv(bootstrap)), |
| 29 | httpOverCapnpFactory(httpOverCapnpFactory), |
| 30 | byteStreamFactory(byteStreamFactory), |
| 31 | connectionState(kj::mv(connectionState)) {} |
| 32 | |
| 33 | kj::Own<WorkerInterface> startRequest(IoChannelFactory::SubrequestMetadata metadata) override { |
| 34 | // Pass cfBlobJson as an RPC parameter on startEvent so the server can include it |
| 35 | // in SubrequestMetadata when creating the WorkerInterface. |
| 36 | auto req = bootstrap.startEventRequest(); |
| 37 | KJ_IF_SOME(cf, metadata.cfBlobJson) { |
| 38 | req.setCfBlobJson(cf); |
| 39 | } |
| 40 | auto dispatcher = req.send().getDispatcher(); |
| 41 | // Attach connection ref for deferred proxying - the HTTP response body/WebSocket |
| 42 | // will get this WorkerInterface attached, keeping the connection alive. |
| 43 | return kj::heap<RpcWorkerInterface>(httpOverCapnpFactory, byteStreamFactory, kj::mv(dispatcher)) |
| 44 | .attach(connectionState->addRef()); |
| 45 | } |
| 46 | |
| 47 | void requireAllowsTransfer() override { |
| 48 | JSG_FAIL_REQUIRE(Error, "WorkerdDebugPort bindings cannot be transferred to other workers"); |
| 49 | } |
| 50 | |
| 51 | private: |
| 52 | rpc::WorkerdBootstrap::Client bootstrap; |
| 53 | capnp::HttpOverCapnpFactory& httpOverCapnpFactory; |
| 54 | capnp::ByteStreamFactory& byteStreamFactory; |
| 55 | kj::Own<DebugPortConnectionState> connectionState; |
| 56 | }; |
| 57 | |
| 58 | jsg::Ref<api::Fetcher> wrapBootstrapAsFetcher(jsg::Lock& js, |
| 59 | IoContext& context, |
| 60 | rpc::WorkerdBootstrap::Client bootstrap, |
| 61 | kj::Own<DebugPortConnectionState> connectionState) { |
| 62 | kj::Own<IoChannelFactory::SubrequestChannel> subrequestChannel = |
| 63 | kj::refcounted<WorkerdBootstrapSubrequestChannel>(kj::mv(bootstrap), |
| 64 | context.getHttpOverCapnpFactory(), context.getByteStreamFactory(), |
| 65 | kj::mv(connectionState)); |
| 66 | return js.alloc<api::Fetcher>( |
| 67 | context.addObject(kj::mv(subrequestChannel)), api::Fetcher::RequiresHostAndProtocol::NO); |
| 68 | } |
| 69 | } // namespace |
| 70 | |
| 71 | jsg::Ref<api::Fetcher> WorkerdDebugPortClient::getEntrypoint(jsg::Lock& js, |
| 72 | kj::String service, |
| 73 | jsg::Optional<kj::String> entrypoint, |
| 74 | jsg::Optional<jsg::JsRef<jsg::JsObject>> props) { |
| 75 | auto& context = IoContext::current(); |
| 76 | |
| 77 | auto req = state->debugPort.getEntrypointRequest(); |
| 78 | req.setService(service); |
| 79 | KJ_IF_SOME(e, entrypoint) { |
| 80 | req.setEntrypoint(e); |
| 81 | } |
| 82 | KJ_IF_SOME(p, props) { |
| 83 | Frankenvalue::fromJs(js, p.getHandle(js)).toCapnp(req.initProps()); |
| 84 | } |
| 85 | |
| 86 | // Use Cap'n Proto pipelining: extract the entrypoint capability from the in-flight |
| 87 | // RPC response without waiting for it to resolve. The capability is a lazy proxy that |
| 88 | // only triggers the actual network round-trip when first used (e.g. fetch()). |
| 89 | auto bootstrap = req.send().getEntrypoint(); |
| 90 | return wrapBootstrapAsFetcher(js, context, kj::mv(bootstrap), state->addRef()); |
| 91 | } |
| 92 | |
| 93 | jsg::Ref<api::Fetcher> WorkerdDebugPortClient::getActor( |
| 94 | jsg::Lock& js, kj::String service, kj::String entrypoint, kj::String actorId) { |
| 95 | auto& context = IoContext::current(); |
| 96 | |
| 97 | auto req = state->debugPort.getActorRequest(); |
| 98 | req.setService(service); |
| 99 | req.setEntrypoint(entrypoint); |
| 100 | req.setActorId(actorId); |
| 101 | |
| 102 | // Use Cap'n Proto pipelining: extract the actor capability from the in-flight |
| 103 | // RPC response without waiting for it to resolve. |
| 104 | auto bootstrap = req.send().getActor(); |
| 105 | return wrapBootstrapAsFetcher(js, context, kj::mv(bootstrap), state->addRef()); |
| 106 | } |
| 107 | |
| 108 | jsg::Ref<WorkerdDebugPortClient> WorkerdDebugPortConnector::connect( |
| 109 | jsg::Lock& js, kj::String address) { |
| 110 | auto& context = IoContext::current(); |
| 111 | auto connectPromise = |
| 112 | context.getIoChannelFactory().getWorkerdDebugPortNetwork().parseAddress(address).then( |
| 113 | [](kj::Own<kj::NetworkAddress> addr) { return addr->connect(); }); |
| 114 | |
| 115 | // Use kj::newPromisedStream() to get an AsyncIoStream immediately. The actual TCP |
| 116 | // connection is deferred — Cap'n Proto pipelining queues all RPC calls until connected. |
| 117 | auto stream = kj::newPromisedStream(kj::mv(connectPromise)); |
| 118 | auto rpcClient = kj::heap<capnp::TwoPartyClient>(*stream); |
| 119 | auto debugPort = rpcClient->bootstrap().castAs<rpc::WorkerdDebugPort>(); |
| 120 | auto state = kj::refcounted<DebugPortConnectionState>( |
| 121 | kj::mv(stream), kj::mv(rpcClient), kj::mv(debugPort)); |
| 122 | return js.alloc<WorkerdDebugPortClient>(context.addObject(kj::mv(state))); |
| 123 | } |
| 124 | |
| 125 | } // namespace workerd::server |