Skip to content
File

Blob: src/workerd/server/workerd-debug-port-client.c++

5.2 KB
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 
14namespace workerd::server {
15 
16namespace {
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.
22class 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 
58jsg::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 
71jsg::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 
93jsg::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 
108jsg::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