Skip to content
File

Blob: src/workerd/api/worker-rpc.c++

98.8 KB
1// Copyright (c) 2017-2023 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/api/actor-state.h>
6#include <workerd/api/global-scope.h>
7#include <workerd/api/worker-rpc.h>
8#include <workerd/io/features.h>
9#include <workerd/io/tracer.h>
10#include <workerd/jsg/ser.h>
11#include <workerd/util/autogate.h>
12#include <workerd/util/completion-membrane.h>
13 
14#include <capnp/membrane.h>
15 
16namespace workerd::api {
17 
18namespace {
19 
20using StreamSinkFulfiller = kj::Own<kj::PromiseFulfiller<rpc::JsValue::StreamSink::Client>>;
21 
22} // namespace
23 
24// Implementation of StreamSink RPC interface. The stream sender calls `startStream()` when
25// serializing each stream, and the recipient calls `setSlot()` when deserializing streams to
26// provide the appropriate destination capability. This class is designed to allow these two
27// calls to happen in either order for each slot.
28class StreamSinkImpl final: public rpc::JsValue::StreamSink::Server, public kj::Refcounted {
29 public:
30 ~StreamSinkImpl() noexcept(false) {
31 for (auto& slot: table) {
32 KJ_IF_SOME(f, slot.tryGet<StreamFulfiller>()) {
33 f->reject(KJ_EXCEPTION(FAILED, "expected startStream() was never received"));
34 }
35 }
36 }
37 
38 void setSlot(uint i, capnp::Capability::Client stream) {
39 if (table.size() <= i) table.resize(i + 1);
40 
41 if (table[i] == nullptr) {
42 table[i] = kj::mv(stream);
43 } else KJ_SWITCH_ONEOF(table[i]) {
44 KJ_CASE_ONEOF(stream, capnp::Capability::Client) {
45 KJ_FAIL_REQUIRE("setSlot() tried to set the same slot twice", i);
46 }
47 KJ_CASE_ONEOF(fulfiller, StreamFulfiller) {
48 fulfiller->fulfill(kj::mv(stream));
49 table[i] = Consumed();
50 }
51 KJ_CASE_ONEOF(_, Consumed) {
52 KJ_FAIL_REQUIRE("setSlot() tried to set the same slot twice", i);
53 }
54 }
55 }
56 
57 kj::Promise<void> startStream(StartStreamContext context) override {
58 uint i = context.getParams().getExternalIndex();
59 
60 if (table.size() <= i) {
61 // guard against ridiculous table allocation
62 JSG_REQUIRE(i < 1024, Error, "Too many streams in one message.");
63 table.resize(i + 1);
64 }
65 
66 if (table[i] == nullptr) {
67 auto paf = kj::newPromiseAndFulfiller<capnp::Capability::Client>();
68 table[i] = kj::mv(paf.fulfiller);
69 context.getResults(capnp::MessageSize{4, 1}).setStream(kj::mv(paf.promise));
70 } else KJ_SWITCH_ONEOF(table[i]) {
71 KJ_CASE_ONEOF(stream, capnp::Capability::Client) {
72 context.getResults(capnp::MessageSize{4, 1}).setStream(kj::mv(stream));
73 table[i] = Consumed();
74 }
75 KJ_CASE_ONEOF(fulfiller, StreamFulfiller) {
76 KJ_FAIL_REQUIRE("startStream() tried to start the same stream twice", i);
77 }
78 KJ_CASE_ONEOF(_, Consumed) {
79 KJ_FAIL_REQUIRE("startStream() tried to start the same stream twice", i);
80 }
81 }
82 
83 return kj::READY_NOW;
84 }
85 
86 private:
87 using StreamFulfiller = kj::Own<kj::PromiseFulfiller<capnp::Capability::Client>>;
88 struct Consumed {};
89 
90 // Each slot starts out null (uninitialized). It becomes a Capability::Client if setSlot() is
91 // called first, or a StreamFulfiller if startStream() is called first. It becomes `Consumed`
92 // when the other method is called.
93 // HACK: Slots in the table take advantage of the little-known fact that OneOf has a "null"
94 // value, which is the value a OneOf has when default-initialized. This is useful because we
95 // don't want to explicitly initialize skipped slots. Maybe<OneOf> would be another option
96 // here, but would add 8 bytes to every slot just to store a boolean... feels bloated. There
97 // are only two methods in this class so I think it's OK.
98 using Slot = kj::OneOf<capnp::Capability::Client, StreamFulfiller, Consumed>;
99 
100 kj::Vector<Slot> table;
101};
102 
103kj::Maybe<rpc::JsValue::ExternalPusher::Client> RpcSerializerExternalHandler::getExternalPusher() {
104 KJ_IF_SOME(ep, externalPusher) {
105 return ep;
106 } else KJ_IF_SOME(func, getStreamHandlerFunc.tryGet<GetExternalPusherFunc>()) {
107 // First call, set up ExternalPusher.
108 return externalPusher.emplace(func());
109 } else {
110 // Using StreamSink.
111 return kj::none;
112 }
113}
114 
115capnp::Capability::Client RpcSerializerExternalHandler::writeStream(BuilderCallback callback) {
116 rpc::JsValue::StreamSink::Client* streamSinkPtr;
117 KJ_IF_SOME(ss, streamSink) {
118 streamSinkPtr = &ss;
119 } else {
120 // First stream written, set up the StreamSink.
121 auto& func = KJ_REQUIRE_NONNULL(getStreamHandlerFunc.tryGet<GetStreamSinkFunc>(),
122 "this serialization is not using StreamSink; use getExternalPusher() instead");
123 streamSinkPtr = &streamSink.emplace(func());
124 }
125 
126 auto result = ({
127 auto req = streamSinkPtr->startStreamRequest(capnp::MessageSize{4, 0});
128 req.setExternalIndex(externals.size());
129 req.send().getStream();
130 });
131 
132 write(kj::mv(callback));
133 
134 return result;
135}
136 
137capnp::Orphan<capnp::List<rpc::JsValue::External>> RpcSerializerExternalHandler::build(
138 capnp::Orphanage orphanage) {
139 auto result = orphanage.newOrphan<capnp::List<rpc::JsValue::External>>(externals.size());
140 auto builder = result.get();
141 for (auto i: kj::indices(externals)) {
142 externals[i](builder[i]);
143 }
144 return result;
145}
146 
147RpcDeserializerExternalHandler::~RpcDeserializerExternalHandler() noexcept(false) {
148 if (!unwindDetector.isUnwinding()) {
149 KJ_ASSERT(i == externals.size(), "deserialization did not consume all of the externals");
150 }
151}
152 
153rpc::JsValue::External::Reader RpcDeserializerExternalHandler::read() {
154 KJ_ASSERT(i < externals.size());
155 return externals[i++];
156}
157 
158void RpcDeserializerExternalHandler::setLastStream(capnp::Capability::Client stream) {
159 KJ_IF_SOME(ss, streamSink) {
160 ss.setSlot(i - 1, kj::mv(stream));
161 } else {
162 auto ss = kj::refcounted<StreamSinkImpl>();
163 ss->setSlot(i - 1, kj::mv(stream));
164 streamSink = *ss;
165 streamSinkCap = rpc::JsValue::StreamSink::Client(kj::mv(ss));
166 }
167}
168 
169namespace {
170 
171// Call to construct an `rpc::JsValue` from a JS value.
172//
173// `makeBuilder` is a function which takes a capnp::MessageSize hint and returns the
174// rpc::JsValue::Builder to fill in.
175template <typename Func>
176void serializeJsValue(jsg::Lock& js,
177 jsg::JsValue value,
178 RpcSerializerExternalHandler& externalHandler,
179 Func makeBuilder) {
180 jsg::Serializer serializer(js,
181 jsg::Serializer::Options{
182 .version = 15,
183 .omitHeader = false,
184 .treatClassInstancesAsPlainObjects = false,
185 .externalHandler = externalHandler,
186 });
187 serializer.write(js, value);
188 kj::Array<const byte> data = serializer.release().data;
189 JSG_ASSERT(data.size() <= MAX_JS_RPC_MESSAGE_SIZE, Error,
190 "Serialized RPC arguments or return values are limited to 32MiB, but the size of this value "
191 "was: ",
192 data.size(), " bytes.");
193 
194 capnp::MessageSize hint{0, 0};
195 hint.wordCount += (data.size() + sizeof(capnp::word) - 1) / sizeof(capnp::word);
196 hint.wordCount += capnp::sizeInWords<rpc::JsValue>();
197 hint.wordCount += externalHandler.size() * capnp::sizeInWords<rpc::JsValue::External>();
198 hint.capCount += externalHandler.size();
199 
200 rpc::JsValue::Builder builder = makeBuilder(hint);
201 
202 // TODO(perf): It would be nice if we could serialize directly into the capnp message to avoid
203 // a redundant copy of the bytes here. Maybe we could even cancel serialization early if it
204 // goes over the size limit.
205 builder.setV8Serialized(data);
206 
207 if (externalHandler.size() > 0) {
208 builder.adoptExternals(
209 externalHandler.build(capnp::Orphanage::getForMessageContaining(builder)));
210 }
211}
212 
213struct DeserializeResult {
214 jsg::JsValue value;
215 kj::Own<RpcStubDisposalGroup> disposalGroup;
216 kj::Maybe<rpc::JsValue::StreamSink::Client> streamSink;
217};
218 
219// Call to construct a JS value from an `rpc::JsValue`.
220DeserializeResult deserializeJsValue(jsg::Lock& js,
221 rpc::JsValue::Reader reader,
222 kj::LiteralStringConst debugContext,
223 kj::Maybe<StreamSinkImpl&> streamSink = kj::none) {
224 auto disposalGroup = kj::heap<RpcStubDisposalGroup>();
225 
226 RpcDeserializerExternalHandler externalHandler(
227 reader.getExternals(), *disposalGroup, streamSink, debugContext);
228 
229 jsg::Deserializer deserializer(js, reader.getV8Serialized(), kj::none, kj::none,
230 jsg::Deserializer::Options{
231 .version = 15,
232 .readHeader = true,
233 // Previously, while these are passing over an RPC boundary, we preserved stack
234 // traces in errors that happened to get passed through rather than thrown.
235 // This was mainly due, I believe, to a misunderstanding about whether or not
236 // v8 serialization preserved the stacks or not. When enhanced error serialization
237 // is disabled, stacks are preserved and this flag has no effect. When enhanced
238 // error serialization is enabled, then we'll switch to not preserving stacks in
239 // passed-through errors.
240 .preserveStackInErrors = false,
241 .externalHandler = externalHandler,
242 });
243 
244 return {
245 .value = deserializer.readValue(js),
246 .disposalGroup = kj::mv(disposalGroup),
247 .streamSink = externalHandler.getStreamSink(),
248 };
249}
250 
251// Does deserializeJsValue() and then adds a `dispose()` method to the returned object (if it is
252// an object) which disposes all stubs therein.
253jsg::JsValue deserializeRpcReturnValue(jsg::Lock& js,
254 rpc::JsRpcTarget::CallResults::Reader callResults,
255 kj::Maybe<StreamSinkImpl&> streamSink) {
256 auto [value, disposalGroup, ss] =
257 deserializeJsValue(js, callResults.getResult(), "return"_kjc, streamSink);
258 
259 if (streamSink == kj::none) {
260 KJ_REQUIRE(ss == kj::none,
261 "RPC returned result using StreamSink even though ExternalPusher was provided");
262 }
263 
264 // If the object had a disposer on the callee side, it will run when we discard the callPipeline,
265 // so attach that to the disposal group on the caller side. If the returned object did NOT have
266 // a disposer then we should discard callPipeline so that we don't hold open the callee's
267 // context for no reason.
268 if (callResults.getHasDisposer()) {
269 disposalGroup->setCallPipeline(
270 IoContext::current().addObject(kj::heap(callResults.getCallPipeline())));
271 }
272 
273 KJ_IF_SOME(obj, value.tryCast<jsg::JsObject>()) {
274 if (obj.isInstanceOf<JsRpcStub>(js)) {
275 // We're returning a plain stub. We don't need to override its `dispose` method.
276 disposalGroup->disownAll();
277 } else {
278 // Add a dispose method to the return object that disposes the DisposalGroup.
279 v8::Local<v8::Value> func = js.wrapSimpleFunction(js.v8Context(),
280 [disposalGroup = kj::mv(disposalGroup)](jsg::Lock&,
281 const v8::FunctionCallbackInfo<v8::Value>&) mutable { disposalGroup->disposeAll(); });
282 obj.setNonEnumerable(js, js.symbolDispose(), jsg::JsValue(func));
283 }
284 } else {
285 // Result wasn't an object, so it must not contain any stubs.
286 KJ_ASSERT(disposalGroup->empty());
287 }
288 
289 return value;
290}
291 
292// A membrane which attaches some object until it is destroyed.
293//
294// TODO(cleanup): This is generally useful, should it be part of capnp?
295class AttachmentMembrane final: public capnp::MembranePolicy, public kj::Refcounted {
296 public:
297 explicit AttachmentMembrane(kj::Own<void> attachment): attachment(kj::mv(attachment)) {}
298 
299 kj::Maybe<capnp::Capability::Client> inboundCall(
300 uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override {
301 return kj::none;
302 }
303 
304 kj::Maybe<capnp::Capability::Client> outboundCall(
305 uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override {
306 return kj::none;
307 }
308 
309 kj::Own<MembranePolicy> addRef() override {
310 return kj::addRef(*this);
311 }
312 
313 private:
314 kj::Own<void> attachment;
315};
316 
317// Given a value, check if it has a dispose method and, if so, invoke it.
318void tryCallDisposeMethod(jsg::Lock& js, jsg::JsValue value) {
319 js.withinHandleScope([&]() {
320 KJ_IF_SOME(obj, value.tryCast<jsg::JsObject>()) {
321 auto dispose = obj.get(js, js.symbolDispose());
322 if (dispose.isFunction()) {
323 jsg::check(v8::Local<v8::Value>(dispose).As<v8::Function>()->Call(
324 js.v8Context(), value, 0, nullptr));
325 }
326 }
327 });
328}
329 
330} // namespace
331 
332JsRpcPromise::JsRpcPromise(jsg::JsRef<jsg::JsPromise> inner,
333 kj::Own<WeakRef> weakRefParam,
334 IoOwn<rpc::JsRpcTarget::CallResults::Pipeline> pipeline)
335 : inner(kj::mv(inner)),
336 weakRef(kj::mv(weakRefParam)),
337 state(Pending{kj::mv(pipeline)}) {
338 KJ_REQUIRE(weakRef->ref == kj::none);
339 weakRef->ref = *this;
340}
341JsRpcPromise::~JsRpcPromise() noexcept(false) {
342 weakRef->ref = kj::none;
343}
344 
345void JsRpcPromise::resolve(jsg::Lock& js, jsg::JsValue result) {
346 if (state.is<Pending>()) {
347 state = Resolved{
348 .result = jsg::Value(js.v8Isolate, result),
349 .ctxCheck = IoContext::current().addObject(*this),
350 };
351 } else {
352 // We'd better dispose this.
353 tryCallDisposeMethod(js, result);
354 }
355}
356 
357void JsRpcPromise::dispose(jsg::Lock& js) {
358 KJ_IF_SOME(resolved, state.tryGet<Resolved>()) {
359 // Disposing the promise implies disposing the final result.
360 tryCallDisposeMethod(js, jsg::JsValue(resolved.result.getHandle(js)));
361 }
362 
363 state = Disposed();
364 weakRef->disposed = true;
365}
366 
367// See comment at call site for explanation.
368static rpc::JsRpcTarget::Client makeJsRpcTargetForSingleLoopbackCall(
369 jsg::Lock& js, jsg::JsObject obj);
370 
371rpc::JsRpcTarget::Client JsRpcPromise::getClientForOneCall(
372 jsg::Lock& js, kj::Vector<kj::StringPtr>& path) {
373 // (Don't extend `path` because we're the root.)
374 
375 KJ_SWITCH_ONEOF(state) {
376 KJ_CASE_ONEOF(pending, Pending) {
377 return pending.pipeline->getCallPipeline();
378 }
379 KJ_CASE_ONEOF(resolved, Resolved) {
380 // Dereference `ctxCheck` just to verify we're running in the correct context. (If not,
381 // this will throw.)
382 *resolved.ctxCheck;
383 
384 // A value was already returned, and we closed the original RPC pipeline. But the application
385 // kept the promise around and is still trying to pipeline on it. What do we do?
386 //
387 // A naive answer would be: We just return the actual value that was returned originally.
388 // Like if someone asked for `promise.foo.bar`, we just give them `returnValue.foo.bar`.
389 //
390 // That doesn't quite work, for a couple reasons:
391 // * If the caller is awaiting a property, they expect the result will have a `dispose()`
392 // method added to it, and that any stubs in the result will be independently disposable.
393 // This essentially means we need to clone the value so that we can dup() all the stubs and
394 // modify the result.
395 // * If the caller is trying to make a pipelined RPC call, they expect this call to go
396 // through all the usual RPC machinery. They do NOT expect that this is going to be a local
397 // call.
398 //
399 // The easiest way to make this all just work is... to actually wrap the value in a one-off
400 // RPC stub, and make a real RPC on it.
401 
402 return js.withinHandleScope([&]() -> rpc::JsRpcTarget::Client {
403 auto value = jsg::JsValue(resolved.result.getHandle(js));
404 
405 KJ_IF_SOME(obj, value.tryCast<jsg::JsObject>()) {
406 KJ_IF_SOME(stub, obj.tryUnwrapAs<JsRpcStub>(js)) {
407 // Oh, the return value is actually a stub itself. Just use it.
408 return stub->getClient();
409 } else {
410 // Must be a plain object.
411 return makeJsRpcTargetForSingleLoopbackCall(js, obj);
412 }
413 } else {
414 JSG_FAIL_REQUIRE(TypeError, "Can't pipeline on RPC that did not return an object.");
415 }
416 });
417 }
418 KJ_CASE_ONEOF(disposed, Disposed) {
419 return JSG_KJ_EXCEPTION(FAILED, Error, "RPC promise used after being disposed.");
420 }
421 }
422 KJ_UNREACHABLE;
423}
424 
425rpc::JsRpcTarget::Client JsRpcProperty::getClientForOneCall(
426 jsg::Lock& js, kj::Vector<kj::StringPtr>& path) {
427 auto result = parent->getClientForOneCall(js, path);
428 path.add(name);
429 return result;
430}
431 
432namespace {
433 
434struct JsRpcPromiseAndPipeline {
435 jsg::JsPromise promise;
436 kj::Own<JsRpcPromise::WeakRef> weakRef;
437 rpc::JsRpcTarget::CallResults::Pipeline pipeline;
438 
439 jsg::Ref<JsRpcPromise> asJsRpcPromise(jsg::Lock& js) && {
440 return js.alloc<JsRpcPromise>(jsg::JsRef<jsg::JsPromise>(js, promise), kj::mv(weakRef),
441 IoContext::current().addObject(kj::heap(kj::mv(pipeline))));
442 }
443};
444 
445// Core implementation of making an RPC call, reusable for many cases below.
446JsRpcPromiseAndPipeline callImpl(jsg::Lock& js,
447 JsRpcClientProvider& parent,
448 kj::Maybe<const kj::String&> name,
449 // If `maybeArgs` is provided, this is a call, otherwise it is a property access.
450 kj::Maybe<const v8::FunctionCallbackInfo<v8::Value>&> maybeArgs) {
451 // Note: We used to enforce that RPC methods had to be called with the correct `this`. That is,
452 // we prevented people from doing:
453 //
454 // let obj = {foo: someRpcStub.foo};
455 // obj.foo();
456 //
457 // This would throw "Illegal invocation", as is the norm when pulling methods of a native object.
458 // That worked as long as RPC methods were implemented as `jsg::Function`. However, when we
459 // switched to RPC methods being implemented as callable objects (JsRpcProperty), this became
460 // impossible, because V8's SetCallAsFunctionHandler() arranges that `this` is bound to the
461 // callable object itself, regardless of how it was invoked. So now we cannot detect the
462 // situation above, because V8 never tells us about `obj` at all.
463 //
464 // Oh well. It's not a big deal. Just annoying that we have to forever support tearing RPC
465 // methods off their source object, even if we change implementations to something where that's
466 // less convenient.
467 
468 try {
469 return js.tryCatch([&]() -> JsRpcPromiseAndPipeline {
470 // `path` will be filled in with the path of property names leading from the stub represented by
471 // `client` to the specific property / method that we're trying to invoke.
472 kj::Vector<kj::StringPtr> path;
473 auto client = parent.getClientForOneCall(js, path);
474 
475 auto& ioContext = IoContext::current();
476 
477 KJ_IF_SOME(lock, ioContext.waitForOutputLocksIfNecessary()) {
478 // Replace the client with a promise client that will delay the call until the output gate
479 // is open.
480 client = lock.then([client = kj::mv(client)]() mutable { return kj::mv(client); });
481 }
482 
483 auto builder = client.callRequest();
484 
485 // This code here is slightly overcomplicated in order to avoid pushing anything to the
486 // kj::Vector in the common case that the parent path is empty. I'm probably trying too hard
487 // but oh well.
488 if (path.empty()) {
489 KJ_IF_SOME(n, name) {
490 builder.setMethodName(n);
491 } else {
492 // No name and no path, must be directly calling a stub.
493 builder.initMethodPath(0);
494 }
495 } else {
496 auto pathBuilder = builder.initMethodPath(path.size() + (name != kj::none));
497 for (auto i: kj::indices(path)) {
498 pathBuilder.set(i, path[i]);
499 }
500 KJ_IF_SOME(n, name) {
501 pathBuilder.set(path.size(), n);
502 }
503 }
504 
505 kj::Maybe<StreamSinkFulfiller> paramsStreamSinkFulfiller;
506 
507 bool useExternalPusher =
508 util::Autogate::isEnabled(util::AutogateKey::RPC_USE_EXTERNAL_PUSHER);
509 
510 KJ_IF_SOME(args, maybeArgs) {
511 // If we have arguments, serialize them.
512 // Note that we may fail to serialize some element, in which case this will throw back to
513 // JS.
514 if (args.Length() > 0) {
515 // This is a function call with arguments.
516 v8::LocalVector<v8::Value> argv(js.v8Isolate, args.Length());
517 for (int n = 0; n < args.Length(); n++) {
518 argv[n] = args[n];
519 }
520 auto arr = v8::Array::New(js.v8Isolate, argv.data(), argv.size());
521 
522 auto stubOwnership = FeatureFlags::get(js).getRpcParamsDupStubs()
523 ? RpcSerializerExternalHandler::DUPLICATE
524 : RpcSerializerExternalHandler::TRANSFER;
525 
526 RpcSerializerExternalHandler::GetStreamHandlerFunc getStreamHandlerFunc;
527 if (useExternalPusher) {
528 getStreamHandlerFunc.init<RpcSerializerExternalHandler::GetExternalPusherFunc>(
529 [&]() -> rpc::JsValue::ExternalPusher::Client { return client; });
530 } else {
531 getStreamHandlerFunc.init<RpcSerializerExternalHandler::GetStreamSinkFunc>([&]() {
532 // A stream was encountered in the params, so we must expect the response to contain
533 // paramsStreamSink. But we don't have the response yet. So, we need to set up a
534 // temporary promise client, which we hook to the response a little bit later.
535 auto paf = kj::newPromiseAndFulfiller<rpc::JsValue::StreamSink::Client>();
536 paramsStreamSinkFulfiller = kj::mv(paf.fulfiller);
537 return kj::mv(paf.promise);
538 });
539 }
540 
541 RpcSerializerExternalHandler externalHandler(stubOwnership, kj::mv(getStreamHandlerFunc));
542 serializeJsValue(js, jsg::JsValue(arr), externalHandler, [&](capnp::MessageSize hint) {
543 // TODO(perf): Actually use the size hint.
544 return builder.getOperation().initCallWithArgs();
545 });
546 }
547 } else {
548 // This is a property access.
549 builder.getOperation().setGetProperty();
550 }
551 
552 kj::Maybe<kj::Own<StreamSinkImpl>> resultStreamSink;
553 if (useExternalPusher) {
554 // Unfortunately, we always have to send the ExternalPusher since we don't know whether the
555 // call will return any streams (or other pushed externals). Luckily, it's a
556 // one-per-IoContext object, not a big deal. (It'll take a slot on the capnp export table
557 // though.)
558 builder.getResultsStreamHandler().setExternalPusher(ioContext.getExternalPusher());
559 } else {
560 // Unfortunately, we always have to send a `resultsStreamSink` because we don't know until
561 // after the call completes whether or not it will return any streams. If it's unused,
562 // though, it should only be a couple allocations.
563 builder.getResultsStreamHandler().setStreamSink(
564 kj::addRef(*resultStreamSink.emplace(kj::refcounted<StreamSinkImpl>())));
565 }
566 
567 auto callResult = builder.send();
568 
569 KJ_IF_SOME(ssf, paramsStreamSinkFulfiller) {
570 ssf->fulfill(callResult.getParamsStreamSink());
571 }
572 
573 // We need to arrange that our JsRpcPromise will updated in-place with the final settlement
574 // of this RPC promise. However, we can't actually construct the JsRpcPromise until we have
575 // the final promise to give it. To resolve the cycle, we only create a JsRpcPromise::WeakRef
576 // here, which is filled in later on to point at the JsRpcPromise, if and when one is created.
577 auto weakRef = kj::atomicRefcounted<JsRpcPromise::WeakRef>();
578 
579 // RemotePromise lets us consume its pipeline and promise portions independently; we consume
580 // the promise here and we consume the pipeline below, both via kj::mv().
581 auto jsPromise = ioContext.awaitIo(js, kj::mv(callResult),
582 [weakRef = kj::atomicAddRef(*weakRef), resultStreamSink = kj::mv(resultStreamSink)](
583 jsg::Lock& js,
584 capnp::Response<rpc::JsRpcTarget::CallResults> response) mutable -> jsg::Value {
585 auto jsResult = deserializeRpcReturnValue(js, response, resultStreamSink);
586 
587 if (weakRef->disposed) {
588 // The promise was explicitly disposed before it even resolved. This means we must dispose
589 // the returned object as well.
590 tryCallDisposeMethod(js, jsResult);
591 } else {
592 KJ_IF_SOME(r, weakRef->ref) {
593 r.resolve(js, jsResult);
594 }
595 }
596 
597 return jsg::Value(js.v8Isolate, jsResult);
598 });
599 
600 return {
601 .promise = jsg::JsPromise(js.wrapSimplePromise(kj::mv(jsPromise))),
602 .weakRef = kj::mv(weakRef),
603 .pipeline = kj::mv(callResult),
604 };
605 }, [&](jsg::Value error) -> JsRpcPromiseAndPipeline {
606 // Probably a serialization error. Need to convert to an async error since we never throw
607 // synchronously from async functions.
608 auto jsError = jsg::JsValue(error.getHandle(js));
609 auto pipeline = capnp::newBrokenPipeline(js.exceptionToKj(jsError));
610 return {.promise = js.rejectedJsPromise(jsError),
611 .weakRef = kj::atomicRefcounted<JsRpcPromise::WeakRef>(),
612 .pipeline =
613 rpc::JsRpcTarget::CallResults::Pipeline(capnp::AnyPointer::Pipeline(kj::mv(pipeline)))};
614 });
615 } catch (jsg::JsExceptionThrown&) {
616 // This must be a termination exception, or we would have caught it above.
617 throw;
618 } catch (...) {
619 // Catch KJ exceptions and make them async, since we don't want async calls to throw
620 // synchronously.
621 auto e = kj::getCaughtExceptionAsKj();
622 auto pipeline = capnp::newBrokenPipeline(e.clone());
623 return {
624 .promise = jsg::JsPromise(js.wrapSimplePromise(js.rejectedPromise<jsg::Value>(kj::mv(e)))),
625 .weakRef = kj::atomicRefcounted<JsRpcPromise::WeakRef>(),
626 .pipeline =
627 rpc::JsRpcTarget::CallResults::Pipeline(capnp::AnyPointer::Pipeline(kj::mv(pipeline)))};
628 }
629}
630 
631} // namespace
632 
633jsg::Ref<JsRpcPromise> JsRpcProperty::call(const v8::FunctionCallbackInfo<v8::Value>& args) {
634 jsg::Lock& js = jsg::Lock::from(args.GetIsolate());
635 
636 return callImpl(js, *parent, name, args).asJsRpcPromise(js);
637}
638 
639jsg::Ref<JsRpcPromise> JsRpcStub::call(const v8::FunctionCallbackInfo<v8::Value>& args) {
640 jsg::Lock& js = jsg::Lock::from(args.GetIsolate());
641 
642 return callImpl(js, *this, kj::none, args).asJsRpcPromise(js);
643}
644 
645jsg::Ref<JsRpcPromise> JsRpcPromise::call(const v8::FunctionCallbackInfo<v8::Value>& args) {
646 jsg::Lock& js = jsg::Lock::from(args.GetIsolate());
647 
648 return callImpl(js, *this, kj::none, args).asJsRpcPromise(js);
649}
650 
651namespace {
652 
653jsg::JsValue thenImpl(jsg::Lock& js,
654 v8::Local<v8::Promise> promise,
655 v8::Local<v8::Function> handler,
656 jsg::Optional<v8::Local<v8::Function>> errorHandler) {
657 KJ_IF_SOME(e, errorHandler) {
658 // Note that we intentionally propagate any exception from promise->Then() synchronously since
659 // if V8's native Promise threw synchronously from `then()`, we might as well too. Anyway it's
660 // probably a termination exception.
661 return jsg::JsPromise(jsg::check(promise->Then(js.v8Context(), handler, e)));
662 } else {
663 return jsg::JsPromise(jsg::check(promise->Then(js.v8Context(), handler)));
664 }
665}
666 
667jsg::JsValue catchImpl(
668 jsg::Lock& js, v8::Local<v8::Promise> promise, v8::Local<v8::Function> errorHandler) {
669 return jsg::JsPromise(jsg::check(promise->Catch(js.v8Context(), errorHandler)));
670}
671 
672jsg::JsValue finallyImpl(
673 jsg::Lock& js, v8::Local<v8::Promise> promise, v8::Local<v8::Function> onFinally) {
674 // HACK: `finally()` is not exposed as a C++ API, so we have to manually read it from JS.
675 jsg::JsObject obj(promise);
676 auto func = obj.get(js, "finally");
677 KJ_ASSERT(func.isFunction());
678 v8::Local<v8::Value> param = onFinally;
679 return jsg::JsValue(jsg::check(
680 v8::Local<v8::Value>(func).As<v8::Function>()->Call(js.v8Context(), obj, 1, &param)));
681}
682 
683} // namespace
684 
685jsg::JsValue JsRpcProperty::then(jsg::Lock& js,
686 v8::Local<v8::Function> handler,
687 jsg::Optional<v8::Local<v8::Function>> errorHandler) {
688 auto promise = callImpl(js, *parent, name, kj::none).promise;
689 
690 return thenImpl(js, promise, handler, errorHandler);
691}
692 
693jsg::JsValue JsRpcProperty::catch_(jsg::Lock& js, v8::Local<v8::Function> errorHandler) {
694 auto promise = callImpl(js, *parent, name, kj::none).promise;
695 
696 return catchImpl(js, promise, errorHandler);
697}
698 
699jsg::JsValue JsRpcProperty::finally(jsg::Lock& js, v8::Local<v8::Function> onFinally) {
700 auto promise = callImpl(js, *parent, name, kj::none).promise;
701 
702 return finallyImpl(js, promise, onFinally);
703}
704 
705jsg::JsValue JsRpcPromise::then(jsg::Lock& js,
706 v8::Local<v8::Function> handler,
707 jsg::Optional<v8::Local<v8::Function>> errorHandler) {
708 return thenImpl(js, inner.getHandle(js), handler, errorHandler);
709}
710 
711jsg::JsValue JsRpcPromise::catch_(jsg::Lock& js, v8::Local<v8::Function> errorHandler) {
712 return catchImpl(js, inner.getHandle(js), errorHandler);
713}
714 
715jsg::JsValue JsRpcPromise::finally(jsg::Lock& js, v8::Local<v8::Function> onFinally) {
716 return finallyImpl(js, inner.getHandle(js), onFinally);
717}
718 
719kj::Maybe<jsg::Ref<JsRpcProperty>> JsRpcProperty::getProperty(jsg::Lock& js, kj::String name) {
720 return js.alloc<JsRpcProperty>(JSG_THIS, kj::mv(name));
721}
722 
723kj::Maybe<jsg::Ref<JsRpcProperty>> JsRpcPromise::getProperty(jsg::Lock& js, kj::String name) {
724 return js.alloc<JsRpcProperty>(JSG_THIS, kj::mv(name));
725}
726 
727JsRpcStub::JsRpcStub(IoOwn<rpc::JsRpcTarget::Client> capnpClient,
728 RpcStubDisposalGroup& disposalGroup,
729 jsg::ExternalMemoryAdjustment externalMemoryAdjustment)
730 : capnpClient(kj::mv(capnpClient)),
731 disposalGroup(disposalGroup),
732 externalMemoryAdjustment(kj::mv(externalMemoryAdjustment)) {
733 disposalGroup.list.add(*this);
734}
735 
736JsRpcStub::~JsRpcStub() noexcept(false) {
737 KJ_IF_SOME(d, disposalGroup) {
738 d.list.remove(*this);
739 }
740 
741 KJ_IF_SOME(c, capnpClient) {
742 // The app failed to dispose the stub; it leaked. We'd rather not make GC observable, so we
743 // must pass the capnp capability off to the I/O context to be dropped when the I/O context
744 // itself shuts down.
745 kj::mv(c).deferGcToContext();
746 
747 // In preview, let's try to warn the developer about the problem.
748 //
749 // TODO(cleanup): Instead of logging this warning at GC time, it would be better if we logged
750 // it at the time that the client is destroyed, i.e. when the IoContext is torn down,
751 // which is usually sooner (and more deterministic). But logging a warning during
752 // IoContext tear-down is problematic since logWarningOnce() is a method on
753 // IoContext...
754 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
755 ioContext.logWarningOnce(
756 "An RPC stub was not disposed properly. You must call dispose() on all stubs in order to "
757 "let the other side know that you are no longer using them. You cannot rely on "
758 "the garbage collector for this because it may take arbitrarily long before actually "
759 "collecting unreachable objects. As a shortcut, calling dispose() on the result of "
760 "an RPC call disposes all stubs within it."_kj);
761 }
762 }
763}
764 
765RpcStubDisposalGroup::~RpcStubDisposalGroup() noexcept(false) {
766 if (jsg::isInGcDestructor()) {
767 // If the disposal group was dropped as a result of garbage collection, we should NOT actually
768 // dispose any stubs. In particular:
769 // * If an application never invokes dispose() on an RPC result and the result is GC'ed, the
770 // app could still be holding onto stubs that came from that result. We don't want to
771 // dispose those unexpectedly.
772 // * If an incoming RPC call does something like `await new Promise(() => {})` to hang
773 // forever, the promise reaction can be GC'ed even though the call didn't really complete.
774 // We don't want to dispose param stubs in this case.
775 disownAll();
776 
777 // If we have a `callPipeline`, it means we called an RPC that returned an object, and that
778 // object had a dispose method defined on the server side. We don't want it to observe GC,
779 // so we'll defer dropping the pipeline until the IoContext is destroyed.
780 //
781 // (We don't do this as part of disownAll() because the one other call site of disownAll()
782 // is only invoked in cases where there shouldn't be a `callPipeline` anyway...)
783 KJ_IF_SOME(c, callPipeline) {
784 kj::mv(c).deferGcToContext();
785 
786 // In preview, let's try to warn the developer about the problem.
787 //
788 // TODO(cleanup): Same comment as in ~JsRpcStub().
789 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
790 ioContext.logWarningOnce(
791 "An RPC result was not disposed properly. One of the RPC calls you made expects you "
792 "to call dispose() on the return value, but you didn't do so. You cannot rely on "
793 "the garbage collector for this because it may take arbitrarily long before actually "
794 "collecting unreachable objects."_kj);
795 }
796 }
797 } else {
798 // However, if we're destroying the RpcStubDisposalGroup NOT as a result of GC, this probably
799 // means one of:
800 // * This is the disposal group for an incoming RPC call, and that call completed. The group
801 // was attached to the completion continuation, which executed, and is now being destroyed.
802 // This is the normal completion case, and we should dispose all the param stubs.
803 // * An exception was thrown in the RPC implementation before stubs could be passed to
804 // JavaScript in the first place, resulting in the disposal group being destroyed during
805 // exception unwind. The stubs should be disposed proactively since they were never
806 // received.
807 disposeAll();
808 }
809}
810 
811rpc::JsRpcTarget::Client JsRpcStub::getClient() {
812 KJ_IF_SOME(c, capnpClient) {
813 return *c;
814 } else {
815 // TODO(soon): Improve the error message to describe why it was disposed.
816 return JSG_KJ_EXCEPTION(FAILED, Error, "RPC stub used after being disposed.");
817 }
818}
819 
820rpc::JsRpcTarget::Client JsRpcStub::getClientForOneCall(
821 jsg::Lock& js, kj::Vector<kj::StringPtr>& path) {
822 // (Don't extend `path` because we're the root.)
823 return getClient();
824}
825 
826jsg::Ref<JsRpcStub> JsRpcStub::dup(jsg::Lock& js) {
827 return js.alloc<JsRpcStub>(IoContext::current().addObject(kj::heap(getClient())));
828}
829 
830void JsRpcStub::dispose() {
831 capnpClient = kj::none;
832 externalMemoryAdjustment = kj::none;
833 KJ_IF_SOME(d, disposalGroup) {
834 d.list.remove(*this);
835 disposalGroup = kj::none;
836 }
837}
838 
839void RpcStubDisposalGroup::disownAll() {
840 for (auto& stub: list) {
841 stub.disposalGroup = kj::none;
842 list.remove(stub);
843 }
844}
845 
846void RpcStubDisposalGroup::disposeAll() {
847 for (auto& stub: list) {
848 stub.dispose();
849 }
850 callPipeline = kj::none;
851 
852 // Each stub should have removed itself.
853 KJ_ASSERT(list.empty());
854}
855 
856kj::Maybe<jsg::Ref<JsRpcProperty>> JsRpcStub::getRpcMethod(jsg::Lock& js, kj::String name) {
857 // Do not return a method for `then`, otherwise JavaScript decides this is a thenable, i.e. a
858 // custom Promise, which will mean a Promise that resolves to this object will attempt to chain
859 // with it, which is not what you want!
860 if (name == "then"_kj) return kj::none;
861 
862 return js.alloc<JsRpcProperty>(JSG_THIS, kj::mv(name));
863}
864 
865void JsRpcStub::serialize(jsg::Lock& js, jsg::Serializer& serializer) {
866 auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError,
867 "Remote RPC references can only be serialized for RPC.");
868 auto externalHandler = dynamic_cast<RpcSerializerExternalHandler*>(&handler);
869 JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError,
870 "Remote RPC references can only be serialized for RPC.");
871 
872 // We may be forwarding a stub that points to some other isolate. Consider the case where we
873 // are returning the stub to our client. The RPC session remains live as long as the client is
874 // holding any remaining stubs obtained from this session, due to CompletionMembrane. However, if
875 // the only remaining stubs point on to different isolates, and we don't have anything left to
876 // do in this IoContext, then the pending event mechanism would abort the IoContext early with
877 // "The script will never generate a response." To avoid that, we need to attach a pending event
878 // to this stub, using a membrane.
879 //
880 // TODO(someday): Ideally, we would not need to keep the IoContext live just because stubs pass
881 // through it. It would be nice to implement a sort of "deferred proxying" for RPC, where we
882 // shut down the IoContext when it has nothing left to do. Note, though, that if the IoContext
883 // is explicitly *aborted*, we probably should revoke all capabilities obtained through it.
884 // That actually doesn't quite happen today: aborting the IoContext is likely to cancel all
885 // subrequests which probably has the effect of breaking any stubs obtained from them, but
886 // not necessarily (the subrequests could use waitUntil() to extend themselves). Anyway, this
887 // will be trickier to get right, so I'm punting with this work-around for now.
888 auto cap = capnp::membrane(
889 getClient(), kj::refcounted<AttachmentMembrane>(IoContext::current().registerPendingEvent()));
890 
891 externalHandler->write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable {
892 builder.setRpcTarget(kj::mv(cap));
893 });
894 
895 if (externalHandler->getStubOwnership() == RpcSerializerExternalHandler::TRANSFER) {
896 // Instead of disposing the stub immediately, we add a disposer to the serializer
897 // that will be executed when the pipeline is finished. This ensures the stub
898 // remains valid for the duration of any pipelined operations.
899 externalHandler->addStubDisposer(
900 kj::heap(kj::defer([self = JSG_THIS]() mutable { self->dispose(); })));
901 }
902}
903 
904jsg::Ref<JsRpcStub> JsRpcStub::deserialize(
905 jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) {
906 auto& handler = KJ_REQUIRE_NONNULL(
907 deserializer.getExternalHandler(), "got JsRpcStub on non-RPC serialized object?");
908 auto externalHandler = dynamic_cast<RpcDeserializerExternalHandler*>(&handler);
909 KJ_REQUIRE(externalHandler != nullptr, "got JsRpcStub on non-RPC serialized object?");
910 
911 auto reader = externalHandler->read();
912 KJ_REQUIRE(reader.isRpcTarget(), "external table slot type doesn't match serialization tag");
913 
914 auto& ioctx = IoContext::current();
915 
916 // Account for membrane/promise memory in the KJ heap (~1600 bytes per stub from profiling).
917 static constexpr size_t ESTIMATED_EXTERNAL_MEMORY_PER_STUB = 1600;
918 auto externalMemory = js.getExternalMemoryAdjustment(ESTIMATED_EXTERNAL_MEMORY_PER_STUB);
919 
920 return js.alloc<JsRpcStub>(ioctx.addObject(kj::heap(reader.getRpcTarget())),
921 externalHandler->getDisposalGroup(), kj::mv(externalMemory));
922}
923 
924static bool isFunctionForRpc(jsg::Lock& js, v8::Local<v8::Function> func) {
925 jsg::JsObject obj(func);
926 if (obj.isInstanceOf<JsRpcProperty>(js) || obj.isInstanceOf<JsRpcPromise>(js)) {
927 // Don't allow JsRpcProperty or JsRpcPromise to be treated as plain functions, even though they
928 // are technically callable. These types need to be treated specially (if we decide to let
929 // them be passed over RPC at all).
930 return false;
931 }
932 return true;
933}
934 
935static bool isFunctionForRpc(jsg::Lock& js, jsg::JsValue value) {
936 if (!value.isFunction()) return false;
937 return isFunctionForRpc(js, v8::Local<v8::Value>(value).As<v8::Function>());
938}
939 
940// `makeCallPipeline()` has a bit of a complicated result type..
941namespace MakeCallPipeline {
942// The value is an object, which may have stubs inside it.
943struct Object {
944 rpc::JsRpcTarget::Client cap;
945 
946 // Was the value a plain JavaScript object which had a custom dispose() method?
947 bool hasDispose;
948};
949 
950// The value was something that should serialize to a single stub (e.g. it was an RpcTarget, a
951// plain function, or already a stub). The callPipeline should simply be a copy of that stub.
952struct SingleStub {};
953 
954// The value is not a type that supports pipelining. It may still be serializable, and it could
955// even contain stubs (e.g. in a Map).
956struct NonPipelinable {
957 // callPipeline to return just for error-handling purposes.
958 rpc::JsRpcTarget::Client errorPipeline;
959};
960 
961using Result = kj::OneOf<Object, SingleStub, NonPipelinable>;
962}; // namespace MakeCallPipeline
963 
964template <typename Func>
965MakeCallPipeline::Result serializeJsValueWithPipeline(jsg::Lock& js,
966 jsg::JsValue value,
967 Func makeBuilder,
968 RpcSerializerExternalHandler::GetStreamHandlerFunc getStreamSinkFunc);
969 
970// Callee-side implementation of JsRpcTarget.
971//
972// Most of the implementation is in this base class. There are subclasses specializing for the case
973// of a top-level entrypoint vs. a transient object introduced by a previous RPC in the same
974// session.
975class JsRpcTargetBase: public rpc::JsRpcTarget::Server {
976 public:
977 struct MayOutliveIncomingRequest {};
978 struct CantOutliveIncomingRequest {};
979 
980 // Constructor used by TransientJsRpcTarget, which does not own the context. It needs to use
981 // makeReentryCallback() to guard against the possibility that the IoContext is canceled before
982 // or during a call.
983 JsRpcTargetBase(IoContext& ctx, MayOutliveIncomingRequest)
984 : enterIsolateAndCall(ctx.makeReentryCallback<IoContext::TOP_UP>(
985 [this, &ctx](Worker::Lock& lock, CallContext callContext) {
986 return callImpl(lock, ctx, callContext);
987 })),
988 externalPusher(ctx.getExternalPusher()) {}
989 
990 // Constructor use by EntrypointJsRpcTarget, which is revoked and destroyed before the IoContext
991 // can possibly be canceled. It can just use ctx.run().
992 JsRpcTargetBase(IoContext& ctx, CantOutliveIncomingRequest)
993 : enterIsolateAndCall([this, &ctx](CallContext callContext) {
994 // Note: No need to topUpActor() since this is the start of a top-level request, so the
995 // actor will already have been topped up by IncomingRequest::delivered().
996 return ctx.run([this, &ctx, callContext](Worker::Lock& lock) mutable {
997 return callImpl(lock, ctx, callContext);
998 });
999 }),
1000 externalPusher(ctx.getExternalPusher()) {}
1001 
1002 struct EnvCtx {
1003 v8::Local<v8::Value> env;
1004 jsg::JsObject ctx;
1005 };
1006 
1007 struct TargetInfo {
1008 // The object on which the RPC method should be invoked.
1009 jsg::JsObject target;
1010 
1011 // If `env` and `ctx` need to be delivered as arguments to the method, these are the values
1012 // to deliver.
1013 kj::Maybe<EnvCtx> envCtx;
1014 
1015 bool allowInstanceProperties;
1016 };
1017 
1018 // Get the object on which the method is to be invoked. This is virtual so that we can have
1019 // separate subclasses handling the case of an entrypoint vs. a transient RPC object.
1020 virtual TargetInfo getTargetInfo(Worker::Lock& lock, IoContext& ioCtx) = 0;
1021 
1022 // Handles the delivery of JS RPC method calls.
1023 kj::Promise<void> call(CallContext callContext) override {
1024 co_await kj::yield();
1025 
1026 // Try to execute the requested method.
1027 co_return co_await enterIsolateAndCall(callContext).catch_([](kj::Exception&& e) {
1028 if (jsg::isTunneledException(e.getDescription())) {
1029 // Annotate exceptions in RPC worker calls as remote exceptions.
1030 auto description = jsg::stripRemoteExceptionPrefix(e.getDescription());
1031 if (!description.startsWith("remote.")) {
1032 // If we already were annotated as remote from some other worker entrypoint, no point
1033 // adding an additional prefix.
1034 e.setDescription(kj::str("remote.", description));
1035 }
1036 }
1037 kj::throwFatalException(kj::mv(e));
1038 });
1039 }
1040 
1041 // Implements ExternalPusher by forwarding to the shared implementation.
1042 //
1043 // Note JsRpcTarget has to implement `ExternalPusher` directly rather than providing a method
1044 // like `getExternalPusher()` because it's important that the pushes arrive before the call, and
1045 // the ordering can only be guaranteed if they're on the same object.
1046 kj::Promise<void> pushByteStream(PushByteStreamContext context) override {
1047 return externalPusher->pushByteStream(context);
1048 }
1049 kj::Promise<void> pushAbortSignal(PushAbortSignalContext context) override {
1050 return externalPusher->pushAbortSignal(context);
1051 }
1052 
1053 KJ_DISALLOW_COPY_AND_MOVE(JsRpcTargetBase);
1054 
1055 private:
1056 virtual void maybeSetJsRpcInfo(IoContext& ctx, const kj::ConstString& methodNameForTrace) = 0;
1057 
1058 // Function which enters the isolate lock and IoContext and then invokes callImpl(). Created
1059 // using IoContext::makeReentryCallback().
1060 kj::Function<kj::Promise<void>(CallContext callContext)> enterIsolateAndCall;
1061 
1062 kj::Rc<ExternalPusherImpl> externalPusher;
1063 
1064 // Returns true if the given name cannot be used as a method on this type.
1065 virtual bool isReservedName(kj::StringPtr name) = 0;
1066 
1067 kj::Promise<void> callImpl(Worker::Lock& lock, IoContext& ctx, CallContext callContext) {
1068 jsg::Lock& js = lock;
1069 auto params = callContext.getParams();
1070 // Method name suitable for use in trace and error messages. May be a pointer into the RPC
1071 // params reader.
1072 kj::ConstString methodNameForTrace;
1073 
1074 // Retrieve the method name and report onset event info if tracing is enabled.
1075 switch (params.which()) {
1076 case rpc::JsRpcTarget::CallParams::METHOD_NAME: {
1077 methodNameForTrace = kj::ConstString(kj::str(params.getMethodName()));
1078 break;
1079 }
1080 case rpc::JsRpcTarget::CallParams::METHOD_PATH: {
1081 auto path = params.getMethodPath();
1082 auto n = path.size();
1083 
1084 if (n == 0) {
1085 // Call the target itself as a function.
1086 methodNameForTrace = "(this)"_kjc;
1087 } else {
1088 methodNameForTrace = kj::ConstString(kj::strArray(path, "."));
1089 }
1090 break;
1091 }
1092 }
1093 
1094 maybeSetJsRpcInfo(ctx, methodNameForTrace);
1095 
1096 auto targetInfo = getTargetInfo(lock, ctx);
1097 
1098 // We will try to get the function, if we can't we'll throw an error to the client.
1099 auto [propHandle, thisArg] =
1100 tryGetProperty(lock, targetInfo.target, params, targetInfo.allowInstanceProperties, ctx);
1101 
1102 auto op = params.getOperation();
1103 
1104 auto handleResult = [&](InvocationResult&& invocationResult) {
1105 // Given a handle for the result, if it's a promise, await the promise, then serialize the
1106 // final result for return.
1107 
1108 RpcSerializerExternalHandler::GetStreamHandlerFunc getResultsStreamHandlerFunc;
1109 auto resultStreamHandler = params.getResultsStreamHandler();
1110 switch (resultStreamHandler.which()) {
1111 case rpc::JsRpcTarget::CallParams::ResultsStreamHandler::EXTERNAL_PUSHER:
1112 getResultsStreamHandlerFunc.init<RpcSerializerExternalHandler::GetExternalPusherFunc>(
1113 [cap = resultStreamHandler.getExternalPusher()]() mutable { return kj::mv(cap); });
1114 break;
1115 case rpc::JsRpcTarget::CallParams::ResultsStreamHandler::STREAM_SINK:
1116 getResultsStreamHandlerFunc.init<RpcSerializerExternalHandler::GetStreamSinkFunc>(
1117 [cap = resultStreamHandler.getStreamSink()]() mutable { return kj::mv(cap); });
1118 break;
1119 }
1120 
1121 kj::Maybe<kj::Own<kj::PromiseFulfiller<rpc::JsRpcTarget::Client>>> callPipelineFulfiller;
1122 
1123 // We need another ref to this fulfiller for the error callback. It can rely on being
1124 // destroyed at the same time as the success callback.
1125 kj::Maybe<kj::PromiseFulfiller<rpc::JsRpcTarget::Client>&> callPipelineFulfillerRef;
1126 
1127 KJ_IF_SOME(ss, invocationResult.streamSink) {
1128 // Since we have a StreamSink, it's important that we hook up the pipeline for that
1129 // immediately. Annoyingly, that also means we need to hook up a pipeline for
1130 // callPipeline, which we don't actually have yet, so we need to promise-ify it.
1131 
1132 // If the caller requested using ExternalPusher for the results, then it should also use
1133 // ExternalPusher for the params. (Theoretically we could support mix-and-match but...
1134 // let's keep it simple.)
1135 KJ_REQUIRE(resultStreamHandler.isStreamSink(),
1136 "RPC params used StreamSink when result is supposed to use ExternalPusher");
1137 
1138 auto paf = kj::newPromiseAndFulfiller<rpc::JsRpcTarget::Client>();
1139 callPipelineFulfillerRef = *paf.fulfiller;
1140 callPipelineFulfiller = kj::mv(paf.fulfiller);
1141 
1142 capnp::PipelineBuilder<rpc::JsRpcTarget::CallResults> builder(16);
1143 builder.setCallPipeline(kj::mv(paf.promise));
1144 builder.setParamsStreamSink(ss);
1145 callContext.setPipeline(builder.build());
1146 }
1147 
1148 // HACK: Cap'n Proto call contexts are documented as being pointer-like types where the
1149 // backing object's lifetime is that of the RPC call, but in reality they are refcounted
1150 // under the hood. Since we'll be executing the call in the JS microtask queue, we have no
1151 // ability to actually cancel execution if a cancellation arrives over RPC, and at the end of
1152 // that execution we're going to access the call context to write the results. We could
1153 // invent some complicated way to skip initializing results in the case the call has been
1154 // canceled, but it's easier and safer to just grab a refcount on the call context object
1155 // itself, which fully protects us. So... do that.
1156 auto ownCallContext = capnp::CallContextHook::from(callContext).addRef();
1157 
1158 auto result = ctx.awaitJs(js,
1159 js.toPromise(invocationResult.returnValue)
1160 .then(js,
1161 ctx.addFunctor(
1162 // Warning: Be careful about captures here! If the incoming RPC is canceled,
1163 // this continuation will still execute, sice it's a JS promise continuation.
1164 // But `this` could have been destroyed in the meantime. So all our captures
1165 // must take full ownership.
1166 [callContext, ownCallContext = kj::mv(ownCallContext),
1167 paramDisposalGroup = kj::mv(invocationResult.paramDisposalGroup),
1168 paramsStreamSink = kj::mv(invocationResult.streamSink),
1169 getResultsStreamHandlerFunc = kj::mv(getResultsStreamHandlerFunc),
1170 callPipelineFulfiller = kj::mv(callPipelineFulfiller)](
1171 jsg::Lock& js, jsg::Value value) mutable {
1172 jsg::JsValue resultValue(value.getHandle(js));
1173 
1174 rpc::JsRpcTarget::CallResults::Builder results = nullptr;
1175 auto maybePipeline =
1176 serializeJsValueWithPipeline(js, resultValue, [&](capnp::MessageSize hint) {
1177 hint.wordCount += capnp::sizeInWords<rpc::JsRpcTarget::CallResults>();
1178 hint.capCount += 1; // for callPipeline
1179 results = callContext.initResults(hint);
1180 return results.initResult();
1181 }, kj::mv(getResultsStreamHandlerFunc));
1182 
1183 KJ_SWITCH_ONEOF(maybePipeline) {
1184 KJ_CASE_ONEOF(obj, MakeCallPipeline::Object) {
1185 results.setCallPipeline(kj::mv(obj.cap));
1186 
1187 // Note that hasDisposer is ONLY meant to indicate the presence of an
1188 // application-level disposer. It need not be true if we only have stub disposers.
1189 results.setHasDisposer(obj.hasDispose);
1190 }
1191 KJ_CASE_ONEOF(obj, MakeCallPipeline::SingleStub) {
1192 // Serialization should have produced a single stub. We can use that same stub as
1193 // the callPipeline.
1194 auto externals = results.asReader().getResult().getExternals();
1195 KJ_ASSERT(externals.size() == 1);
1196 auto external = externals[0];
1197 KJ_ASSERT(external.isRpcTarget());
1198 results.setCallPipeline(external.getRpcTarget());
1199 }
1200 KJ_CASE_ONEOF(nonPipelinable, MakeCallPipeline::NonPipelinable) {
1201 results.setCallPipeline(kj::mv(nonPipelinable.errorPipeline));
1202 // leave hasDisposer false
1203 }
1204 }
1205 
1206 KJ_IF_SOME(cpf, callPipelineFulfiller) {
1207 cpf->fulfill(results.getCallPipeline());
1208 }
1209 
1210 KJ_IF_SOME(ss, paramsStreamSink) {
1211 results.setParamsStreamSink(kj::mv(ss));
1212 }
1213 
1214 // paramDisposalGroup will be destroyed when we return (or when this lambda is destroyed
1215 // as a result of the promise being rejected). This will implicitly dispose the param
1216 // stubs.
1217 }),
1218 ctx.addFunctor([callPipelineFulfillerRef](jsg::Lock& js, jsg::Value&& error) {
1219 // If we set up a `callPipeline` early, we have to make sure it propagates the error.
1220 // (Otherwise we get a PromiseFulfiller error instead, which is pretty useless...)
1221 KJ_IF_SOME(cpf, callPipelineFulfillerRef) {
1222 cpf.reject(js.exceptionToKj(error.addRef(js)));
1223 }
1224 js.throwException(kj::mv(error));
1225 })));
1226 
1227 if (ctx.hasOutputGate()) {
1228 // Note: If `ctx` is destroyed, the entire call to `callImpl()` will be canceled
1229 // (makeReentryCallback() ensures this). This does NOT cancel the JavaScript (because JS
1230 // promises are not RAII-cancelable), but it will cancel this trailing .then(), which is
1231 // why it's safe to capture `&ctx` here.
1232 return result.then([&ctx]() mutable { return ctx.waitForOutputLocks(); });
1233 } else {
1234 return result;
1235 }
1236 };
1237 
1238 switch (op.which()) {
1239 case rpc::JsRpcTarget::CallParams::Operation::CALL_WITH_ARGS: {
1240 // Note that using isFunctionForRpc(js, propHandle) here would be incorrect, since that
1241 // decides whether it is a function *that can be serialized as a stub*. JsRpcProperty
1242 // is (at present) considered non-serializable in itself, but when traversing the
1243 // pipeline path, we may have descended into a stub and its properties, thus we could
1244 // actually be invoking a JsRpcProperty here. As long as it is in fact callable, we will
1245 // allow it.
1246 JSG_REQUIRE(propHandle->IsFunction(), TypeError,
1247 kj::str("\"", methodNameForTrace, "\" is not a function."));
1248 auto fn = propHandle.As<v8::Function>();
1249 
1250 kj::Maybe<rpc::JsValue::Reader> args;
1251 if (op.hasCallWithArgs()) {
1252 args = op.getCallWithArgs();
1253 }
1254 
1255 InvocationResult invocationResult;
1256 KJ_IF_SOME(envCtx, targetInfo.envCtx) {
1257 invocationResult = invokeFnInsertingEnvCtx(
1258 js, methodNameForTrace, fn, thisArg, args, envCtx.env, envCtx.ctx);
1259 } else {
1260 invocationResult = invokeFn(js, fn, thisArg, args);
1261 }
1262 
1263 // We have a function, so let's call it and serialize the result for RPC.
1264 // If the function returns a promise we will wait for the promise to finish so we can
1265 // serialize the result.
1266 return handleResult(kj::mv(invocationResult));
1267 }
1268 
1269 case rpc::JsRpcTarget::CallParams::Operation::GET_PROPERTY:
1270 return handleResult({.returnValue = propHandle});
1271 }
1272 
1273 KJ_FAIL_ASSERT("unknown JsRpcTarget::CallParams::Operation", (uint)op.which());
1274 }
1275 
1276 struct GetPropResult {
1277 v8::Local<v8::Value> handle;
1278 v8::Local<v8::Object> thisArg;
1279 };
1280 
1281 [[noreturn]] static void failLookup(kj::StringPtr kjName) {
1282 JSG_FAIL_REQUIRE(
1283 TypeError, kj::str("The RPC receiver does not implement the method \"", kjName, "\"."));
1284 }
1285 
1286 GetPropResult tryGetProperty(jsg::Lock& js,
1287 jsg::JsObject object,
1288 rpc::JsRpcTarget::CallParams::Reader callParams,
1289 bool allowInstanceProperties,
1290 IoContext& ctx) {
1291 auto prototypeOfObject = KJ_ASSERT_NONNULL(js.obj().getPrototype(js).tryCast<jsg::JsObject>());
1292 
1293 // Get the named property of `object`.
1294 auto getProperty = [&](kj::StringPtr kjName) {
1295 JSG_REQUIRE(!isReservedName(kjName), TypeError,
1296 kj::str("'", kjName, "' is a reserved method and cannot be called over RPC."));
1297 
1298 jsg::JsValue jsName = js.strIntern(kjName);
1299 
1300 if (allowInstanceProperties) {
1301 // This is a simple object. Its own properties are considered to be accessible over RPC, but
1302 // inherited properties (i.e. from Object.prototype) are not.
1303 if (!object.has(js, jsName, jsg::JsObject::HasOption::OWN)) {
1304 failLookup(kjName);
1305 }
1306 return object.get(js, jsName);
1307 } else {
1308 // This is an instance of a valid RPC target class.
1309 if (object.has(js, jsName, jsg::JsObject::HasOption::OWN)) {
1310 // We do NOT allow own properties, only class properties.
1311 failLookup(kjName);
1312 }
1313 
1314 auto value = object.get(js, jsName);
1315 if (value == prototypeOfObject.get(js, jsName)) {
1316 // This property is inherited from the prototype of `Object`. Don't allow.
1317 failLookup(kjName);
1318 }
1319 
1320 return value;
1321 }
1322 };
1323 
1324 kj::Maybe<jsg::JsValue> result;
1325 
1326 switch (callParams.which()) {
1327 case rpc::JsRpcTarget::CallParams::METHOD_NAME: {
1328 result = getProperty(callParams.getMethodName());
1329 break;
1330 }
1331 
1332 case rpc::JsRpcTarget::CallParams::METHOD_PATH: {
1333 auto path = callParams.getMethodPath();
1334 auto n = path.size();
1335 
1336 if (n == 0) {
1337 // Call the target itself as a function.
1338 result = object;
1339 } else {
1340 bool inStub = false;
1341 for (auto i: kj::zeroTo(n - 1)) {
1342 // For each property name except the last, look up the property and replace `object`
1343 // with it.
1344 kj::StringPtr name = path[i];
1345 auto next = getProperty(name);
1346 
1347 KJ_IF_SOME(o, next.tryCast<jsg::JsObject>()) {
1348 object = o;
1349 } else {
1350 // Not an object, doesn't have further properties.
1351 failLookup(name);
1352 }
1353 
1354 // If the object is a Proxy, then `isInstanceOf<JsRpcTarget>()` won't actually work,
1355 // because the Proxy is not an instance of any native type. But for our purposes,
1356 // RpcTarget is only a marker used to indicate what semantics are desired.
1357 bool isProxyOfRpcTarget = false;
1358 if (jsg::JsValue(object).isProxy()) {
1359 // Unfortunatley in this case we need to follow the prototype chain manually, looking
1360 // for `JsRpcTarget`.
1361 js.withinHandleScope([&]() {
1362 auto proto = object.getPrototype(js);
1363 auto prototypeOfRpcTarget = js.getPrototypeFor<JsRpcTarget>();
1364 
1365 for (;;) {
1366 auto objProto = KJ_UNWRAP_OR(proto.tryCast<jsg::JsObject>(), break);
1367 if (objProto == prototypeOfRpcTarget) {
1368 isProxyOfRpcTarget = true;
1369 break;
1370 }
1371 proto = objProto.getPrototype(js);
1372 }
1373 });
1374 }
1375 
1376 // Decide whether the new object is a suitable RPC target.
1377 if (object.getPrototype(js) == prototypeOfObject) {
1378 // Yes. It's a simple object.
1379 allowInstanceProperties = true;
1380 } else if (isProxyOfRpcTarget || object.isInstanceOf<JsRpcTarget>(js)) {
1381 // Yes. It's a JsRpcTarget.
1382 allowInstanceProperties = false;
1383 } else if (object.isInstanceOf<JsRpcStub>(js) || object.isInstanceOf<Fetcher>(js) ||
1384 (inStub && object.isInstanceOf<JsRpcProperty>(js))) {
1385 // Yes. It's a JsRpcStub or Fetcher. We should allow descending into the stub.
1386 // Note that the wildcard property of a stub is a prototype property, not an instance
1387 // property, so setting allowInstanceProperties = false here gets the behavior we
1388 // want.
1389 // TODO(someday): We'll need to support JsRpcPromise here if someday we allow it to
1390 // be serialized.
1391 allowInstanceProperties = false;
1392 
1393 // We will only traverse JsRpcProperty if we got there by descending through a
1394 // JsRpcStub. At present you can't just pull a property of a stub and return it.
1395 inStub = true;
1396 } else if (isFunctionForRpc(js, object)) {
1397 // Yes. It's a function.
1398 allowInstanceProperties = true;
1399 } else {
1400 failLookup(name);
1401 }
1402 }
1403 
1404 result = getProperty(path[n - 1]);
1405 }
1406 
1407 break;
1408 }
1409 }
1410 
1411 return {
1412 .handle = KJ_ASSERT_NONNULL(result, "unknown CallParams type", (uint)callParams.which()),
1413 .thisArg = object,
1414 };
1415 }
1416 
1417 struct InvocationResult {
1418 v8::Local<v8::Value> returnValue;
1419 kj::Maybe<kj::Own<RpcStubDisposalGroup>> paramDisposalGroup;
1420 kj::Maybe<rpc::JsValue::StreamSink::Client> streamSink;
1421 };
1422 
1423 // Deserializes the arguments and passes them to the given function.
1424 static InvocationResult invokeFn(jsg::Lock& js,
1425 v8::Local<v8::Function> fn,
1426 v8::Local<v8::Object> thisArg,
1427 kj::Maybe<rpc::JsValue::Reader> args) {
1428 // We received arguments from the client, deserialize them back to JS.
1429 KJ_IF_SOME(a, args) {
1430 auto [value, disposalGroup, streamSink] = deserializeJsValue(js, a, "params"_kjc);
1431 auto args = KJ_REQUIRE_NONNULL(
1432 value.tryCast<jsg::JsArray>(), "expected JsArray when deserializing arguments.");
1433 // Call() expects a `Local<Value> []`... so we populate an array.
1434 
1435 v8::LocalVector<v8::Value> arguments(js.v8Isolate, args.size());
1436 for (size_t i = 0; i < args.size(); ++i) {
1437 arguments[i] = args.get(js, i);
1438 }
1439 
1440 InvocationResult result{
1441 .returnValue =
1442 jsg::check(fn->Call(js.v8Context(), thisArg, arguments.size(), arguments.data())),
1443 .streamSink = kj::mv(streamSink),
1444 };
1445 if (!disposalGroup->empty()) {
1446 result.paramDisposalGroup = kj::mv(disposalGroup);
1447 }
1448 return result;
1449 } else {
1450 return {.returnValue = jsg::check(fn->Call(js.v8Context(), thisArg, 0, nullptr))};
1451 }
1452 };
1453 
1454 // Like `invokeFn`, but inject the `env` and `ctx` values between the first and second
1455 // parameters. Used for service bindings that use functional syntax.
1456 static InvocationResult invokeFnInsertingEnvCtx(jsg::Lock& js,
1457 kj::StringPtr methodName,
1458 v8::Local<v8::Function> fn,
1459 v8::Local<v8::Object> thisArg,
1460 kj::Maybe<rpc::JsValue::Reader> args,
1461 v8::Local<v8::Value> env,
1462 jsg::JsObject ctx) {
1463 // Determine the function arity (how many parameters it was declared to accept) by reading the
1464 // `.length` attribute.
1465 auto arity = js.withinHandleScope([&]() {
1466 auto length = jsg::check(fn->Get(js.v8Context(), js.strIntern("length")));
1467 return jsg::check(length->IntegerValue(js.v8Context()));
1468 });
1469 
1470 // Avoid excessive allocation from a maliciously-set `length`.
1471 JSG_REQUIRE(arity >= 0 && arity < 256, TypeError,
1472 "RPC function has unreasonable length attribute: ", arity);
1473 
1474 if (arity < 3) {
1475 // If a function has fewer than three arguments, reproduce the historical behavior where
1476 // we'd pass the main argument followed by `env` and `ctx` and the undeclared parameters
1477 // would just be truncated.
1478 arity = 3;
1479 }
1480 
1481 kj::Maybe<kj::Own<RpcStubDisposalGroup>> paramDisposalGroup;
1482 kj::Maybe<rpc::JsValue::StreamSink::Client> streamSink;
1483 
1484 // We're going to pass all the arguments from the client to the function, but we are going to
1485 // insert `env` and `ctx`. We assume the last two arguments that the function declared are
1486 // `env` and `ctx`, so we can determine where to insert them based on the function's arity.
1487 kj::Maybe<jsg::JsArray> argsArrayFromClient;
1488 size_t argCountFromClient = 0;
1489 KJ_IF_SOME(a, args) {
1490 auto [value, disposalGroup, ss] = deserializeJsValue(js, a, "paramsNonClass"_kjc);
1491 streamSink = kj::mv(ss);
1492 
1493 auto array = KJ_REQUIRE_NONNULL(
1494 value.tryCast<jsg::JsArray>(), "expected JsArray when deserializing arguments.");
1495 argCountFromClient = array.size();
1496 argsArrayFromClient = kj::mv(array);
1497 
1498 if (!disposalGroup->empty()) {
1499 paramDisposalGroup = kj::mv(disposalGroup);
1500 }
1501 }
1502 
1503 // For now, we are disallowing multiple arguments with bare function syntax, due to a footgun:
1504 // if you forget to add `env, ctx` to your arg list, then the last arguments from the client
1505 // will be replaced with `env` and `ctx`. Probably this would be quickly noticed in testing,
1506 // but if you were to accidentally reflect `env` back to the client, it would be a severe
1507 // security flaw.
1508 JSG_REQUIRE(arity == 3, TypeError, "Cannot call handler function \"", methodName,
1509 "\" over RPC because it has the wrong "
1510 "number of arguments. A simple function handler can only be called over RPC if it has "
1511 "exactly the arguments (arg, env, ctx), where only the first argument comes from the "
1512 "client. To support multi-argument RPC functions, use class-based syntax (extending "
1513 "WorkerEntrypoint) instead.");
1514 JSG_REQUIRE(argCountFromClient == 1, TypeError, "Attempted to call RPC function \"", methodName,
1515 "\" with the wrong number of arguments. "
1516 "When calling a top-level handler function that is not declared as part of a class, you "
1517 "must always send exactly one argument. In order to support variable numbers of "
1518 "arguments, the server must use class-based syntax (extending WorkerEntrypoint) "
1519 "instead.");
1520 
1521 v8::LocalVector<v8::Value> arguments(js.v8Isolate, kj::max(argCountFromClient + 2, arity));
1522 
1523 for (auto i: kj::zeroTo(arity - 2)) {
1524 if (argCountFromClient > i) {
1525 arguments[i] = KJ_ASSERT_NONNULL(argsArrayFromClient).get(js, i);
1526 } else {
1527 arguments[i] = js.undefined();
1528 }
1529 }
1530 
1531 arguments[arity - 2] = env;
1532 arguments[arity - 1] = ctx;
1533 
1534 KJ_IF_SOME(a, argsArrayFromClient) {
1535 for (size_t i = arity - 2; i < argCountFromClient; ++i) {
1536 arguments[i + 2] = a.get(js, i);
1537 }
1538 }
1539 
1540 return {
1541 .returnValue =
1542 jsg::check(fn->Call(js.v8Context(), thisArg, arguments.size(), arguments.data())),
1543 .paramDisposalGroup = kj::mv(paramDisposalGroup),
1544 .streamSink = kj::mv(streamSink),
1545 };
1546 };
1547};
1548 
1549class TransientJsRpcTarget final: public JsRpcTargetBase {
1550 public:
1551 TransientJsRpcTarget(
1552 jsg::Lock& js, IoContext& ioCtx, jsg::JsObject object, bool allowInstanceProperties = false)
1553 : JsRpcTargetBase(ioCtx, MayOutliveIncomingRequest()),
1554 handles(ioCtx.addObjectReverse(kj::heap<Handles>(js, object))),
1555 allowInstanceProperties(allowInstanceProperties) {
1556 // Check for the existence of a dispose function now so that the destructor doesn't have to
1557 // take an isolate lock if there isn't one.
1558 auto getResult = object.get(js, js.symbolDispose());
1559 if (getResult.isFunction()) {
1560 auto dispose = jsg::V8Ref<v8::Function>(
1561 js.v8Isolate, v8::Local<v8::Value>(getResult).As<v8::Function>());
1562 disposeFulfiller = addDisposeTask(js, ioCtx, object, kj::mv(dispose), {});
1563 }
1564 }
1565 
1566 // Use this version of the constructor to pass the dispose function separately.
1567 TransientJsRpcTarget(jsg::Lock& js,
1568 IoContext& ioCtx,
1569 jsg::JsObject object,
1570 kj::Maybe<jsg::V8Ref<v8::Function>> dispose,
1571 kj::Vector<kj::Own<void>> stubDisposers,
1572 bool allowInstanceProperties = false)
1573 : JsRpcTargetBase(ioCtx, MayOutliveIncomingRequest()),
1574 handles(ioCtx.addObjectReverse(kj::heap<Handles>(js, object))),
1575 disposeFulfiller(addDisposeTask(js, ioCtx, object, kj::mv(dispose), kj::mv(stubDisposers))),
1576 allowInstanceProperties(allowInstanceProperties) {}
1577 
1578 ~TransientJsRpcTarget() noexcept(false) {
1579 KJ_IF_SOME(f, kj::mv(disposeFulfiller)) {
1580 f->fulfill();
1581 }
1582 }
1583 
1584 TargetInfo getTargetInfo(Worker::Lock& lock, IoContext& ioCtx) override {
1585 return {
1586 .target = handles->object.getHandle(lock),
1587 .envCtx = kj::none,
1588 .allowInstanceProperties = allowInstanceProperties,
1589 };
1590 }
1591 
1592 private:
1593 struct Handles {
1594 jsg::JsRef<jsg::JsObject> object;
1595 
1596 Handles(jsg::Lock& js, jsg::JsObject object): object(js, object) {}
1597 };
1598 
1599 // This object could outlive the IoContext (that's why `JsRpcTargetBase` holds a `WeakRef` to the
1600 // context). That means hypothetically it could also outlive the isolate. We therefore need to
1601 // place these handles in a `ReverseIoOwn` so that if the `IoContext` dies before we do, they are
1602 // dropped at that point.
1603 ReverseIoOwn<Handles> handles;
1604 
1605 // When fulfilled, calls the original object's dispose function.
1606 kj::Maybe<kj::Own<kj::PromiseFulfiller<void>>> disposeFulfiller;
1607 
1608 static kj::Maybe<kj::Own<kj::PromiseFulfiller<void>>> addDisposeTask(jsg::Lock& js,
1609 IoContext& ctx,
1610 jsg::JsObject object,
1611 kj::Maybe<jsg::V8Ref<v8::Function>> dispose,
1612 kj::Vector<kj::Own<void>> stubDisposers) {
1613 if (dispose == kj::none && stubDisposers.empty()) {
1614 // Don't bother scheduling disposal if we have neither.
1615 return kj::none;
1616 }
1617 
1618 auto obj = jsg::JsRef<jsg::JsObject>(js, object);
1619 auto [promise, fulfiller] = kj::newPromiseAndFulfiller<void>();
1620 auto jsPromise = ctx.awaitIo(js, kj::mv(promise),
1621 [obj = kj::mv(obj), dispose = kj::mv(dispose), stubDiposers = kj::mv(stubDisposers)](
1622 jsg::Lock& js) {
1623 KJ_IF_SOME(d, dispose) {
1624 jsg::check(d.getHandle(js)->Call(js.v8Context(), obj.getHandle(js), 0, nullptr));
1625 }
1626 
1627 // Our stub disposers are dropped at the end of this task.
1628 });
1629 ctx.addTask(ctx.awaitJs(js, kj::mv(jsPromise)));
1630 return kj::mv(fulfiller);
1631 }
1632 
1633 bool allowInstanceProperties;
1634 
1635 bool isReservedName(kj::StringPtr name) override {
1636 if ( // dup() is reserved to duplicate the stub itself, pointing to the same object.
1637 name == "dup" ||
1638 
1639 // All JS classes define a method `constructor` on the prototype, but we don't actually
1640 // want this to be callable over RPC!
1641 name == "constructor") {
1642 return true;
1643 }
1644 return false;
1645 }
1646 
1647 void maybeSetJsRpcInfo(IoContext& ctx, const kj::ConstString& methodNameForTrace) override {}
1648};
1649 
1650// See comment at call site for explanation.
1651static rpc::JsRpcTarget::Client makeJsRpcTargetForSingleLoopbackCall(
1652 jsg::Lock& js, jsg::JsObject obj) {
1653 // We intentionally do not want to hook up the disposer here since we're not taking ownership
1654 // of the object.
1655 return rpc::JsRpcTarget::Client(kj::heap<TransientJsRpcTarget>(
1656 js, IoContext::current(), obj, kj::none, kj::Vector<kj::Own<void>>(), true));
1657}
1658 
1659template <typename Func>
1660MakeCallPipeline::Result serializeJsValueWithPipeline(jsg::Lock& js,
1661 jsg::JsValue value,
1662 Func makeBuilder,
1663 RpcSerializerExternalHandler::GetStreamHandlerFunc getStreamHandlerFunc) {
1664 auto maybeDispose = js.withinHandleScope([&]() -> kj::Maybe<jsg::V8Ref<v8::Function>> {
1665 jsg::JsObject obj = KJ_UNWRAP_OR(value.tryCast<jsg::JsObject>(), { return kj::none; });
1666 
1667 if (obj.getPrototype(js) == js.obj().getPrototype(js)) {
1668 // It's a plain object.
1669 jsg::JsValue disposeProperty = obj.get(js, js.symbolDispose());
1670 
1671 // We don't want the disposer to be serialized, so delete it from the object. (Remember
1672 // that a new `dispose()` method will always be added on the client side).
1673 obj.delete_(js, js.symbolDispose());
1674 
1675 if (disposeProperty.isFunction()) {
1676 auto localDispose = v8::Local<v8::Value>(disposeProperty).As<v8::Function>();
1677 return jsg::V8Ref<v8::Function>(js.v8Isolate, localDispose);
1678 }
1679 }
1680 
1681 return kj::none;
1682 });
1683 auto hasDispose = maybeDispose != kj::none;
1684 
1685 // Now that we've extracted our dispose function, we can serialize our value.
1686 RpcSerializerExternalHandler externalHandler(
1687 RpcSerializerExternalHandler::TRANSFER, kj::mv(getStreamHandlerFunc));
1688 serializeJsValue(js, value, externalHandler, kj::mv(makeBuilder));
1689 
1690 auto stubDisposers = externalHandler.releaseStubDisposers();
1691 
1692 return js.withinHandleScope([&]() -> MakeCallPipeline::Result {
1693 jsg::JsObject obj = KJ_UNWRAP_OR(value.tryCast<jsg::JsObject>(), {
1694 // Primitive value. Return a fake pipeline just so that we get nice errors if someone tries
1695 // to pipeline on it. (If we return null, we'll get "called null capability" out of
1696 // Cap'n Proto, which will be treated as an internal error.)
1697 return MakeCallPipeline::NonPipelinable{
1698 .errorPipeline = rpc::JsRpcTarget::Client(kj::heap<TransientJsRpcTarget>(
1699 js, IoContext::current(), js.obj(), kj::none, kj::Vector<kj::Own<void>>(), true))};
1700 });
1701 
1702 if (obj.getPrototype(js) == js.obj().getPrototype(js)) {
1703 // It's a plain object.
1704 auto pipeline = kj::heap<TransientJsRpcTarget>(
1705 js, IoContext::current(), obj, kj::mv(maybeDispose), kj::mv(stubDisposers), true);
1706 
1707 return MakeCallPipeline::Object{
1708 .cap = rpc::JsRpcTarget::Client(kj::mv(pipeline)), .hasDispose = hasDispose};
1709 } else if (obj.isInstanceOf<JsRpcStub>(js)) {
1710 // It's just a stub. It'll serialize as a single stub, obviously.
1711 return MakeCallPipeline::SingleStub();
1712 } else if (obj.isInstanceOf<JsRpcTarget>(js)) {
1713 // It's an RPC target. It will be serialized as a single stub.
1714 return MakeCallPipeline::SingleStub();
1715 } else if (isFunctionForRpc(js, obj)) {
1716 // It's a plain function. It will be serialized as a single stub.
1717 return MakeCallPipeline::SingleStub();
1718 } else if (obj.isInstanceOf<Fetcher>(js)) {
1719 // It's a plain fetcher. We want to allow pipelining on it, but we also actually need to
1720 // serialize it, so we can't use `SingleStub()`. Note we set `allowInstanceProperties` to
1721 // `false` here because the wildcard property of a `Fetcher` is a prototype property, and
1722 // that's what we want to expose for pipelining.
1723 auto pipeline = kj::heap<TransientJsRpcTarget>(
1724 js, IoContext::current(), obj, kj::mv(maybeDispose), kj::mv(stubDisposers), false);
1725 
1726 return MakeCallPipeline::Object{
1727 .cap = rpc::JsRpcTarget::Client(kj::mv(pipeline)), .hasDispose = hasDispose};
1728 } else {
1729 // Not an RPC object. Could be a String or other serializable types that derive from Object.
1730 // Similar to primitive types, we return a fake pipeline for error-handling reasons.
1731 // TODO(soon): What if someone returns e.g. a Map with a disposer on it? Should we honor that
1732 // disposer?
1733 return MakeCallPipeline::NonPipelinable{
1734 .errorPipeline = rpc::JsRpcTarget::Client(kj::heap<TransientJsRpcTarget>(
1735 js, IoContext::current(), js.obj(), kj::none, kj::Vector<kj::Own<void>>(), true))};
1736 }
1737 });
1738}
1739 
1740// RpcStub are allowed to wrap:
1741// * RpcTargets
1742// * Functions
1743// * Plain objects (only when created explicitly via `new RpcStub`)
1744//
1745// This function checks for these and returns:
1746// * kj::none if it's not a valid type to be wrapped in as tub.
1747// * The value for allowInstanceProperties if it is.
1748kj::Maybe<bool> checkStubType(jsg::Lock& js, jsg::JsObject handle) {
1749 return js.withinHandleScope([&]() -> kj::Maybe<bool> {
1750 // TODO(perf): We should really cache `prototypeOfObject` somewhere so we don't have to create
1751 // an object to get it. (We do this other places in this file, too...)
1752 auto prototypeOfObject = KJ_ASSERT_NONNULL(js.obj().getPrototype(js).tryCast<jsg::JsObject>());
1753 auto prototypeOfRpcTarget = js.getPrototypeFor<JsRpcTarget>();
1754 auto proto = handle.getPrototype(js);
1755 if (proto == prototypeOfObject) {
1756 // A regular object. Allow access to instance properties.
1757 return true;
1758 } else {
1759 // Walk the prototype chain looking for RpcTarget.
1760 //
1761 // (Note we can't simply use handle.isInstanceOf<JsRpcTarget>() because that doesn't work
1762 // correctly for proxies. Since RpcTarget is only used as a marker, we don't really need
1763 // the object to be an instance of it -- we just care if it's in the prototype chain, even
1764 // if the prototype chain is faked by the Proxy.)
1765 //
1766 // TODO(someday): Consider whether `new RpcStub(obj)` should work on arbitrary types. This
1767 // could be a useful way to say: "I am explicitly opting into treating this like an
1768 // RpcTarget even though I do not have the ability to make its type extend RpcTarget."
1769 for (;;) {
1770 if (proto == prototypeOfRpcTarget) {
1771 // An RpcTarget, don't allow instance properties.
1772 return false;
1773 }
1774 
1775 KJ_IF_SOME(protoObj, proto.tryCast<jsg::JsObject>()) {
1776 proto = protoObj.getPrototype(js);
1777 } else if (isFunctionForRpc(js, handle)) {
1778 // This is NOT an RpcTarget, but it IS callable as a function, so treat it as such.
1779 return true;
1780 } else {
1781 // End of prototype chain, and didn't find RpcTarget.
1782 return kj::none;
1783 }
1784 }
1785 }
1786 });
1787}
1788 
1789jsg::Ref<JsRpcStub> JsRpcStub::constructor(jsg::Lock& js, jsg::JsObject object) {
1790 auto& ioctx = IoContext::current();
1791 
1792 bool allowInstanceProperties = JSG_REQUIRE_NONNULL(checkStubType(js, object), TypeError,
1793 "RpcStubs can only wrap plain objects, functions, and RpcTarget derivatives.");
1794 
1795 rpc::JsRpcTarget::Client cap =
1796 kj::heap<TransientJsRpcTarget>(js, ioctx, object, allowInstanceProperties);
1797 
1798 return js.alloc<JsRpcStub>(ioctx.addObject(kj::heap(kj::mv(cap))));
1799}
1800 
1801void JsRpcTarget::serialize(jsg::Lock& js, jsg::Serializer& serializer) {
1802 // Serialize by effectively creating a `JsRpcStub` around this object and serializing that.
1803 // Except we don't actually want to do _exactly_ that, because we do not want to actually create
1804 // a `JsRpcStub` locally. So do the important parts of `JsRpcStub::constructor()` followed by
1805 // `JsRpcStub::serialize()`.
1806 
1807 auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError,
1808 "Remote RPC references can only be serialized for RPC.");
1809 auto externalHandler = dynamic_cast<RpcSerializerExternalHandler*>(&handler);
1810 JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError,
1811 "Remote RPC references can only be serialized for RPC.");
1812 
1813 // Handle can't possibly be missing during serialization, it's how we got here.
1814 auto handle = jsg::JsObject(KJ_ASSERT_NONNULL(JSG_THIS.tryGetHandle(js)));
1815 
1816 if (externalHandler->getStubOwnership() == RpcSerializerExternalHandler::DUPLICATE) {
1817 // This message isn't supposed to take ownership of stubs. What does that mean for an
1818 // RpcTarget? You might argue that it means we should never call the disposer. But that's not
1819 // really enough: what if the real owner *does* call the disposer, before our stub is done
1820 // with it? How do we make sure the RpcTarget stays alive?
1821 //
1822 // Things get clearer if we look at a real use case: pure-JS Cap'n Web stubs. We don't see
1823 // them as stubs (since they are not instances of JsRpcStub). Instead, we see them as
1824 // RpcTargets. But we need the semantics to come out the same: when passed as a parameter
1825 // to a native RPC call, we need to duplicate the stub, because the original copy might very
1826 // well be disposed before we use it.
1827 //
1828 // How do we duplicate this non-native stub? Well... proper way to duplicate a pure-JS Cap'n
1829 // Web stub is, of course, to call its `dup()` method.
1830 //
1831 // So how about we just do that? If the target has a `dup()` method, we call it, and we take
1832 // ownership of the result, instead of taking ownership of the original object.
1833 auto dup = handle.get(js, "dup");
1834 KJ_IF_SOME(dupFunc, dup.tryCast<jsg::JsFunction>()) {
1835 auto replacement = dupFunc.call(js, handle);
1836 bool replaced = false;
1837 
1838 // We got a duplicate. Is it still an RpcTarget?
1839 KJ_IF_SOME(replacementObj, replacement.tryCast<jsg::JsObject>()) {
1840 if (replacementObj.isInstanceOf<JsRpcTarget>(js)) {
1841 // It is! Let's replace our handle with the duplicate!
1842 handle = replacementObj;
1843 replaced = true;
1844 }
1845 }
1846 
1847 JSG_REQUIRE(replaced, DOMDataCloneError,
1848 "Couldn't create a stub for the RcpTarget because it has a dup() method which did not "
1849 "return another RpcTarget. Either remove the dup() method or make sure it returns an "
1850 "RpcTarget.");
1851 } else {
1852 // If no dup() method was present, then what?
1853 //
1854 // The pedantic argument would say: we need to throw an exception. But that would lead to a
1855 // pretty poor development experience as people would have to mess with adding dup()
1856 // methods to all their RpcTargets.
1857 //
1858 // Another argument might say: we should just use the RpcTarget but never call the disposer
1859 // since we don't own it. But that would probably be confusing. People would wonder why their
1860 // disposers are never called.
1861 //
1862 // If someone passes an RpcTarget with no dup() method, but which does have a disposer, as
1863 // the argument to an RPC method, *probably* they just want the disposer to be called when
1864 // the callee is done with the object. That is, they want us to take ownership after all. If
1865 // that is *not* what they want, then they can always implement a dup() method to make it
1866 // clear.
1867 //
1868 // So, we will just "take ownership" of the target after all, and call its disposer.
1869 }
1870 }
1871 
1872 rpc::JsRpcTarget::Client cap = kj::heap<TransientJsRpcTarget>(js, IoContext::current(), handle);
1873 
1874 externalHandler->write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable {
1875 builder.setRpcTarget(kj::mv(cap));
1876 });
1877}
1878 
1879void RpcSerializerExternalHandler::serializeFunction(
1880 jsg::Lock& js, jsg::Serializer& serializer, v8::Local<v8::Function> func) {
1881 serializer.writeRawUint32(static_cast<uint>(rpc::SerializationTag::JS_RPC_STUB));
1882 
1883 auto handle = jsg::JsObject(func);
1884 
1885 // Similar to JsRpcTarget::serialize(), we may need to dup() the function.
1886 if (stubOwnership == RpcSerializerExternalHandler::DUPLICATE) {
1887 auto dup = handle.get(js, "dup");
1888 KJ_IF_SOME(dupFunc, dup.tryCast<jsg::JsFunction>()) {
1889 auto replacement = dupFunc.call(js, handle);
1890 bool replaced = false;
1891 
1892 // We got a duplicate. Is it still a Function?
1893 KJ_IF_SOME(replacementObj, replacement.tryCast<jsg::JsObject>()) {
1894 if (isFunctionForRpc(js, replacementObj)) {
1895 // It is! Let's replace our handle with the duplicate!
1896 handle = replacementObj;
1897 replaced = true;
1898 }
1899 }
1900 
1901 JSG_REQUIRE(replaced, DOMDataCloneError,
1902 "Couldn't create a stub for the function because it has a dup() method which did not "
1903 "return another function. Either remove the dup() method or make sure it returns a "
1904 "function.");
1905 }
1906 }
1907 
1908 rpc::JsRpcTarget::Client cap =
1909 kj::heap<TransientJsRpcTarget>(js, IoContext::current(), handle, true);
1910 write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable {
1911 builder.setRpcTarget(kj::mv(cap));
1912 });
1913}
1914 
1915void RpcSerializerExternalHandler::serializeProxy(
1916 jsg::Lock& js, jsg::Serializer& serializer, v8::Local<v8::Proxy> proxy) {
1917 auto handle = jsg::JsObject(proxy);
1918 
1919 // Proxies are allowed to present themselves as anything that you could pass to `new RpcStub`.
1920 //
1921 // Note there's an intentional quirk here: If the Proxy presents itself as a plain object, we
1922 // wrap it in a stub, rather than serialize the object. This enables the Proxy to continue
1923 // intercepting property accesses when they happen, rather than have all the properties accessed
1924 // and serialized upfront. However, in retrospect, this may haev been a bad choice, as it means a
1925 // Proxy on a plain object cannot have exactly the same behavior as a plain object would have.
1926 // Note that apps which explicitly want to prevent a plain object from being serialized over
1927 // RPC can simply use `new RpcStub(object)` to explicitly wrap it in a stub -- no need to use
1928 // a Proxy for that.
1929 auto allowInstanceProperties = JSG_REQUIRE_NONNULL(checkStubType(js, handle), DOMDataCloneError,
1930 "Proxy could not be serialized because it is not a valid RPC receiver type. The "
1931 "Proxy must emulate either a plain object or an RpcTarget, as indicated by the "
1932 "Proxy's prototype chain.");
1933 
1934 // Similar to JsRpcTarget::serialize(), we may need to dup() the proxy.
1935 if (stubOwnership == RpcSerializerExternalHandler::DUPLICATE) {
1936 auto dup = handle.get(js, "dup");
1937 KJ_IF_SOME(dupFunc, dup.tryCast<jsg::JsFunction>()) {
1938 auto replacement = dupFunc.call(js, handle);
1939 bool replaced = false;
1940 
1941 // We got a duplicate. Is it still the same type?
1942 KJ_IF_SOME(replacementObj, replacement.tryCast<jsg::JsObject>()) {
1943 KJ_IF_SOME(stubType, checkStubType(js, replacementObj)) {
1944 if (stubType == allowInstanceProperties) {
1945 // It is! Let's replace our handle with the duplicate!
1946 handle = replacementObj;
1947 replaced = true;
1948 }
1949 }
1950 }
1951 
1952 JSG_REQUIRE(replaced, DOMDataCloneError,
1953 "Couldn't create a stub for the Proxy because it has a dup() method which did not "
1954 "return the same underlying type (RpcTarget or Function) as the Proxy itself represents. "
1955 "Either remove the dup() method or make sure it returns an RpcTarget.");
1956 }
1957 }
1958 
1959 // Great, we've concluded we can indeed point a stub at this proxy.
1960 serializer.writeRawUint32(static_cast<uint>(rpc::SerializationTag::JS_RPC_STUB));
1961 
1962 rpc::JsRpcTarget::Client cap =
1963 kj::heap<TransientJsRpcTarget>(js, IoContext::current(), handle, allowInstanceProperties);
1964 write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable {
1965 builder.setRpcTarget(kj::mv(cap));
1966 });
1967}
1968 
1969// JsRpcTarget implementation specific to entrypoints. This is used to deliver the first, top-level
1970// call of an RPC session.
1971class EntrypointJsRpcTarget final: public JsRpcTargetBase {
1972 public:
1973 EntrypointJsRpcTarget(IoContext& ioCtx,
1974 kj::Maybe<kj::StringPtr> entrypointName,
1975 kj::Maybe<Worker::VersionInfo> versionInfo,
1976 Frankenvalue props,
1977 kj::Maybe<kj::String> wrapperModule,
1978 kj::Maybe<kj::Own<BaseTracer>> tracer,
1979 bool isDynamicDispatch)
1980 : JsRpcTargetBase(ioCtx, CantOutliveIncomingRequest()),
1981 ioCtx(ioCtx),
1982 // Most of the time we don't really have to clone this but it's hard to fully prove, so
1983 // let's be safe.
1984 entrypointName(entrypointName.map([](kj::StringPtr s) { return kj::str(s); })),
1985 versionInfo(kj::mv(versionInfo)),
1986 props(kj::mv(props)),
1987 wrapperModule(kj::mv(wrapperModule)),
1988 tracer(kj::mv(tracer)),
1989 isDynamicDispatch(isDynamicDispatch) {}
1990 
1991 // Override call() to emit the Return event when the top-level RPC call completes.
1992 // This marks when the handler returned a value, NOT when all data has been streamed or all
1993 // capabilities released.
1994 kj::Promise<void> call(CallContext callContext) override {
1995 return JsRpcTargetBase::call(kj::mv(callContext)).then([this]() {
1996 KJ_IF_SOME(t, ioCtx.getWorkerTracer()) {
1997 t.setReturn(ioCtx.now());
1998 }
1999 });
2000 }
2001 
2002 TargetInfo getTargetInfo(Worker::Lock& lock, IoContext& ioCtx) override {
2003 jsg::Lock& js = lock;
2004 
2005 auto handler = KJ_REQUIRE_NONNULL(lock.getExportedHandler(entrypointName, kj::mv(versionInfo),
2006 kj::mv(props), ioCtx.getActor(), isDynamicDispatch),
2007 "Failed to get handler to worker.");
2008 
2009 if (handler->missingSuperclass && wrapperModule == kj::none) {
2010 // JS RPC is not enabled on the server side, we cannot call any methods.
2011 JSG_REQUIRE(FeatureFlags::get(js).getJsRpc(), TypeError,
2012 "The receiving Durable Object does not support RPC, because its class was not declared "
2013 "with `extends DurableObject`. In order to enable RPC, make sure your class "
2014 "extends the special class `DurableObject`, which can be imported from the module "
2015 "\"cloudflare:workers\".");
2016 }
2017 
2018 auto target = jsg::JsObject(handler->self.getHandle(lock));
2019 
2020 KJ_IF_SOME(moduleName, wrapperModule) {
2021 // We've been asked to apply a wrapper module to the handler. This is a builtin module whose
2022 // default export is a class. The class is constructed with the constructor arguments being
2023 // the ctx and env objects and the original DO instance.
2024 
2025 // This mechanism probably won't work very well on anything other than Durable Objects, so
2026 // block such usage for now. We could reconsider this if we have a use case in the future.
2027 auto& actor = JSG_REQUIRE_NONNULL(
2028 ioCtx.getActor(), Error, "Wrapper modules can only be applied to Durable Objects.");
2029 
2030 auto module = JSG_REQUIRE_NONNULL(
2031 js.resolveInternalModule(moduleName), Error, "Unknown internal module: ", moduleName);
2032 v8::Local<v8::Value> defaultExport = module.get(js, "default"_kj);
2033 JSG_REQUIRE(defaultExport->IsFunction(), TypeError,
2034 "Internal module's default export is not a function.");
2035 auto func = defaultExport.As<v8::Function>();
2036 
2037 v8::Local<v8::Value> args[3] = {actor.getCtx(js), actor.getEnv(js), target};
2038 auto jsContext = js.v8Context();
2039 v8::Local<v8::Value> result = jsg::check(func->NewInstance(jsContext, 3, args));
2040 JSG_REQUIRE(result->IsObject(), TypeError,
2041 "Internal module wrapper function did not return an object.");
2042 target = jsg::JsObject(result.As<v8::Object>());
2043 }
2044 
2045 // clang-format off
2046 TargetInfo targetInfo{
2047 .target = target,
2048 .envCtx = handler->ctx.map([&](jsg::Ref<ExecutionContext>& execCtx) -> EnvCtx {
2049 return {
2050 .env = handler->env.getHandle(js),
2051 .ctx = lock.getWorker().getIsolate().getApi().wrapExecutionContext(js, execCtx.addRef()),
2052 };
2053 })
2054 };
2055 // clang-format on
2056 
2057 // `targetInfo.envCtx` is present when we're invoking a freestanding function, and therefore
2058 // `env` and `ctx` need to be passed as parameters. In that case, we our method lookup
2059 // should obviously permit instance properties, since we expect the export is a plain object.
2060 // Otherwise, though, the export is a class. In that case, we have set the rule that we will
2061 // only allow class properties (aka prototype properties) to be accessed, to avoid
2062 // programmers shooting themselves in the foot by forgetting to make their members private.
2063 targetInfo.allowInstanceProperties = targetInfo.envCtx != kj::none;
2064 
2065 return targetInfo;
2066 }
2067 
2068 private:
2069 IoContext& ioCtx;
2070 kj::Maybe<kj::String> entrypointName;
2071 kj::Maybe<Worker::VersionInfo> versionInfo;
2072 Frankenvalue props;
2073 kj::Maybe<kj::String> wrapperModule;
2074 kj::Maybe<kj::Own<BaseTracer>> tracer;
2075 bool isDynamicDispatch;
2076 
2077 bool isReservedName(kj::StringPtr name) override {
2078 if ( // "fetch" and "connect" are treated specially on entrypoints.
2079 name == "fetch" || name == "connect" ||
2080 
2081 // These methods are reserved by the Durable Objects implementation.
2082 // TODO(someday): Should they be reserved only for Durable Objects, not WorkerEntrypoint?
2083 name == "alarm" || name == "webSocketMessage" || name == "webSocketClose" ||
2084 name == "webSocketError" ||
2085 
2086 // dup() is reserved to duplicate the stub itself, pointing to the same object.
2087 name == "dup" ||
2088 
2089 // All JS classes define a method `constructor` on the prototype, but we don't actually
2090 // want this to be callable over RPC!
2091 name == "constructor") {
2092 return true;
2093 }
2094 return false;
2095 }
2096 
2097 void maybeSetJsRpcInfo(IoContext& ctx, const kj::ConstString& methodNameForTrace) override {
2098 KJ_IF_SOME(tracer, ctx.getWorkerTracer()) {
2099 tracer.setJsRpcInfo(ctx.getInvocationSpanContext(), ctx.now(), methodNameForTrace);
2100 }
2101 }
2102};
2103 
2104// A membrane which wraps the top-level JsRpcTarget of an RPC session on the server side. The
2105// purpose of this membrane is to allow only a single top-level call, which then gets a
2106// `CompletionMembrane` wrapped around it. Note that we can't just wrap `CompletionMembrane` around
2107// the top-level object directly because that capability will not be dropped until the RPC session
2108// completes, since it is actually returned as the result of the top-level RPC call, but that
2109// call doesn't return until the `CompletionMembrane` says all capabilities were dropped, so this
2110// would create a cycle.
2111class JsRpcSessionCustomEvent::ServerTopLevelMembrane final: public capnp::MembranePolicy,
2112 public kj::Refcounted {
2113 public:
2114 explicit ServerTopLevelMembrane(kj::Own<kj::PromiseFulfiller<void>> doneFulfiller)
2115 : completionMembrane(kj::refcounted<CompletionMembrane>(kj::mv(doneFulfiller))) {}
2116 
2117 ~ServerTopLevelMembrane() noexcept(false) {
2118 KJ_IF_SOME(cm, completionMembrane) {
2119 cm->reject(
2120 KJ_EXCEPTION(DISCONNECTED, "JS RPC session canceled without calling an RPC method."));
2121 }
2122 }
2123 
2124 kj::Maybe<capnp::Capability::Client> inboundCall(
2125 uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override {
2126 if (interfaceId == capnp::typeId<rpc::JsRpcTarget>()) {
2127 // JsRpcTarget::call()
2128 auto cm = kj::mv(JSG_REQUIRE_NONNULL(
2129 completionMembrane, Error, "Only one RPC method call is allowed on this object."));
2130 completionMembrane = kj::none;
2131 return capnp::membrane(kj::mv(target), kj::mv(cm));
2132 } else if (interfaceId == capnp::typeId<rpc::JsValue::ExternalPusher>()) {
2133 // ExternalPusher methods
2134 //
2135 // It's important that we use the same membrane that we'll use for call(), so that
2136 // capabilities returned by the ExternalPusher will be wrapped in the membrane, hence they
2137 // will be unwrapped when passed back through the membrane again to call().
2138 auto& cm = *JSG_REQUIRE_NONNULL(
2139 completionMembrane, Error, "getExternalPusher() must be called before call()");
2140 return capnp::membrane(kj::mv(target), kj::addRef(cm));
2141 } else {
2142 KJ_FAIL_ASSERT("unkown interface ID for JsRpcTarget");
2143 }
2144 }
2145 
2146 kj::Maybe<capnp::Capability::Client> outboundCall(
2147 uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override {
2148 KJ_FAIL_ASSERT("ServerTopLevelMembrane shouldn't have outgoing capabilities");
2149 }
2150 
2151 kj::Own<MembranePolicy> addRef() override {
2152 return kj::addRef(*this);
2153 }
2154 
2155 private:
2156 kj::Maybe<kj::Own<CompletionMembrane>> completionMembrane;
2157};
2158 
2159kj::Promise<WorkerInterface::CustomEvent::Result> JsRpcSessionCustomEvent::run(
2160 kj::Own<IoContext::IncomingRequest> incomingRequest,
2161 kj::Maybe<kj::StringPtr> entrypointName,
2162 kj::Maybe<Worker::VersionInfo> versionInfo,
2163 Frankenvalue props,
2164 kj::TaskSet& waitUntilTasks,
2165 bool isDynamicDispatch) {
2166 IoContext& ioctx = incomingRequest->getContext();
2167 
2168 incomingRequest->delivered();
2169 
2170 KJ_DEFER({
2171 // waitUntil() should allow extending execution on the server side even when the client
2172 // disconnects.
2173 waitUntilTasks.add(incomingRequest->drain().attach(kj::mv(incomingRequest)));
2174 });
2175 
2176 EntrypointJsRpcTarget target(ioctx, entrypointName, kj::mv(versionInfo), kj::mv(props),
2177 kj::mv(wrapperModule), mapAddRef(incomingRequest->getWorkerTracer()), isDynamicDispatch);
2178 capnp::RevocableServer<rpc::JsRpcTarget> revcableTarget(target);
2179 
2180 try {
2181 auto [donePromise, doneFulfiller] = kj::newPromiseAndFulfiller<void>();
2182 
2183 kj::Own<capnp::MembranePolicy> topMembrane;
2184 if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_SESSION_HANDLE)) {
2185 // When using the session handle approach, we don't need the convoluted
2186 // `ServerTopLevelMembrane` because the the top-level `JsRpcTarget` is not unnaturally held
2187 // open, so it can be treated the same as any other capability in the session.
2188 topMembrane = kj::refcounted<CompletionMembrane>(kj::mv(doneFulfiller));
2189 } else {
2190 topMembrane = kj::refcounted<ServerTopLevelMembrane>(kj::mv(doneFulfiller));
2191 }
2192 
2193 capFulfiller->fulfill(capnp::membrane(revcableTarget.getClient(), kj::mv(topMembrane)));
2194 
2195 // `donePromise` resolves once there are no longer any capabilities pointing between the client
2196 // and server as part of this session.
2197 co_await donePromise.exclusiveJoin(ioctx.onAbort());
2198 
2199 co_return WorkerInterface::CustomEvent::Result{.outcome = EventOutcome::OK};
2200 } catch (...) {
2201 // Make sure the top-level capability is revoked with the same exception that `run()` is
2202 // throwing, rather than some generic revocation exception.
2203 auto e = kj::getCaughtExceptionAsKj();
2204 revcableTarget.revoke(e.clone());
2205 kj::throwFatalException(kj::mv(e));
2206 }
2207}
2208 
2209kj::Promise<WorkerInterface::CustomEvent::Result> JsRpcSessionCustomEvent::sendRpc(
2210 capnp::HttpOverCapnpFactory& httpOverCapnpFactory,
2211 capnp::ByteStreamFactory& byteStreamFactory,
2212 rpc::EventDispatcher::Client dispatcher) {
2213 // We arrange to revoke all capabilities in this session as soon as `sendRpc()` completes or is
2214 // canceled. Normally, the server side doesn't return if any capabilities still exist, so this
2215 // only makes a difference in the case that some sort of an error occurred. We don't strictly
2216 // have to revoke the capabilities as they are probably already broken anyway, but revoking them
2217 // helps to ensure that the underlying transport isn't "held open" waiting for the JS garbage
2218 // collector to actually collect the JsRpcStub objects.
2219 auto revokePaf = kj::newPromiseAndFulfiller<void>();
2220 
2221 KJ_DEFER({
2222 if (revokePaf.fulfiller->isWaiting()) {
2223 revokePaf.fulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "JS-RPC session canceled"));
2224 }
2225 });
2226 
2227 auto req = dispatcher.jsRpcSessionRequest();
2228 auto sent = req.send();
2229 
2230 rpc::JsRpcTarget::Client cap = sent.getTopLevel();
2231 
2232 cap = capnp::membrane(kj::mv(cap), kj::refcounted<RevokerMembrane>(kj::mv(revokePaf.promise)));
2233 
2234 // When no more capabilities exist on the connection, we want to proactively cancel the RPC.
2235 // This is needed in particular for the case where the client is dropped without making any calls
2236 // at all, e.g. because serializing the arguments failed. Unfortunately, simply dropping the
2237 // capability obtained through `sent.getTopLevel()` above will not be detected by the server,
2238 // because this is a pipeline capability on a call that is still running. So, if we don't
2239 // actually cancel the connection client-side, the server will hang open waiting for the initial
2240 // top-level call to arrive, and the event will appear never to complete at our end.
2241 //
2242 // TODO(cleanup): It feels like there's something wrong with the design here. Can we make this
2243 // less ugly?
2244 auto completionPaf = kj::newPromiseAndFulfiller<void>();
2245 cap = capnp::membrane(
2246 kj::mv(cap), kj::refcounted<CompletionMembrane>(kj::mv(completionPaf.fulfiller)));
2247 
2248 this->capFulfiller->fulfill(kj::mv(cap));
2249 
2250 auto session = sent.getSession();
2251 
2252 // We don't need to await the call itself as `session.whenResolved()` will already propagate any
2253 // errors from it. So we can drop the call promise now.
2254 //
2255 // Note that it would NOT work to use `req.sendForPipeline()` above, since pipelined capabilities
2256 // cannot resolve until the call returns, but `sendForPipeline()` explicitly inhibits the return
2257 // message.
2258 { auto drop = kj::mv(sent); }
2259 
2260 try {
2261 // Wait for `session` to resolve to a null capability.
2262 //
2263 // Note that this works even if the server is using the "old approach" where it doesn't return
2264 // a `session` at all, because in that case the return itself represents the end of the
2265 // session, and the response contains a null pointer for `session`, so this does the expected
2266 // thing: resolves `session` to null.
2267 co_await session.whenResolved().exclusiveJoin(kj::mv(completionPaf.promise));
2268 } catch (...) {
2269 auto e = kj::getCaughtExceptionAsKj();
2270 if (revokePaf.fulfiller->isWaiting()) {
2271 revokePaf.fulfiller->reject(e.clone());
2272 }
2273 kj::throwFatalException(kj::mv(e));
2274 }
2275 
2276 co_return WorkerInterface::CustomEvent::Result{.outcome = EventOutcome::OK};
2277}
2278 
2279kj::Promise<void> JsRpcSessionCustomEvent::receiveRpc(JsRpcSessionContext context,
2280 WorkerInterface& worker,
2281 kj::Own<void> ownWorker,
2282 kj::Maybe<kj::String> wrapperModule) {
2283 // Client wants to start a JS RPC session, we'll dispatch to the WorkerEntrypoint
2284 // here and read the capability off the event itself.
2285 auto customEvent =
2286 kj::heap<api::JsRpcSessionCustomEvent>(WORKER_RPC_EVENT_TYPE, kj::mv(wrapperModule));
2287 
2288 auto cap = customEvent->getCap();
2289 
2290 if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_SESSION_HANDLE)) {
2291 auto promise = worker.customEvent(kj::mv(customEvent));
2292 
2293 auto results = context.getResults(capnp::MessageSize{4, 2});
2294 results.setTopLevel(kj::mv(cap));
2295 
2296 // Set the returned session capability to resolve to a null capability when the event is
2297 // complete. This also neatly arranges that if the session is dropped early, the
2298 // `customEvent()` promise is canceled, thus canceling the session.
2299 results.setSession(promise.then([ownWorker = kj::mv(ownWorker)](auto outcome) {
2300 return rpc::JsRpcSession::Client(nullptr);
2301 }));
2302 } else {
2303 capnp::PipelineBuilder<rpc::EventDispatcher::JsRpcSessionResults> pipelineBuilder;
2304 pipelineBuilder.setTopLevel(cap);
2305 context.setPipeline(pipelineBuilder.build());
2306 context.getResults().setTopLevel(kj::mv(cap));
2307 
2308 co_await worker.customEvent(kj::mv(customEvent));
2309 }
2310}
2311 
2312}; // namespace workerd::api