// Copyright (c) 2017-2023 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include #include #include #include #include #include #include #include #include namespace workerd::api { namespace { using StreamSinkFulfiller = kj::Own>; } // namespace // Implementation of StreamSink RPC interface. The stream sender calls `startStream()` when // serializing each stream, and the recipient calls `setSlot()` when deserializing streams to // provide the appropriate destination capability. This class is designed to allow these two // calls to happen in either order for each slot. class StreamSinkImpl final: public rpc::JsValue::StreamSink::Server, public kj::Refcounted { public: ~StreamSinkImpl() noexcept(false) { for (auto& slot: table) { KJ_IF_SOME(f, slot.tryGet()) { f->reject(KJ_EXCEPTION(FAILED, "expected startStream() was never received")); } } } void setSlot(uint i, capnp::Capability::Client stream) { if (table.size() <= i) table.resize(i + 1); if (table[i] == nullptr) { table[i] = kj::mv(stream); } else KJ_SWITCH_ONEOF(table[i]) { KJ_CASE_ONEOF(stream, capnp::Capability::Client) { KJ_FAIL_REQUIRE("setSlot() tried to set the same slot twice", i); } KJ_CASE_ONEOF(fulfiller, StreamFulfiller) { fulfiller->fulfill(kj::mv(stream)); table[i] = Consumed(); } KJ_CASE_ONEOF(_, Consumed) { KJ_FAIL_REQUIRE("setSlot() tried to set the same slot twice", i); } } } kj::Promise startStream(StartStreamContext context) override { uint i = context.getParams().getExternalIndex(); if (table.size() <= i) { // guard against ridiculous table allocation JSG_REQUIRE(i < 1024, Error, "Too many streams in one message."); table.resize(i + 1); } if (table[i] == nullptr) { auto paf = kj::newPromiseAndFulfiller(); table[i] = kj::mv(paf.fulfiller); context.getResults(capnp::MessageSize{4, 1}).setStream(kj::mv(paf.promise)); } else KJ_SWITCH_ONEOF(table[i]) { KJ_CASE_ONEOF(stream, capnp::Capability::Client) { context.getResults(capnp::MessageSize{4, 1}).setStream(kj::mv(stream)); table[i] = Consumed(); } KJ_CASE_ONEOF(fulfiller, StreamFulfiller) { KJ_FAIL_REQUIRE("startStream() tried to start the same stream twice", i); } KJ_CASE_ONEOF(_, Consumed) { KJ_FAIL_REQUIRE("startStream() tried to start the same stream twice", i); } } return kj::READY_NOW; } private: using StreamFulfiller = kj::Own>; struct Consumed {}; // Each slot starts out null (uninitialized). It becomes a Capability::Client if setSlot() is // called first, or a StreamFulfiller if startStream() is called first. It becomes `Consumed` // when the other method is called. // HACK: Slots in the table take advantage of the little-known fact that OneOf has a "null" // value, which is the value a OneOf has when default-initialized. This is useful because we // don't want to explicitly initialize skipped slots. Maybe would be another option // here, but would add 8 bytes to every slot just to store a boolean... feels bloated. There // are only two methods in this class so I think it's OK. using Slot = kj::OneOf; kj::Vector table; }; kj::Maybe RpcSerializerExternalHandler::getExternalPusher() { KJ_IF_SOME(ep, externalPusher) { return ep; } else KJ_IF_SOME(func, getStreamHandlerFunc.tryGet()) { // First call, set up ExternalPusher. return externalPusher.emplace(func()); } else { // Using StreamSink. return kj::none; } } capnp::Capability::Client RpcSerializerExternalHandler::writeStream(BuilderCallback callback) { rpc::JsValue::StreamSink::Client* streamSinkPtr; KJ_IF_SOME(ss, streamSink) { streamSinkPtr = &ss; } else { // First stream written, set up the StreamSink. auto& func = KJ_REQUIRE_NONNULL(getStreamHandlerFunc.tryGet(), "this serialization is not using StreamSink; use getExternalPusher() instead"); streamSinkPtr = &streamSink.emplace(func()); } auto result = ({ auto req = streamSinkPtr->startStreamRequest(capnp::MessageSize{4, 0}); req.setExternalIndex(externals.size()); req.send().getStream(); }); write(kj::mv(callback)); return result; } capnp::Orphan> RpcSerializerExternalHandler::build( capnp::Orphanage orphanage) { auto result = orphanage.newOrphan>(externals.size()); auto builder = result.get(); for (auto i: kj::indices(externals)) { externals[i](builder[i]); } return result; } RpcDeserializerExternalHandler::~RpcDeserializerExternalHandler() noexcept(false) { if (!unwindDetector.isUnwinding()) { KJ_ASSERT(i == externals.size(), "deserialization did not consume all of the externals"); } } rpc::JsValue::External::Reader RpcDeserializerExternalHandler::read() { KJ_ASSERT(i < externals.size()); return externals[i++]; } void RpcDeserializerExternalHandler::setLastStream(capnp::Capability::Client stream) { KJ_IF_SOME(ss, streamSink) { ss.setSlot(i - 1, kj::mv(stream)); } else { auto ss = kj::refcounted(); ss->setSlot(i - 1, kj::mv(stream)); streamSink = *ss; streamSinkCap = rpc::JsValue::StreamSink::Client(kj::mv(ss)); } } namespace { // Call to construct an `rpc::JsValue` from a JS value. // // `makeBuilder` is a function which takes a capnp::MessageSize hint and returns the // rpc::JsValue::Builder to fill in. template void serializeJsValue(jsg::Lock& js, jsg::JsValue value, RpcSerializerExternalHandler& externalHandler, Func makeBuilder) { jsg::Serializer serializer(js, jsg::Serializer::Options{ .version = 15, .omitHeader = false, .treatClassInstancesAsPlainObjects = false, .externalHandler = externalHandler, }); serializer.write(js, value); kj::Array data = serializer.release().data; JSG_ASSERT(data.size() <= MAX_JS_RPC_MESSAGE_SIZE, Error, "Serialized RPC arguments or return values are limited to 32MiB, but the size of this value " "was: ", data.size(), " bytes."); capnp::MessageSize hint{0, 0}; hint.wordCount += (data.size() + sizeof(capnp::word) - 1) / sizeof(capnp::word); hint.wordCount += capnp::sizeInWords(); hint.wordCount += externalHandler.size() * capnp::sizeInWords(); hint.capCount += externalHandler.size(); rpc::JsValue::Builder builder = makeBuilder(hint); // TODO(perf): It would be nice if we could serialize directly into the capnp message to avoid // a redundant copy of the bytes here. Maybe we could even cancel serialization early if it // goes over the size limit. builder.setV8Serialized(data); if (externalHandler.size() > 0) { builder.adoptExternals( externalHandler.build(capnp::Orphanage::getForMessageContaining(builder))); } } struct DeserializeResult { jsg::JsValue value; kj::Own disposalGroup; kj::Maybe streamSink; }; // Call to construct a JS value from an `rpc::JsValue`. DeserializeResult deserializeJsValue(jsg::Lock& js, rpc::JsValue::Reader reader, kj::LiteralStringConst debugContext, kj::Maybe streamSink = kj::none) { auto disposalGroup = kj::heap(); RpcDeserializerExternalHandler externalHandler( reader.getExternals(), *disposalGroup, streamSink, debugContext); jsg::Deserializer deserializer(js, reader.getV8Serialized(), kj::none, kj::none, jsg::Deserializer::Options{ .version = 15, .readHeader = true, // Previously, while these are passing over an RPC boundary, we preserved stack // traces in errors that happened to get passed through rather than thrown. // This was mainly due, I believe, to a misunderstanding about whether or not // v8 serialization preserved the stacks or not. When enhanced error serialization // is disabled, stacks are preserved and this flag has no effect. When enhanced // error serialization is enabled, then we'll switch to not preserving stacks in // passed-through errors. .preserveStackInErrors = false, .externalHandler = externalHandler, }); return { .value = deserializer.readValue(js), .disposalGroup = kj::mv(disposalGroup), .streamSink = externalHandler.getStreamSink(), }; } // Does deserializeJsValue() and then adds a `dispose()` method to the returned object (if it is // an object) which disposes all stubs therein. jsg::JsValue deserializeRpcReturnValue(jsg::Lock& js, rpc::JsRpcTarget::CallResults::Reader callResults, kj::Maybe streamSink) { auto [value, disposalGroup, ss] = deserializeJsValue(js, callResults.getResult(), "return"_kjc, streamSink); if (streamSink == kj::none) { KJ_REQUIRE(ss == kj::none, "RPC returned result using StreamSink even though ExternalPusher was provided"); } // If the object had a disposer on the callee side, it will run when we discard the callPipeline, // so attach that to the disposal group on the caller side. If the returned object did NOT have // a disposer then we should discard callPipeline so that we don't hold open the callee's // context for no reason. if (callResults.getHasDisposer()) { disposalGroup->setCallPipeline( IoContext::current().addObject(kj::heap(callResults.getCallPipeline()))); } KJ_IF_SOME(obj, value.tryCast()) { if (obj.isInstanceOf(js)) { // We're returning a plain stub. We don't need to override its `dispose` method. disposalGroup->disownAll(); } else { // Add a dispose method to the return object that disposes the DisposalGroup. v8::Local func = js.wrapSimpleFunction(js.v8Context(), [disposalGroup = kj::mv(disposalGroup)](jsg::Lock&, const v8::FunctionCallbackInfo&) mutable { disposalGroup->disposeAll(); }); obj.setNonEnumerable(js, js.symbolDispose(), jsg::JsValue(func)); } } else { // Result wasn't an object, so it must not contain any stubs. KJ_ASSERT(disposalGroup->empty()); } return value; } // A membrane which attaches some object until it is destroyed. // // TODO(cleanup): This is generally useful, should it be part of capnp? class AttachmentMembrane final: public capnp::MembranePolicy, public kj::Refcounted { public: explicit AttachmentMembrane(kj::Own attachment): attachment(kj::mv(attachment)) {} kj::Maybe inboundCall( uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override { return kj::none; } kj::Maybe outboundCall( uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override { return kj::none; } kj::Own addRef() override { return kj::addRef(*this); } private: kj::Own attachment; }; // Given a value, check if it has a dispose method and, if so, invoke it. void tryCallDisposeMethod(jsg::Lock& js, jsg::JsValue value) { js.withinHandleScope([&]() { KJ_IF_SOME(obj, value.tryCast()) { auto dispose = obj.get(js, js.symbolDispose()); if (dispose.isFunction()) { jsg::check(v8::Local(dispose).As()->Call( js.v8Context(), value, 0, nullptr)); } } }); } } // namespace JsRpcPromise::JsRpcPromise(jsg::JsRef inner, kj::Own weakRefParam, IoOwn pipeline) : inner(kj::mv(inner)), weakRef(kj::mv(weakRefParam)), state(Pending{kj::mv(pipeline)}) { KJ_REQUIRE(weakRef->ref == kj::none); weakRef->ref = *this; } JsRpcPromise::~JsRpcPromise() noexcept(false) { weakRef->ref = kj::none; } void JsRpcPromise::resolve(jsg::Lock& js, jsg::JsValue result) { if (state.is()) { state = Resolved{ .result = jsg::Value(js.v8Isolate, result), .ctxCheck = IoContext::current().addObject(*this), }; } else { // We'd better dispose this. tryCallDisposeMethod(js, result); } } void JsRpcPromise::dispose(jsg::Lock& js) { KJ_IF_SOME(resolved, state.tryGet()) { // Disposing the promise implies disposing the final result. tryCallDisposeMethod(js, jsg::JsValue(resolved.result.getHandle(js))); } state = Disposed(); weakRef->disposed = true; } // See comment at call site for explanation. static rpc::JsRpcTarget::Client makeJsRpcTargetForSingleLoopbackCall( jsg::Lock& js, jsg::JsObject obj); rpc::JsRpcTarget::Client JsRpcPromise::getClientForOneCall( jsg::Lock& js, kj::Vector& path) { // (Don't extend `path` because we're the root.) KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(pending, Pending) { return pending.pipeline->getCallPipeline(); } KJ_CASE_ONEOF(resolved, Resolved) { // Dereference `ctxCheck` just to verify we're running in the correct context. (If not, // this will throw.) *resolved.ctxCheck; // A value was already returned, and we closed the original RPC pipeline. But the application // kept the promise around and is still trying to pipeline on it. What do we do? // // A naive answer would be: We just return the actual value that was returned originally. // Like if someone asked for `promise.foo.bar`, we just give them `returnValue.foo.bar`. // // That doesn't quite work, for a couple reasons: // * If the caller is awaiting a property, they expect the result will have a `dispose()` // method added to it, and that any stubs in the result will be independently disposable. // This essentially means we need to clone the value so that we can dup() all the stubs and // modify the result. // * If the caller is trying to make a pipelined RPC call, they expect this call to go // through all the usual RPC machinery. They do NOT expect that this is going to be a local // call. // // The easiest way to make this all just work is... to actually wrap the value in a one-off // RPC stub, and make a real RPC on it. return js.withinHandleScope([&]() -> rpc::JsRpcTarget::Client { auto value = jsg::JsValue(resolved.result.getHandle(js)); KJ_IF_SOME(obj, value.tryCast()) { KJ_IF_SOME(stub, obj.tryUnwrapAs(js)) { // Oh, the return value is actually a stub itself. Just use it. return stub->getClient(); } else { // Must be a plain object. return makeJsRpcTargetForSingleLoopbackCall(js, obj); } } else { JSG_FAIL_REQUIRE(TypeError, "Can't pipeline on RPC that did not return an object."); } }); } KJ_CASE_ONEOF(disposed, Disposed) { return JSG_KJ_EXCEPTION(FAILED, Error, "RPC promise used after being disposed."); } } KJ_UNREACHABLE; } rpc::JsRpcTarget::Client JsRpcProperty::getClientForOneCall( jsg::Lock& js, kj::Vector& path) { auto result = parent->getClientForOneCall(js, path); path.add(name); return result; } namespace { struct JsRpcPromiseAndPipeline { jsg::JsPromise promise; kj::Own weakRef; rpc::JsRpcTarget::CallResults::Pipeline pipeline; jsg::Ref asJsRpcPromise(jsg::Lock& js) && { return js.alloc(jsg::JsRef(js, promise), kj::mv(weakRef), IoContext::current().addObject(kj::heap(kj::mv(pipeline)))); } }; // Core implementation of making an RPC call, reusable for many cases below. JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, JsRpcClientProvider& parent, kj::Maybe name, // If `maybeArgs` is provided, this is a call, otherwise it is a property access. kj::Maybe&> maybeArgs) { // Note: We used to enforce that RPC methods had to be called with the correct `this`. That is, // we prevented people from doing: // // let obj = {foo: someRpcStub.foo}; // obj.foo(); // // This would throw "Illegal invocation", as is the norm when pulling methods of a native object. // That worked as long as RPC methods were implemented as `jsg::Function`. However, when we // switched to RPC methods being implemented as callable objects (JsRpcProperty), this became // impossible, because V8's SetCallAsFunctionHandler() arranges that `this` is bound to the // callable object itself, regardless of how it was invoked. So now we cannot detect the // situation above, because V8 never tells us about `obj` at all. // // Oh well. It's not a big deal. Just annoying that we have to forever support tearing RPC // methods off their source object, even if we change implementations to something where that's // less convenient. try { return js.tryCatch([&]() -> JsRpcPromiseAndPipeline { // `path` will be filled in with the path of property names leading from the stub represented by // `client` to the specific property / method that we're trying to invoke. kj::Vector path; auto client = parent.getClientForOneCall(js, path); auto& ioContext = IoContext::current(); KJ_IF_SOME(lock, ioContext.waitForOutputLocksIfNecessary()) { // Replace the client with a promise client that will delay the call until the output gate // is open. client = lock.then([client = kj::mv(client)]() mutable { return kj::mv(client); }); } auto builder = client.callRequest(); // This code here is slightly overcomplicated in order to avoid pushing anything to the // kj::Vector in the common case that the parent path is empty. I'm probably trying too hard // but oh well. if (path.empty()) { KJ_IF_SOME(n, name) { builder.setMethodName(n); } else { // No name and no path, must be directly calling a stub. builder.initMethodPath(0); } } else { auto pathBuilder = builder.initMethodPath(path.size() + (name != kj::none)); for (auto i: kj::indices(path)) { pathBuilder.set(i, path[i]); } KJ_IF_SOME(n, name) { pathBuilder.set(path.size(), n); } } kj::Maybe paramsStreamSinkFulfiller; bool useExternalPusher = util::Autogate::isEnabled(util::AutogateKey::RPC_USE_EXTERNAL_PUSHER); KJ_IF_SOME(args, maybeArgs) { // If we have arguments, serialize them. // Note that we may fail to serialize some element, in which case this will throw back to // JS. if (args.Length() > 0) { // This is a function call with arguments. v8::LocalVector argv(js.v8Isolate, args.Length()); for (int n = 0; n < args.Length(); n++) { argv[n] = args[n]; } auto arr = v8::Array::New(js.v8Isolate, argv.data(), argv.size()); auto stubOwnership = FeatureFlags::get(js).getRpcParamsDupStubs() ? RpcSerializerExternalHandler::DUPLICATE : RpcSerializerExternalHandler::TRANSFER; RpcSerializerExternalHandler::GetStreamHandlerFunc getStreamHandlerFunc; if (useExternalPusher) { getStreamHandlerFunc.init( [&]() -> rpc::JsValue::ExternalPusher::Client { return client; }); } else { getStreamHandlerFunc.init([&]() { // A stream was encountered in the params, so we must expect the response to contain // paramsStreamSink. But we don't have the response yet. So, we need to set up a // temporary promise client, which we hook to the response a little bit later. auto paf = kj::newPromiseAndFulfiller(); paramsStreamSinkFulfiller = kj::mv(paf.fulfiller); return kj::mv(paf.promise); }); } RpcSerializerExternalHandler externalHandler(stubOwnership, kj::mv(getStreamHandlerFunc)); serializeJsValue(js, jsg::JsValue(arr), externalHandler, [&](capnp::MessageSize hint) { // TODO(perf): Actually use the size hint. return builder.getOperation().initCallWithArgs(); }); } } else { // This is a property access. builder.getOperation().setGetProperty(); } kj::Maybe> resultStreamSink; if (useExternalPusher) { // Unfortunately, we always have to send the ExternalPusher since we don't know whether the // call will return any streams (or other pushed externals). Luckily, it's a // one-per-IoContext object, not a big deal. (It'll take a slot on the capnp export table // though.) builder.getResultsStreamHandler().setExternalPusher(ioContext.getExternalPusher()); } else { // Unfortunately, we always have to send a `resultsStreamSink` because we don't know until // after the call completes whether or not it will return any streams. If it's unused, // though, it should only be a couple allocations. builder.getResultsStreamHandler().setStreamSink( kj::addRef(*resultStreamSink.emplace(kj::refcounted()))); } auto callResult = builder.send(); KJ_IF_SOME(ssf, paramsStreamSinkFulfiller) { ssf->fulfill(callResult.getParamsStreamSink()); } // We need to arrange that our JsRpcPromise will updated in-place with the final settlement // of this RPC promise. However, we can't actually construct the JsRpcPromise until we have // the final promise to give it. To resolve the cycle, we only create a JsRpcPromise::WeakRef // here, which is filled in later on to point at the JsRpcPromise, if and when one is created. auto weakRef = kj::atomicRefcounted(); // RemotePromise lets us consume its pipeline and promise portions independently; we consume // the promise here and we consume the pipeline below, both via kj::mv(). auto jsPromise = ioContext.awaitIo(js, kj::mv(callResult), [weakRef = kj::atomicAddRef(*weakRef), resultStreamSink = kj::mv(resultStreamSink)]( jsg::Lock& js, capnp::Response response) mutable -> jsg::Value { auto jsResult = deserializeRpcReturnValue(js, response, resultStreamSink); if (weakRef->disposed) { // The promise was explicitly disposed before it even resolved. This means we must dispose // the returned object as well. tryCallDisposeMethod(js, jsResult); } else { KJ_IF_SOME(r, weakRef->ref) { r.resolve(js, jsResult); } } return jsg::Value(js.v8Isolate, jsResult); }); return { .promise = jsg::JsPromise(js.wrapSimplePromise(kj::mv(jsPromise))), .weakRef = kj::mv(weakRef), .pipeline = kj::mv(callResult), }; }, [&](jsg::Value error) -> JsRpcPromiseAndPipeline { // Probably a serialization error. Need to convert to an async error since we never throw // synchronously from async functions. auto jsError = jsg::JsValue(error.getHandle(js)); auto pipeline = capnp::newBrokenPipeline(js.exceptionToKj(jsError)); return {.promise = js.rejectedJsPromise(jsError), .weakRef = kj::atomicRefcounted(), .pipeline = rpc::JsRpcTarget::CallResults::Pipeline(capnp::AnyPointer::Pipeline(kj::mv(pipeline)))}; }); } catch (jsg::JsExceptionThrown&) { // This must be a termination exception, or we would have caught it above. throw; } catch (...) { // Catch KJ exceptions and make them async, since we don't want async calls to throw // synchronously. auto e = kj::getCaughtExceptionAsKj(); auto pipeline = capnp::newBrokenPipeline(e.clone()); return { .promise = jsg::JsPromise(js.wrapSimplePromise(js.rejectedPromise(kj::mv(e)))), .weakRef = kj::atomicRefcounted(), .pipeline = rpc::JsRpcTarget::CallResults::Pipeline(capnp::AnyPointer::Pipeline(kj::mv(pipeline)))}; } } } // namespace jsg::Ref JsRpcProperty::call(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); return callImpl(js, *parent, name, args).asJsRpcPromise(js); } jsg::Ref JsRpcStub::call(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); return callImpl(js, *this, kj::none, args).asJsRpcPromise(js); } jsg::Ref JsRpcPromise::call(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); return callImpl(js, *this, kj::none, args).asJsRpcPromise(js); } namespace { jsg::JsValue thenImpl(jsg::Lock& js, v8::Local promise, v8::Local handler, jsg::Optional> errorHandler) { KJ_IF_SOME(e, errorHandler) { // Note that we intentionally propagate any exception from promise->Then() synchronously since // if V8's native Promise threw synchronously from `then()`, we might as well too. Anyway it's // probably a termination exception. return jsg::JsPromise(jsg::check(promise->Then(js.v8Context(), handler, e))); } else { return jsg::JsPromise(jsg::check(promise->Then(js.v8Context(), handler))); } } jsg::JsValue catchImpl( jsg::Lock& js, v8::Local promise, v8::Local errorHandler) { return jsg::JsPromise(jsg::check(promise->Catch(js.v8Context(), errorHandler))); } jsg::JsValue finallyImpl( jsg::Lock& js, v8::Local promise, v8::Local onFinally) { // HACK: `finally()` is not exposed as a C++ API, so we have to manually read it from JS. jsg::JsObject obj(promise); auto func = obj.get(js, "finally"); KJ_ASSERT(func.isFunction()); v8::Local param = onFinally; return jsg::JsValue(jsg::check( v8::Local(func).As()->Call(js.v8Context(), obj, 1, ¶m))); } } // namespace jsg::JsValue JsRpcProperty::then(jsg::Lock& js, v8::Local handler, jsg::Optional> errorHandler) { auto promise = callImpl(js, *parent, name, kj::none).promise; return thenImpl(js, promise, handler, errorHandler); } jsg::JsValue JsRpcProperty::catch_(jsg::Lock& js, v8::Local errorHandler) { auto promise = callImpl(js, *parent, name, kj::none).promise; return catchImpl(js, promise, errorHandler); } jsg::JsValue JsRpcProperty::finally(jsg::Lock& js, v8::Local onFinally) { auto promise = callImpl(js, *parent, name, kj::none).promise; return finallyImpl(js, promise, onFinally); } jsg::JsValue JsRpcPromise::then(jsg::Lock& js, v8::Local handler, jsg::Optional> errorHandler) { return thenImpl(js, inner.getHandle(js), handler, errorHandler); } jsg::JsValue JsRpcPromise::catch_(jsg::Lock& js, v8::Local errorHandler) { return catchImpl(js, inner.getHandle(js), errorHandler); } jsg::JsValue JsRpcPromise::finally(jsg::Lock& js, v8::Local onFinally) { return finallyImpl(js, inner.getHandle(js), onFinally); } kj::Maybe> JsRpcProperty::getProperty(jsg::Lock& js, kj::String name) { return js.alloc(JSG_THIS, kj::mv(name)); } kj::Maybe> JsRpcPromise::getProperty(jsg::Lock& js, kj::String name) { return js.alloc(JSG_THIS, kj::mv(name)); } JsRpcStub::JsRpcStub(IoOwn capnpClient, RpcStubDisposalGroup& disposalGroup, jsg::ExternalMemoryAdjustment externalMemoryAdjustment) : capnpClient(kj::mv(capnpClient)), disposalGroup(disposalGroup), externalMemoryAdjustment(kj::mv(externalMemoryAdjustment)) { disposalGroup.list.add(*this); } JsRpcStub::~JsRpcStub() noexcept(false) { KJ_IF_SOME(d, disposalGroup) { d.list.remove(*this); } KJ_IF_SOME(c, capnpClient) { // The app failed to dispose the stub; it leaked. We'd rather not make GC observable, so we // must pass the capnp capability off to the I/O context to be dropped when the I/O context // itself shuts down. kj::mv(c).deferGcToContext(); // In preview, let's try to warn the developer about the problem. // // TODO(cleanup): Instead of logging this warning at GC time, it would be better if we logged // it at the time that the client is destroyed, i.e. when the IoContext is torn down, // which is usually sooner (and more deterministic). But logging a warning during // IoContext tear-down is problematic since logWarningOnce() is a method on // IoContext... KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { ioContext.logWarningOnce( "An RPC stub was not disposed properly. You must call dispose() on all stubs in order to " "let the other side know that you are no longer using them. You cannot rely on " "the garbage collector for this because it may take arbitrarily long before actually " "collecting unreachable objects. As a shortcut, calling dispose() on the result of " "an RPC call disposes all stubs within it."_kj); } } } RpcStubDisposalGroup::~RpcStubDisposalGroup() noexcept(false) { if (jsg::isInGcDestructor()) { // If the disposal group was dropped as a result of garbage collection, we should NOT actually // dispose any stubs. In particular: // * If an application never invokes dispose() on an RPC result and the result is GC'ed, the // app could still be holding onto stubs that came from that result. We don't want to // dispose those unexpectedly. // * If an incoming RPC call does something like `await new Promise(() => {})` to hang // forever, the promise reaction can be GC'ed even though the call didn't really complete. // We don't want to dispose param stubs in this case. disownAll(); // If we have a `callPipeline`, it means we called an RPC that returned an object, and that // object had a dispose method defined on the server side. We don't want it to observe GC, // so we'll defer dropping the pipeline until the IoContext is destroyed. // // (We don't do this as part of disownAll() because the one other call site of disownAll() // is only invoked in cases where there shouldn't be a `callPipeline` anyway...) KJ_IF_SOME(c, callPipeline) { kj::mv(c).deferGcToContext(); // In preview, let's try to warn the developer about the problem. // // TODO(cleanup): Same comment as in ~JsRpcStub(). KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { ioContext.logWarningOnce( "An RPC result was not disposed properly. One of the RPC calls you made expects you " "to call dispose() on the return value, but you didn't do so. You cannot rely on " "the garbage collector for this because it may take arbitrarily long before actually " "collecting unreachable objects."_kj); } } } else { // However, if we're destroying the RpcStubDisposalGroup NOT as a result of GC, this probably // means one of: // * This is the disposal group for an incoming RPC call, and that call completed. The group // was attached to the completion continuation, which executed, and is now being destroyed. // This is the normal completion case, and we should dispose all the param stubs. // * An exception was thrown in the RPC implementation before stubs could be passed to // JavaScript in the first place, resulting in the disposal group being destroyed during // exception unwind. The stubs should be disposed proactively since they were never // received. disposeAll(); } } rpc::JsRpcTarget::Client JsRpcStub::getClient() { KJ_IF_SOME(c, capnpClient) { return *c; } else { // TODO(soon): Improve the error message to describe why it was disposed. return JSG_KJ_EXCEPTION(FAILED, Error, "RPC stub used after being disposed."); } } rpc::JsRpcTarget::Client JsRpcStub::getClientForOneCall( jsg::Lock& js, kj::Vector& path) { // (Don't extend `path` because we're the root.) return getClient(); } jsg::Ref JsRpcStub::dup(jsg::Lock& js) { return js.alloc(IoContext::current().addObject(kj::heap(getClient()))); } void JsRpcStub::dispose() { capnpClient = kj::none; externalMemoryAdjustment = kj::none; KJ_IF_SOME(d, disposalGroup) { d.list.remove(*this); disposalGroup = kj::none; } } void RpcStubDisposalGroup::disownAll() { for (auto& stub: list) { stub.disposalGroup = kj::none; list.remove(stub); } } void RpcStubDisposalGroup::disposeAll() { for (auto& stub: list) { stub.dispose(); } callPipeline = kj::none; // Each stub should have removed itself. KJ_ASSERT(list.empty()); } kj::Maybe> JsRpcStub::getRpcMethod(jsg::Lock& js, kj::String name) { // Do not return a method for `then`, otherwise JavaScript decides this is a thenable, i.e. a // custom Promise, which will mean a Promise that resolves to this object will attempt to chain // with it, which is not what you want! if (name == "then"_kj) return kj::none; return js.alloc(JSG_THIS, kj::mv(name)); } void JsRpcStub::serialize(jsg::Lock& js, jsg::Serializer& serializer) { auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError, "Remote RPC references can only be serialized for RPC."); auto externalHandler = dynamic_cast(&handler); JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError, "Remote RPC references can only be serialized for RPC."); // We may be forwarding a stub that points to some other isolate. Consider the case where we // are returning the stub to our client. The RPC session remains live as long as the client is // holding any remaining stubs obtained from this session, due to CompletionMembrane. However, if // the only remaining stubs point on to different isolates, and we don't have anything left to // do in this IoContext, then the pending event mechanism would abort the IoContext early with // "The script will never generate a response." To avoid that, we need to attach a pending event // to this stub, using a membrane. // // TODO(someday): Ideally, we would not need to keep the IoContext live just because stubs pass // through it. It would be nice to implement a sort of "deferred proxying" for RPC, where we // shut down the IoContext when it has nothing left to do. Note, though, that if the IoContext // is explicitly *aborted*, we probably should revoke all capabilities obtained through it. // That actually doesn't quite happen today: aborting the IoContext is likely to cancel all // subrequests which probably has the effect of breaking any stubs obtained from them, but // not necessarily (the subrequests could use waitUntil() to extend themselves). Anyway, this // will be trickier to get right, so I'm punting with this work-around for now. auto cap = capnp::membrane( getClient(), kj::refcounted(IoContext::current().registerPendingEvent())); externalHandler->write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable { builder.setRpcTarget(kj::mv(cap)); }); if (externalHandler->getStubOwnership() == RpcSerializerExternalHandler::TRANSFER) { // Instead of disposing the stub immediately, we add a disposer to the serializer // that will be executed when the pipeline is finished. This ensures the stub // remains valid for the duration of any pipelined operations. externalHandler->addStubDisposer( kj::heap(kj::defer([self = JSG_THIS]() mutable { self->dispose(); }))); } } jsg::Ref JsRpcStub::deserialize( jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) { auto& handler = KJ_REQUIRE_NONNULL( deserializer.getExternalHandler(), "got JsRpcStub on non-RPC serialized object?"); auto externalHandler = dynamic_cast(&handler); KJ_REQUIRE(externalHandler != nullptr, "got JsRpcStub on non-RPC serialized object?"); auto reader = externalHandler->read(); KJ_REQUIRE(reader.isRpcTarget(), "external table slot type doesn't match serialization tag"); auto& ioctx = IoContext::current(); // Account for membrane/promise memory in the KJ heap (~1600 bytes per stub from profiling). static constexpr size_t ESTIMATED_EXTERNAL_MEMORY_PER_STUB = 1600; auto externalMemory = js.getExternalMemoryAdjustment(ESTIMATED_EXTERNAL_MEMORY_PER_STUB); return js.alloc(ioctx.addObject(kj::heap(reader.getRpcTarget())), externalHandler->getDisposalGroup(), kj::mv(externalMemory)); } static bool isFunctionForRpc(jsg::Lock& js, v8::Local func) { jsg::JsObject obj(func); if (obj.isInstanceOf(js) || obj.isInstanceOf(js)) { // Don't allow JsRpcProperty or JsRpcPromise to be treated as plain functions, even though they // are technically callable. These types need to be treated specially (if we decide to let // them be passed over RPC at all). return false; } return true; } static bool isFunctionForRpc(jsg::Lock& js, jsg::JsValue value) { if (!value.isFunction()) return false; return isFunctionForRpc(js, v8::Local(value).As()); } // `makeCallPipeline()` has a bit of a complicated result type.. namespace MakeCallPipeline { // The value is an object, which may have stubs inside it. struct Object { rpc::JsRpcTarget::Client cap; // Was the value a plain JavaScript object which had a custom dispose() method? bool hasDispose; }; // The value was something that should serialize to a single stub (e.g. it was an RpcTarget, a // plain function, or already a stub). The callPipeline should simply be a copy of that stub. struct SingleStub {}; // The value is not a type that supports pipelining. It may still be serializable, and it could // even contain stubs (e.g. in a Map). struct NonPipelinable { // callPipeline to return just for error-handling purposes. rpc::JsRpcTarget::Client errorPipeline; }; using Result = kj::OneOf; }; // namespace MakeCallPipeline template MakeCallPipeline::Result serializeJsValueWithPipeline(jsg::Lock& js, jsg::JsValue value, Func makeBuilder, RpcSerializerExternalHandler::GetStreamHandlerFunc getStreamSinkFunc); // Callee-side implementation of JsRpcTarget. // // Most of the implementation is in this base class. There are subclasses specializing for the case // of a top-level entrypoint vs. a transient object introduced by a previous RPC in the same // session. class JsRpcTargetBase: public rpc::JsRpcTarget::Server { public: struct MayOutliveIncomingRequest {}; struct CantOutliveIncomingRequest {}; // Constructor used by TransientJsRpcTarget, which does not own the context. It needs to use // makeReentryCallback() to guard against the possibility that the IoContext is canceled before // or during a call. JsRpcTargetBase(IoContext& ctx, MayOutliveIncomingRequest) : enterIsolateAndCall(ctx.makeReentryCallback( [this, &ctx](Worker::Lock& lock, CallContext callContext) { return callImpl(lock, ctx, callContext); })), externalPusher(ctx.getExternalPusher()) {} // Constructor use by EntrypointJsRpcTarget, which is revoked and destroyed before the IoContext // can possibly be canceled. It can just use ctx.run(). JsRpcTargetBase(IoContext& ctx, CantOutliveIncomingRequest) : enterIsolateAndCall([this, &ctx](CallContext callContext) { // Note: No need to topUpActor() since this is the start of a top-level request, so the // actor will already have been topped up by IncomingRequest::delivered(). return ctx.run([this, &ctx, callContext](Worker::Lock& lock) mutable { return callImpl(lock, ctx, callContext); }); }), externalPusher(ctx.getExternalPusher()) {} struct EnvCtx { v8::Local env; jsg::JsObject ctx; }; struct TargetInfo { // The object on which the RPC method should be invoked. jsg::JsObject target; // If `env` and `ctx` need to be delivered as arguments to the method, these are the values // to deliver. kj::Maybe envCtx; bool allowInstanceProperties; }; // Get the object on which the method is to be invoked. This is virtual so that we can have // separate subclasses handling the case of an entrypoint vs. a transient RPC object. virtual TargetInfo getTargetInfo(Worker::Lock& lock, IoContext& ioCtx) = 0; // Handles the delivery of JS RPC method calls. kj::Promise call(CallContext callContext) override { co_await kj::yield(); // Try to execute the requested method. co_return co_await enterIsolateAndCall(callContext).catch_([](kj::Exception&& e) { if (jsg::isTunneledException(e.getDescription())) { // Annotate exceptions in RPC worker calls as remote exceptions. auto description = jsg::stripRemoteExceptionPrefix(e.getDescription()); if (!description.startsWith("remote.")) { // If we already were annotated as remote from some other worker entrypoint, no point // adding an additional prefix. e.setDescription(kj::str("remote.", description)); } } kj::throwFatalException(kj::mv(e)); }); } // Implements ExternalPusher by forwarding to the shared implementation. // // Note JsRpcTarget has to implement `ExternalPusher` directly rather than providing a method // like `getExternalPusher()` because it's important that the pushes arrive before the call, and // the ordering can only be guaranteed if they're on the same object. kj::Promise pushByteStream(PushByteStreamContext context) override { return externalPusher->pushByteStream(context); } kj::Promise pushAbortSignal(PushAbortSignalContext context) override { return externalPusher->pushAbortSignal(context); } KJ_DISALLOW_COPY_AND_MOVE(JsRpcTargetBase); private: virtual void maybeSetJsRpcInfo(IoContext& ctx, const kj::ConstString& methodNameForTrace) = 0; // Function which enters the isolate lock and IoContext and then invokes callImpl(). Created // using IoContext::makeReentryCallback(). kj::Function(CallContext callContext)> enterIsolateAndCall; kj::Rc externalPusher; // Returns true if the given name cannot be used as a method on this type. virtual bool isReservedName(kj::StringPtr name) = 0; kj::Promise callImpl(Worker::Lock& lock, IoContext& ctx, CallContext callContext) { jsg::Lock& js = lock; auto params = callContext.getParams(); // Method name suitable for use in trace and error messages. May be a pointer into the RPC // params reader. kj::ConstString methodNameForTrace; // Retrieve the method name and report onset event info if tracing is enabled. switch (params.which()) { case rpc::JsRpcTarget::CallParams::METHOD_NAME: { methodNameForTrace = kj::ConstString(kj::str(params.getMethodName())); break; } case rpc::JsRpcTarget::CallParams::METHOD_PATH: { auto path = params.getMethodPath(); auto n = path.size(); if (n == 0) { // Call the target itself as a function. methodNameForTrace = "(this)"_kjc; } else { methodNameForTrace = kj::ConstString(kj::strArray(path, ".")); } break; } } maybeSetJsRpcInfo(ctx, methodNameForTrace); auto targetInfo = getTargetInfo(lock, ctx); // We will try to get the function, if we can't we'll throw an error to the client. auto [propHandle, thisArg] = tryGetProperty(lock, targetInfo.target, params, targetInfo.allowInstanceProperties, ctx); auto op = params.getOperation(); auto handleResult = [&](InvocationResult&& invocationResult) { // Given a handle for the result, if it's a promise, await the promise, then serialize the // final result for return. RpcSerializerExternalHandler::GetStreamHandlerFunc getResultsStreamHandlerFunc; auto resultStreamHandler = params.getResultsStreamHandler(); switch (resultStreamHandler.which()) { case rpc::JsRpcTarget::CallParams::ResultsStreamHandler::EXTERNAL_PUSHER: getResultsStreamHandlerFunc.init( [cap = resultStreamHandler.getExternalPusher()]() mutable { return kj::mv(cap); }); break; case rpc::JsRpcTarget::CallParams::ResultsStreamHandler::STREAM_SINK: getResultsStreamHandlerFunc.init( [cap = resultStreamHandler.getStreamSink()]() mutable { return kj::mv(cap); }); break; } kj::Maybe>> callPipelineFulfiller; // We need another ref to this fulfiller for the error callback. It can rely on being // destroyed at the same time as the success callback. kj::Maybe&> callPipelineFulfillerRef; KJ_IF_SOME(ss, invocationResult.streamSink) { // Since we have a StreamSink, it's important that we hook up the pipeline for that // immediately. Annoyingly, that also means we need to hook up a pipeline for // callPipeline, which we don't actually have yet, so we need to promise-ify it. // If the caller requested using ExternalPusher for the results, then it should also use // ExternalPusher for the params. (Theoretically we could support mix-and-match but... // let's keep it simple.) KJ_REQUIRE(resultStreamHandler.isStreamSink(), "RPC params used StreamSink when result is supposed to use ExternalPusher"); auto paf = kj::newPromiseAndFulfiller(); callPipelineFulfillerRef = *paf.fulfiller; callPipelineFulfiller = kj::mv(paf.fulfiller); capnp::PipelineBuilder builder(16); builder.setCallPipeline(kj::mv(paf.promise)); builder.setParamsStreamSink(ss); callContext.setPipeline(builder.build()); } // HACK: Cap'n Proto call contexts are documented as being pointer-like types where the // backing object's lifetime is that of the RPC call, but in reality they are refcounted // under the hood. Since we'll be executing the call in the JS microtask queue, we have no // ability to actually cancel execution if a cancellation arrives over RPC, and at the end of // that execution we're going to access the call context to write the results. We could // invent some complicated way to skip initializing results in the case the call has been // canceled, but it's easier and safer to just grab a refcount on the call context object // itself, which fully protects us. So... do that. auto ownCallContext = capnp::CallContextHook::from(callContext).addRef(); auto result = ctx.awaitJs(js, js.toPromise(invocationResult.returnValue) .then(js, ctx.addFunctor( // Warning: Be careful about captures here! If the incoming RPC is canceled, // this continuation will still execute, sice it's a JS promise continuation. // But `this` could have been destroyed in the meantime. So all our captures // must take full ownership. [callContext, ownCallContext = kj::mv(ownCallContext), paramDisposalGroup = kj::mv(invocationResult.paramDisposalGroup), paramsStreamSink = kj::mv(invocationResult.streamSink), getResultsStreamHandlerFunc = kj::mv(getResultsStreamHandlerFunc), callPipelineFulfiller = kj::mv(callPipelineFulfiller)]( jsg::Lock& js, jsg::Value value) mutable { jsg::JsValue resultValue(value.getHandle(js)); rpc::JsRpcTarget::CallResults::Builder results = nullptr; auto maybePipeline = serializeJsValueWithPipeline(js, resultValue, [&](capnp::MessageSize hint) { hint.wordCount += capnp::sizeInWords(); hint.capCount += 1; // for callPipeline results = callContext.initResults(hint); return results.initResult(); }, kj::mv(getResultsStreamHandlerFunc)); KJ_SWITCH_ONEOF(maybePipeline) { KJ_CASE_ONEOF(obj, MakeCallPipeline::Object) { results.setCallPipeline(kj::mv(obj.cap)); // Note that hasDisposer is ONLY meant to indicate the presence of an // application-level disposer. It need not be true if we only have stub disposers. results.setHasDisposer(obj.hasDispose); } KJ_CASE_ONEOF(obj, MakeCallPipeline::SingleStub) { // Serialization should have produced a single stub. We can use that same stub as // the callPipeline. auto externals = results.asReader().getResult().getExternals(); KJ_ASSERT(externals.size() == 1); auto external = externals[0]; KJ_ASSERT(external.isRpcTarget()); results.setCallPipeline(external.getRpcTarget()); } KJ_CASE_ONEOF(nonPipelinable, MakeCallPipeline::NonPipelinable) { results.setCallPipeline(kj::mv(nonPipelinable.errorPipeline)); // leave hasDisposer false } } KJ_IF_SOME(cpf, callPipelineFulfiller) { cpf->fulfill(results.getCallPipeline()); } KJ_IF_SOME(ss, paramsStreamSink) { results.setParamsStreamSink(kj::mv(ss)); } // paramDisposalGroup will be destroyed when we return (or when this lambda is destroyed // as a result of the promise being rejected). This will implicitly dispose the param // stubs. }), ctx.addFunctor([callPipelineFulfillerRef](jsg::Lock& js, jsg::Value&& error) { // If we set up a `callPipeline` early, we have to make sure it propagates the error. // (Otherwise we get a PromiseFulfiller error instead, which is pretty useless...) KJ_IF_SOME(cpf, callPipelineFulfillerRef) { cpf.reject(js.exceptionToKj(error.addRef(js))); } js.throwException(kj::mv(error)); }))); if (ctx.hasOutputGate()) { // Note: If `ctx` is destroyed, the entire call to `callImpl()` will be canceled // (makeReentryCallback() ensures this). This does NOT cancel the JavaScript (because JS // promises are not RAII-cancelable), but it will cancel this trailing .then(), which is // why it's safe to capture `&ctx` here. return result.then([&ctx]() mutable { return ctx.waitForOutputLocks(); }); } else { return result; } }; switch (op.which()) { case rpc::JsRpcTarget::CallParams::Operation::CALL_WITH_ARGS: { // Note that using isFunctionForRpc(js, propHandle) here would be incorrect, since that // decides whether it is a function *that can be serialized as a stub*. JsRpcProperty // is (at present) considered non-serializable in itself, but when traversing the // pipeline path, we may have descended into a stub and its properties, thus we could // actually be invoking a JsRpcProperty here. As long as it is in fact callable, we will // allow it. JSG_REQUIRE(propHandle->IsFunction(), TypeError, kj::str("\"", methodNameForTrace, "\" is not a function.")); auto fn = propHandle.As(); kj::Maybe args; if (op.hasCallWithArgs()) { args = op.getCallWithArgs(); } InvocationResult invocationResult; KJ_IF_SOME(envCtx, targetInfo.envCtx) { invocationResult = invokeFnInsertingEnvCtx( js, methodNameForTrace, fn, thisArg, args, envCtx.env, envCtx.ctx); } else { invocationResult = invokeFn(js, fn, thisArg, args); } // We have a function, so let's call it and serialize the result for RPC. // If the function returns a promise we will wait for the promise to finish so we can // serialize the result. return handleResult(kj::mv(invocationResult)); } case rpc::JsRpcTarget::CallParams::Operation::GET_PROPERTY: return handleResult({.returnValue = propHandle}); } KJ_FAIL_ASSERT("unknown JsRpcTarget::CallParams::Operation", (uint)op.which()); } struct GetPropResult { v8::Local handle; v8::Local thisArg; }; [[noreturn]] static void failLookup(kj::StringPtr kjName) { JSG_FAIL_REQUIRE( TypeError, kj::str("The RPC receiver does not implement the method \"", kjName, "\".")); } GetPropResult tryGetProperty(jsg::Lock& js, jsg::JsObject object, rpc::JsRpcTarget::CallParams::Reader callParams, bool allowInstanceProperties, IoContext& ctx) { auto prototypeOfObject = KJ_ASSERT_NONNULL(js.obj().getPrototype(js).tryCast()); // Get the named property of `object`. auto getProperty = [&](kj::StringPtr kjName) { JSG_REQUIRE(!isReservedName(kjName), TypeError, kj::str("'", kjName, "' is a reserved method and cannot be called over RPC.")); jsg::JsValue jsName = js.strIntern(kjName); if (allowInstanceProperties) { // This is a simple object. Its own properties are considered to be accessible over RPC, but // inherited properties (i.e. from Object.prototype) are not. if (!object.has(js, jsName, jsg::JsObject::HasOption::OWN)) { failLookup(kjName); } return object.get(js, jsName); } else { // This is an instance of a valid RPC target class. if (object.has(js, jsName, jsg::JsObject::HasOption::OWN)) { // We do NOT allow own properties, only class properties. failLookup(kjName); } auto value = object.get(js, jsName); if (value == prototypeOfObject.get(js, jsName)) { // This property is inherited from the prototype of `Object`. Don't allow. failLookup(kjName); } return value; } }; kj::Maybe result; switch (callParams.which()) { case rpc::JsRpcTarget::CallParams::METHOD_NAME: { result = getProperty(callParams.getMethodName()); break; } case rpc::JsRpcTarget::CallParams::METHOD_PATH: { auto path = callParams.getMethodPath(); auto n = path.size(); if (n == 0) { // Call the target itself as a function. result = object; } else { bool inStub = false; for (auto i: kj::zeroTo(n - 1)) { // For each property name except the last, look up the property and replace `object` // with it. kj::StringPtr name = path[i]; auto next = getProperty(name); KJ_IF_SOME(o, next.tryCast()) { object = o; } else { // Not an object, doesn't have further properties. failLookup(name); } // If the object is a Proxy, then `isInstanceOf()` won't actually work, // because the Proxy is not an instance of any native type. But for our purposes, // RpcTarget is only a marker used to indicate what semantics are desired. bool isProxyOfRpcTarget = false; if (jsg::JsValue(object).isProxy()) { // Unfortunatley in this case we need to follow the prototype chain manually, looking // for `JsRpcTarget`. js.withinHandleScope([&]() { auto proto = object.getPrototype(js); auto prototypeOfRpcTarget = js.getPrototypeFor(); for (;;) { auto objProto = KJ_UNWRAP_OR(proto.tryCast(), break); if (objProto == prototypeOfRpcTarget) { isProxyOfRpcTarget = true; break; } proto = objProto.getPrototype(js); } }); } // Decide whether the new object is a suitable RPC target. if (object.getPrototype(js) == prototypeOfObject) { // Yes. It's a simple object. allowInstanceProperties = true; } else if (isProxyOfRpcTarget || object.isInstanceOf(js)) { // Yes. It's a JsRpcTarget. allowInstanceProperties = false; } else if (object.isInstanceOf(js) || object.isInstanceOf(js) || (inStub && object.isInstanceOf(js))) { // Yes. It's a JsRpcStub or Fetcher. We should allow descending into the stub. // Note that the wildcard property of a stub is a prototype property, not an instance // property, so setting allowInstanceProperties = false here gets the behavior we // want. // TODO(someday): We'll need to support JsRpcPromise here if someday we allow it to // be serialized. allowInstanceProperties = false; // We will only traverse JsRpcProperty if we got there by descending through a // JsRpcStub. At present you can't just pull a property of a stub and return it. inStub = true; } else if (isFunctionForRpc(js, object)) { // Yes. It's a function. allowInstanceProperties = true; } else { failLookup(name); } } result = getProperty(path[n - 1]); } break; } } return { .handle = KJ_ASSERT_NONNULL(result, "unknown CallParams type", (uint)callParams.which()), .thisArg = object, }; } struct InvocationResult { v8::Local returnValue; kj::Maybe> paramDisposalGroup; kj::Maybe streamSink; }; // Deserializes the arguments and passes them to the given function. static InvocationResult invokeFn(jsg::Lock& js, v8::Local fn, v8::Local thisArg, kj::Maybe args) { // We received arguments from the client, deserialize them back to JS. KJ_IF_SOME(a, args) { auto [value, disposalGroup, streamSink] = deserializeJsValue(js, a, "params"_kjc); auto args = KJ_REQUIRE_NONNULL( value.tryCast(), "expected JsArray when deserializing arguments."); // Call() expects a `Local []`... so we populate an array. v8::LocalVector arguments(js.v8Isolate, args.size()); for (size_t i = 0; i < args.size(); ++i) { arguments[i] = args.get(js, i); } InvocationResult result{ .returnValue = jsg::check(fn->Call(js.v8Context(), thisArg, arguments.size(), arguments.data())), .streamSink = kj::mv(streamSink), }; if (!disposalGroup->empty()) { result.paramDisposalGroup = kj::mv(disposalGroup); } return result; } else { return {.returnValue = jsg::check(fn->Call(js.v8Context(), thisArg, 0, nullptr))}; } }; // Like `invokeFn`, but inject the `env` and `ctx` values between the first and second // parameters. Used for service bindings that use functional syntax. static InvocationResult invokeFnInsertingEnvCtx(jsg::Lock& js, kj::StringPtr methodName, v8::Local fn, v8::Local thisArg, kj::Maybe args, v8::Local env, jsg::JsObject ctx) { // Determine the function arity (how many parameters it was declared to accept) by reading the // `.length` attribute. auto arity = js.withinHandleScope([&]() { auto length = jsg::check(fn->Get(js.v8Context(), js.strIntern("length"))); return jsg::check(length->IntegerValue(js.v8Context())); }); // Avoid excessive allocation from a maliciously-set `length`. JSG_REQUIRE(arity >= 0 && arity < 256, TypeError, "RPC function has unreasonable length attribute: ", arity); if (arity < 3) { // If a function has fewer than three arguments, reproduce the historical behavior where // we'd pass the main argument followed by `env` and `ctx` and the undeclared parameters // would just be truncated. arity = 3; } kj::Maybe> paramDisposalGroup; kj::Maybe streamSink; // We're going to pass all the arguments from the client to the function, but we are going to // insert `env` and `ctx`. We assume the last two arguments that the function declared are // `env` and `ctx`, so we can determine where to insert them based on the function's arity. kj::Maybe argsArrayFromClient; size_t argCountFromClient = 0; KJ_IF_SOME(a, args) { auto [value, disposalGroup, ss] = deserializeJsValue(js, a, "paramsNonClass"_kjc); streamSink = kj::mv(ss); auto array = KJ_REQUIRE_NONNULL( value.tryCast(), "expected JsArray when deserializing arguments."); argCountFromClient = array.size(); argsArrayFromClient = kj::mv(array); if (!disposalGroup->empty()) { paramDisposalGroup = kj::mv(disposalGroup); } } // For now, we are disallowing multiple arguments with bare function syntax, due to a footgun: // if you forget to add `env, ctx` to your arg list, then the last arguments from the client // will be replaced with `env` and `ctx`. Probably this would be quickly noticed in testing, // but if you were to accidentally reflect `env` back to the client, it would be a severe // security flaw. JSG_REQUIRE(arity == 3, TypeError, "Cannot call handler function \"", methodName, "\" over RPC because it has the wrong " "number of arguments. A simple function handler can only be called over RPC if it has " "exactly the arguments (arg, env, ctx), where only the first argument comes from the " "client. To support multi-argument RPC functions, use class-based syntax (extending " "WorkerEntrypoint) instead."); JSG_REQUIRE(argCountFromClient == 1, TypeError, "Attempted to call RPC function \"", methodName, "\" with the wrong number of arguments. " "When calling a top-level handler function that is not declared as part of a class, you " "must always send exactly one argument. In order to support variable numbers of " "arguments, the server must use class-based syntax (extending WorkerEntrypoint) " "instead."); v8::LocalVector arguments(js.v8Isolate, kj::max(argCountFromClient + 2, arity)); for (auto i: kj::zeroTo(arity - 2)) { if (argCountFromClient > i) { arguments[i] = KJ_ASSERT_NONNULL(argsArrayFromClient).get(js, i); } else { arguments[i] = js.undefined(); } } arguments[arity - 2] = env; arguments[arity - 1] = ctx; KJ_IF_SOME(a, argsArrayFromClient) { for (size_t i = arity - 2; i < argCountFromClient; ++i) { arguments[i + 2] = a.get(js, i); } } return { .returnValue = jsg::check(fn->Call(js.v8Context(), thisArg, arguments.size(), arguments.data())), .paramDisposalGroup = kj::mv(paramDisposalGroup), .streamSink = kj::mv(streamSink), }; }; }; class TransientJsRpcTarget final: public JsRpcTargetBase { public: TransientJsRpcTarget( jsg::Lock& js, IoContext& ioCtx, jsg::JsObject object, bool allowInstanceProperties = false) : JsRpcTargetBase(ioCtx, MayOutliveIncomingRequest()), handles(ioCtx.addObjectReverse(kj::heap(js, object))), allowInstanceProperties(allowInstanceProperties) { // Check for the existence of a dispose function now so that the destructor doesn't have to // take an isolate lock if there isn't one. auto getResult = object.get(js, js.symbolDispose()); if (getResult.isFunction()) { auto dispose = jsg::V8Ref( js.v8Isolate, v8::Local(getResult).As()); disposeFulfiller = addDisposeTask(js, ioCtx, object, kj::mv(dispose), {}); } } // Use this version of the constructor to pass the dispose function separately. TransientJsRpcTarget(jsg::Lock& js, IoContext& ioCtx, jsg::JsObject object, kj::Maybe> dispose, kj::Vector> stubDisposers, bool allowInstanceProperties = false) : JsRpcTargetBase(ioCtx, MayOutliveIncomingRequest()), handles(ioCtx.addObjectReverse(kj::heap(js, object))), disposeFulfiller(addDisposeTask(js, ioCtx, object, kj::mv(dispose), kj::mv(stubDisposers))), allowInstanceProperties(allowInstanceProperties) {} ~TransientJsRpcTarget() noexcept(false) { KJ_IF_SOME(f, kj::mv(disposeFulfiller)) { f->fulfill(); } } TargetInfo getTargetInfo(Worker::Lock& lock, IoContext& ioCtx) override { return { .target = handles->object.getHandle(lock), .envCtx = kj::none, .allowInstanceProperties = allowInstanceProperties, }; } private: struct Handles { jsg::JsRef object; Handles(jsg::Lock& js, jsg::JsObject object): object(js, object) {} }; // This object could outlive the IoContext (that's why `JsRpcTargetBase` holds a `WeakRef` to the // context). That means hypothetically it could also outlive the isolate. We therefore need to // place these handles in a `ReverseIoOwn` so that if the `IoContext` dies before we do, they are // dropped at that point. ReverseIoOwn handles; // When fulfilled, calls the original object's dispose function. kj::Maybe>> disposeFulfiller; static kj::Maybe>> addDisposeTask(jsg::Lock& js, IoContext& ctx, jsg::JsObject object, kj::Maybe> dispose, kj::Vector> stubDisposers) { if (dispose == kj::none && stubDisposers.empty()) { // Don't bother scheduling disposal if we have neither. return kj::none; } auto obj = jsg::JsRef(js, object); auto [promise, fulfiller] = kj::newPromiseAndFulfiller(); auto jsPromise = ctx.awaitIo(js, kj::mv(promise), [obj = kj::mv(obj), dispose = kj::mv(dispose), stubDiposers = kj::mv(stubDisposers)]( jsg::Lock& js) { KJ_IF_SOME(d, dispose) { jsg::check(d.getHandle(js)->Call(js.v8Context(), obj.getHandle(js), 0, nullptr)); } // Our stub disposers are dropped at the end of this task. }); ctx.addTask(ctx.awaitJs(js, kj::mv(jsPromise))); return kj::mv(fulfiller); } bool allowInstanceProperties; bool isReservedName(kj::StringPtr name) override { if ( // dup() is reserved to duplicate the stub itself, pointing to the same object. name == "dup" || // All JS classes define a method `constructor` on the prototype, but we don't actually // want this to be callable over RPC! name == "constructor") { return true; } return false; } void maybeSetJsRpcInfo(IoContext& ctx, const kj::ConstString& methodNameForTrace) override {} }; // See comment at call site for explanation. static rpc::JsRpcTarget::Client makeJsRpcTargetForSingleLoopbackCall( jsg::Lock& js, jsg::JsObject obj) { // We intentionally do not want to hook up the disposer here since we're not taking ownership // of the object. return rpc::JsRpcTarget::Client(kj::heap( js, IoContext::current(), obj, kj::none, kj::Vector>(), true)); } template MakeCallPipeline::Result serializeJsValueWithPipeline(jsg::Lock& js, jsg::JsValue value, Func makeBuilder, RpcSerializerExternalHandler::GetStreamHandlerFunc getStreamHandlerFunc) { auto maybeDispose = js.withinHandleScope([&]() -> kj::Maybe> { jsg::JsObject obj = KJ_UNWRAP_OR(value.tryCast(), { return kj::none; }); if (obj.getPrototype(js) == js.obj().getPrototype(js)) { // It's a plain object. jsg::JsValue disposeProperty = obj.get(js, js.symbolDispose()); // We don't want the disposer to be serialized, so delete it from the object. (Remember // that a new `dispose()` method will always be added on the client side). obj.delete_(js, js.symbolDispose()); if (disposeProperty.isFunction()) { auto localDispose = v8::Local(disposeProperty).As(); return jsg::V8Ref(js.v8Isolate, localDispose); } } return kj::none; }); auto hasDispose = maybeDispose != kj::none; // Now that we've extracted our dispose function, we can serialize our value. RpcSerializerExternalHandler externalHandler( RpcSerializerExternalHandler::TRANSFER, kj::mv(getStreamHandlerFunc)); serializeJsValue(js, value, externalHandler, kj::mv(makeBuilder)); auto stubDisposers = externalHandler.releaseStubDisposers(); return js.withinHandleScope([&]() -> MakeCallPipeline::Result { jsg::JsObject obj = KJ_UNWRAP_OR(value.tryCast(), { // Primitive value. Return a fake pipeline just so that we get nice errors if someone tries // to pipeline on it. (If we return null, we'll get "called null capability" out of // Cap'n Proto, which will be treated as an internal error.) return MakeCallPipeline::NonPipelinable{ .errorPipeline = rpc::JsRpcTarget::Client(kj::heap( js, IoContext::current(), js.obj(), kj::none, kj::Vector>(), true))}; }); if (obj.getPrototype(js) == js.obj().getPrototype(js)) { // It's a plain object. auto pipeline = kj::heap( js, IoContext::current(), obj, kj::mv(maybeDispose), kj::mv(stubDisposers), true); return MakeCallPipeline::Object{ .cap = rpc::JsRpcTarget::Client(kj::mv(pipeline)), .hasDispose = hasDispose}; } else if (obj.isInstanceOf(js)) { // It's just a stub. It'll serialize as a single stub, obviously. return MakeCallPipeline::SingleStub(); } else if (obj.isInstanceOf(js)) { // It's an RPC target. It will be serialized as a single stub. return MakeCallPipeline::SingleStub(); } else if (isFunctionForRpc(js, obj)) { // It's a plain function. It will be serialized as a single stub. return MakeCallPipeline::SingleStub(); } else if (obj.isInstanceOf(js)) { // It's a plain fetcher. We want to allow pipelining on it, but we also actually need to // serialize it, so we can't use `SingleStub()`. Note we set `allowInstanceProperties` to // `false` here because the wildcard property of a `Fetcher` is a prototype property, and // that's what we want to expose for pipelining. auto pipeline = kj::heap( js, IoContext::current(), obj, kj::mv(maybeDispose), kj::mv(stubDisposers), false); return MakeCallPipeline::Object{ .cap = rpc::JsRpcTarget::Client(kj::mv(pipeline)), .hasDispose = hasDispose}; } else { // Not an RPC object. Could be a String or other serializable types that derive from Object. // Similar to primitive types, we return a fake pipeline for error-handling reasons. // TODO(soon): What if someone returns e.g. a Map with a disposer on it? Should we honor that // disposer? return MakeCallPipeline::NonPipelinable{ .errorPipeline = rpc::JsRpcTarget::Client(kj::heap( js, IoContext::current(), js.obj(), kj::none, kj::Vector>(), true))}; } }); } // RpcStub are allowed to wrap: // * RpcTargets // * Functions // * Plain objects (only when created explicitly via `new RpcStub`) // // This function checks for these and returns: // * kj::none if it's not a valid type to be wrapped in as tub. // * The value for allowInstanceProperties if it is. kj::Maybe checkStubType(jsg::Lock& js, jsg::JsObject handle) { return js.withinHandleScope([&]() -> kj::Maybe { // TODO(perf): We should really cache `prototypeOfObject` somewhere so we don't have to create // an object to get it. (We do this other places in this file, too...) auto prototypeOfObject = KJ_ASSERT_NONNULL(js.obj().getPrototype(js).tryCast()); auto prototypeOfRpcTarget = js.getPrototypeFor(); auto proto = handle.getPrototype(js); if (proto == prototypeOfObject) { // A regular object. Allow access to instance properties. return true; } else { // Walk the prototype chain looking for RpcTarget. // // (Note we can't simply use handle.isInstanceOf() because that doesn't work // correctly for proxies. Since RpcTarget is only used as a marker, we don't really need // the object to be an instance of it -- we just care if it's in the prototype chain, even // if the prototype chain is faked by the Proxy.) // // TODO(someday): Consider whether `new RpcStub(obj)` should work on arbitrary types. This // could be a useful way to say: "I am explicitly opting into treating this like an // RpcTarget even though I do not have the ability to make its type extend RpcTarget." for (;;) { if (proto == prototypeOfRpcTarget) { // An RpcTarget, don't allow instance properties. return false; } KJ_IF_SOME(protoObj, proto.tryCast()) { proto = protoObj.getPrototype(js); } else if (isFunctionForRpc(js, handle)) { // This is NOT an RpcTarget, but it IS callable as a function, so treat it as such. return true; } else { // End of prototype chain, and didn't find RpcTarget. return kj::none; } } } }); } jsg::Ref JsRpcStub::constructor(jsg::Lock& js, jsg::JsObject object) { auto& ioctx = IoContext::current(); bool allowInstanceProperties = JSG_REQUIRE_NONNULL(checkStubType(js, object), TypeError, "RpcStubs can only wrap plain objects, functions, and RpcTarget derivatives."); rpc::JsRpcTarget::Client cap = kj::heap(js, ioctx, object, allowInstanceProperties); return js.alloc(ioctx.addObject(kj::heap(kj::mv(cap)))); } void JsRpcTarget::serialize(jsg::Lock& js, jsg::Serializer& serializer) { // Serialize by effectively creating a `JsRpcStub` around this object and serializing that. // Except we don't actually want to do _exactly_ that, because we do not want to actually create // a `JsRpcStub` locally. So do the important parts of `JsRpcStub::constructor()` followed by // `JsRpcStub::serialize()`. auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError, "Remote RPC references can only be serialized for RPC."); auto externalHandler = dynamic_cast(&handler); JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError, "Remote RPC references can only be serialized for RPC."); // Handle can't possibly be missing during serialization, it's how we got here. auto handle = jsg::JsObject(KJ_ASSERT_NONNULL(JSG_THIS.tryGetHandle(js))); if (externalHandler->getStubOwnership() == RpcSerializerExternalHandler::DUPLICATE) { // This message isn't supposed to take ownership of stubs. What does that mean for an // RpcTarget? You might argue that it means we should never call the disposer. But that's not // really enough: what if the real owner *does* call the disposer, before our stub is done // with it? How do we make sure the RpcTarget stays alive? // // Things get clearer if we look at a real use case: pure-JS Cap'n Web stubs. We don't see // them as stubs (since they are not instances of JsRpcStub). Instead, we see them as // RpcTargets. But we need the semantics to come out the same: when passed as a parameter // to a native RPC call, we need to duplicate the stub, because the original copy might very // well be disposed before we use it. // // How do we duplicate this non-native stub? Well... proper way to duplicate a pure-JS Cap'n // Web stub is, of course, to call its `dup()` method. // // So how about we just do that? If the target has a `dup()` method, we call it, and we take // ownership of the result, instead of taking ownership of the original object. auto dup = handle.get(js, "dup"); KJ_IF_SOME(dupFunc, dup.tryCast()) { auto replacement = dupFunc.call(js, handle); bool replaced = false; // We got a duplicate. Is it still an RpcTarget? KJ_IF_SOME(replacementObj, replacement.tryCast()) { if (replacementObj.isInstanceOf(js)) { // It is! Let's replace our handle with the duplicate! handle = replacementObj; replaced = true; } } JSG_REQUIRE(replaced, DOMDataCloneError, "Couldn't create a stub for the RcpTarget because it has a dup() method which did not " "return another RpcTarget. Either remove the dup() method or make sure it returns an " "RpcTarget."); } else { // If no dup() method was present, then what? // // The pedantic argument would say: we need to throw an exception. But that would lead to a // pretty poor development experience as people would have to mess with adding dup() // methods to all their RpcTargets. // // Another argument might say: we should just use the RpcTarget but never call the disposer // since we don't own it. But that would probably be confusing. People would wonder why their // disposers are never called. // // If someone passes an RpcTarget with no dup() method, but which does have a disposer, as // the argument to an RPC method, *probably* they just want the disposer to be called when // the callee is done with the object. That is, they want us to take ownership after all. If // that is *not* what they want, then they can always implement a dup() method to make it // clear. // // So, we will just "take ownership" of the target after all, and call its disposer. } } rpc::JsRpcTarget::Client cap = kj::heap(js, IoContext::current(), handle); externalHandler->write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable { builder.setRpcTarget(kj::mv(cap)); }); } void RpcSerializerExternalHandler::serializeFunction( jsg::Lock& js, jsg::Serializer& serializer, v8::Local func) { serializer.writeRawUint32(static_cast(rpc::SerializationTag::JS_RPC_STUB)); auto handle = jsg::JsObject(func); // Similar to JsRpcTarget::serialize(), we may need to dup() the function. if (stubOwnership == RpcSerializerExternalHandler::DUPLICATE) { auto dup = handle.get(js, "dup"); KJ_IF_SOME(dupFunc, dup.tryCast()) { auto replacement = dupFunc.call(js, handle); bool replaced = false; // We got a duplicate. Is it still a Function? KJ_IF_SOME(replacementObj, replacement.tryCast()) { if (isFunctionForRpc(js, replacementObj)) { // It is! Let's replace our handle with the duplicate! handle = replacementObj; replaced = true; } } JSG_REQUIRE(replaced, DOMDataCloneError, "Couldn't create a stub for the function because it has a dup() method which did not " "return another function. Either remove the dup() method or make sure it returns a " "function."); } } rpc::JsRpcTarget::Client cap = kj::heap(js, IoContext::current(), handle, true); write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable { builder.setRpcTarget(kj::mv(cap)); }); } void RpcSerializerExternalHandler::serializeProxy( jsg::Lock& js, jsg::Serializer& serializer, v8::Local proxy) { auto handle = jsg::JsObject(proxy); // Proxies are allowed to present themselves as anything that you could pass to `new RpcStub`. // // Note there's an intentional quirk here: If the Proxy presents itself as a plain object, we // wrap it in a stub, rather than serialize the object. This enables the Proxy to continue // intercepting property accesses when they happen, rather than have all the properties accessed // and serialized upfront. However, in retrospect, this may haev been a bad choice, as it means a // Proxy on a plain object cannot have exactly the same behavior as a plain object would have. // Note that apps which explicitly want to prevent a plain object from being serialized over // RPC can simply use `new RpcStub(object)` to explicitly wrap it in a stub -- no need to use // a Proxy for that. auto allowInstanceProperties = JSG_REQUIRE_NONNULL(checkStubType(js, handle), DOMDataCloneError, "Proxy could not be serialized because it is not a valid RPC receiver type. The " "Proxy must emulate either a plain object or an RpcTarget, as indicated by the " "Proxy's prototype chain."); // Similar to JsRpcTarget::serialize(), we may need to dup() the proxy. if (stubOwnership == RpcSerializerExternalHandler::DUPLICATE) { auto dup = handle.get(js, "dup"); KJ_IF_SOME(dupFunc, dup.tryCast()) { auto replacement = dupFunc.call(js, handle); bool replaced = false; // We got a duplicate. Is it still the same type? KJ_IF_SOME(replacementObj, replacement.tryCast()) { KJ_IF_SOME(stubType, checkStubType(js, replacementObj)) { if (stubType == allowInstanceProperties) { // It is! Let's replace our handle with the duplicate! handle = replacementObj; replaced = true; } } } JSG_REQUIRE(replaced, DOMDataCloneError, "Couldn't create a stub for the Proxy because it has a dup() method which did not " "return the same underlying type (RpcTarget or Function) as the Proxy itself represents. " "Either remove the dup() method or make sure it returns an RpcTarget."); } } // Great, we've concluded we can indeed point a stub at this proxy. serializer.writeRawUint32(static_cast(rpc::SerializationTag::JS_RPC_STUB)); rpc::JsRpcTarget::Client cap = kj::heap(js, IoContext::current(), handle, allowInstanceProperties); write([cap = kj::mv(cap)](rpc::JsValue::External::Builder builder) mutable { builder.setRpcTarget(kj::mv(cap)); }); } // JsRpcTarget implementation specific to entrypoints. This is used to deliver the first, top-level // call of an RPC session. class EntrypointJsRpcTarget final: public JsRpcTargetBase { public: EntrypointJsRpcTarget(IoContext& ioCtx, kj::Maybe entrypointName, kj::Maybe versionInfo, Frankenvalue props, kj::Maybe wrapperModule, kj::Maybe> tracer, bool isDynamicDispatch) : JsRpcTargetBase(ioCtx, CantOutliveIncomingRequest()), ioCtx(ioCtx), // Most of the time we don't really have to clone this but it's hard to fully prove, so // let's be safe. entrypointName(entrypointName.map([](kj::StringPtr s) { return kj::str(s); })), versionInfo(kj::mv(versionInfo)), props(kj::mv(props)), wrapperModule(kj::mv(wrapperModule)), tracer(kj::mv(tracer)), isDynamicDispatch(isDynamicDispatch) {} // Override call() to emit the Return event when the top-level RPC call completes. // This marks when the handler returned a value, NOT when all data has been streamed or all // capabilities released. kj::Promise call(CallContext callContext) override { return JsRpcTargetBase::call(kj::mv(callContext)).then([this]() { KJ_IF_SOME(t, ioCtx.getWorkerTracer()) { t.setReturn(ioCtx.now()); } }); } TargetInfo getTargetInfo(Worker::Lock& lock, IoContext& ioCtx) override { jsg::Lock& js = lock; auto handler = KJ_REQUIRE_NONNULL(lock.getExportedHandler(entrypointName, kj::mv(versionInfo), kj::mv(props), ioCtx.getActor(), isDynamicDispatch), "Failed to get handler to worker."); if (handler->missingSuperclass && wrapperModule == kj::none) { // JS RPC is not enabled on the server side, we cannot call any methods. JSG_REQUIRE(FeatureFlags::get(js).getJsRpc(), TypeError, "The receiving Durable Object does not support RPC, because its class was not declared " "with `extends DurableObject`. In order to enable RPC, make sure your class " "extends the special class `DurableObject`, which can be imported from the module " "\"cloudflare:workers\"."); } auto target = jsg::JsObject(handler->self.getHandle(lock)); KJ_IF_SOME(moduleName, wrapperModule) { // We've been asked to apply a wrapper module to the handler. This is a builtin module whose // default export is a class. The class is constructed with the constructor arguments being // the ctx and env objects and the original DO instance. // This mechanism probably won't work very well on anything other than Durable Objects, so // block such usage for now. We could reconsider this if we have a use case in the future. auto& actor = JSG_REQUIRE_NONNULL( ioCtx.getActor(), Error, "Wrapper modules can only be applied to Durable Objects."); auto module = JSG_REQUIRE_NONNULL( js.resolveInternalModule(moduleName), Error, "Unknown internal module: ", moduleName); v8::Local defaultExport = module.get(js, "default"_kj); JSG_REQUIRE(defaultExport->IsFunction(), TypeError, "Internal module's default export is not a function."); auto func = defaultExport.As(); v8::Local args[3] = {actor.getCtx(js), actor.getEnv(js), target}; auto jsContext = js.v8Context(); v8::Local result = jsg::check(func->NewInstance(jsContext, 3, args)); JSG_REQUIRE(result->IsObject(), TypeError, "Internal module wrapper function did not return an object."); target = jsg::JsObject(result.As()); } // clang-format off TargetInfo targetInfo{ .target = target, .envCtx = handler->ctx.map([&](jsg::Ref& execCtx) -> EnvCtx { return { .env = handler->env.getHandle(js), .ctx = lock.getWorker().getIsolate().getApi().wrapExecutionContext(js, execCtx.addRef()), }; }) }; // clang-format on // `targetInfo.envCtx` is present when we're invoking a freestanding function, and therefore // `env` and `ctx` need to be passed as parameters. In that case, we our method lookup // should obviously permit instance properties, since we expect the export is a plain object. // Otherwise, though, the export is a class. In that case, we have set the rule that we will // only allow class properties (aka prototype properties) to be accessed, to avoid // programmers shooting themselves in the foot by forgetting to make their members private. targetInfo.allowInstanceProperties = targetInfo.envCtx != kj::none; return targetInfo; } private: IoContext& ioCtx; kj::Maybe entrypointName; kj::Maybe versionInfo; Frankenvalue props; kj::Maybe wrapperModule; kj::Maybe> tracer; bool isDynamicDispatch; bool isReservedName(kj::StringPtr name) override { if ( // "fetch" and "connect" are treated specially on entrypoints. name == "fetch" || name == "connect" || // These methods are reserved by the Durable Objects implementation. // TODO(someday): Should they be reserved only for Durable Objects, not WorkerEntrypoint? name == "alarm" || name == "webSocketMessage" || name == "webSocketClose" || name == "webSocketError" || // dup() is reserved to duplicate the stub itself, pointing to the same object. name == "dup" || // All JS classes define a method `constructor` on the prototype, but we don't actually // want this to be callable over RPC! name == "constructor") { return true; } return false; } void maybeSetJsRpcInfo(IoContext& ctx, const kj::ConstString& methodNameForTrace) override { KJ_IF_SOME(tracer, ctx.getWorkerTracer()) { tracer.setJsRpcInfo(ctx.getInvocationSpanContext(), ctx.now(), methodNameForTrace); } } }; // A membrane which wraps the top-level JsRpcTarget of an RPC session on the server side. The // purpose of this membrane is to allow only a single top-level call, which then gets a // `CompletionMembrane` wrapped around it. Note that we can't just wrap `CompletionMembrane` around // the top-level object directly because that capability will not be dropped until the RPC session // completes, since it is actually returned as the result of the top-level RPC call, but that // call doesn't return until the `CompletionMembrane` says all capabilities were dropped, so this // would create a cycle. class JsRpcSessionCustomEvent::ServerTopLevelMembrane final: public capnp::MembranePolicy, public kj::Refcounted { public: explicit ServerTopLevelMembrane(kj::Own> doneFulfiller) : completionMembrane(kj::refcounted(kj::mv(doneFulfiller))) {} ~ServerTopLevelMembrane() noexcept(false) { KJ_IF_SOME(cm, completionMembrane) { cm->reject( KJ_EXCEPTION(DISCONNECTED, "JS RPC session canceled without calling an RPC method.")); } } kj::Maybe inboundCall( uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override { if (interfaceId == capnp::typeId()) { // JsRpcTarget::call() auto cm = kj::mv(JSG_REQUIRE_NONNULL( completionMembrane, Error, "Only one RPC method call is allowed on this object.")); completionMembrane = kj::none; return capnp::membrane(kj::mv(target), kj::mv(cm)); } else if (interfaceId == capnp::typeId()) { // ExternalPusher methods // // It's important that we use the same membrane that we'll use for call(), so that // capabilities returned by the ExternalPusher will be wrapped in the membrane, hence they // will be unwrapped when passed back through the membrane again to call(). auto& cm = *JSG_REQUIRE_NONNULL( completionMembrane, Error, "getExternalPusher() must be called before call()"); return capnp::membrane(kj::mv(target), kj::addRef(cm)); } else { KJ_FAIL_ASSERT("unkown interface ID for JsRpcTarget"); } } kj::Maybe outboundCall( uint64_t interfaceId, uint16_t methodId, capnp::Capability::Client target) override { KJ_FAIL_ASSERT("ServerTopLevelMembrane shouldn't have outgoing capabilities"); } kj::Own addRef() override { return kj::addRef(*this); } private: kj::Maybe> completionMembrane; }; kj::Promise JsRpcSessionCustomEvent::run( kj::Own incomingRequest, kj::Maybe entrypointName, kj::Maybe versionInfo, Frankenvalue props, kj::TaskSet& waitUntilTasks, bool isDynamicDispatch) { IoContext& ioctx = incomingRequest->getContext(); incomingRequest->delivered(); KJ_DEFER({ // waitUntil() should allow extending execution on the server side even when the client // disconnects. waitUntilTasks.add(incomingRequest->drain().attach(kj::mv(incomingRequest))); }); EntrypointJsRpcTarget target(ioctx, entrypointName, kj::mv(versionInfo), kj::mv(props), kj::mv(wrapperModule), mapAddRef(incomingRequest->getWorkerTracer()), isDynamicDispatch); capnp::RevocableServer revcableTarget(target); try { auto [donePromise, doneFulfiller] = kj::newPromiseAndFulfiller(); kj::Own topMembrane; if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_SESSION_HANDLE)) { // When using the session handle approach, we don't need the convoluted // `ServerTopLevelMembrane` because the the top-level `JsRpcTarget` is not unnaturally held // open, so it can be treated the same as any other capability in the session. topMembrane = kj::refcounted(kj::mv(doneFulfiller)); } else { topMembrane = kj::refcounted(kj::mv(doneFulfiller)); } capFulfiller->fulfill(capnp::membrane(revcableTarget.getClient(), kj::mv(topMembrane))); // `donePromise` resolves once there are no longer any capabilities pointing between the client // and server as part of this session. co_await donePromise.exclusiveJoin(ioctx.onAbort()); co_return WorkerInterface::CustomEvent::Result{.outcome = EventOutcome::OK}; } catch (...) { // Make sure the top-level capability is revoked with the same exception that `run()` is // throwing, rather than some generic revocation exception. auto e = kj::getCaughtExceptionAsKj(); revcableTarget.revoke(e.clone()); kj::throwFatalException(kj::mv(e)); } } kj::Promise JsRpcSessionCustomEvent::sendRpc( capnp::HttpOverCapnpFactory& httpOverCapnpFactory, capnp::ByteStreamFactory& byteStreamFactory, rpc::EventDispatcher::Client dispatcher) { // We arrange to revoke all capabilities in this session as soon as `sendRpc()` completes or is // canceled. Normally, the server side doesn't return if any capabilities still exist, so this // only makes a difference in the case that some sort of an error occurred. We don't strictly // have to revoke the capabilities as they are probably already broken anyway, but revoking them // helps to ensure that the underlying transport isn't "held open" waiting for the JS garbage // collector to actually collect the JsRpcStub objects. auto revokePaf = kj::newPromiseAndFulfiller(); KJ_DEFER({ if (revokePaf.fulfiller->isWaiting()) { revokePaf.fulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "JS-RPC session canceled")); } }); auto req = dispatcher.jsRpcSessionRequest(); auto sent = req.send(); rpc::JsRpcTarget::Client cap = sent.getTopLevel(); cap = capnp::membrane(kj::mv(cap), kj::refcounted(kj::mv(revokePaf.promise))); // When no more capabilities exist on the connection, we want to proactively cancel the RPC. // This is needed in particular for the case where the client is dropped without making any calls // at all, e.g. because serializing the arguments failed. Unfortunately, simply dropping the // capability obtained through `sent.getTopLevel()` above will not be detected by the server, // because this is a pipeline capability on a call that is still running. So, if we don't // actually cancel the connection client-side, the server will hang open waiting for the initial // top-level call to arrive, and the event will appear never to complete at our end. // // TODO(cleanup): It feels like there's something wrong with the design here. Can we make this // less ugly? auto completionPaf = kj::newPromiseAndFulfiller(); cap = capnp::membrane( kj::mv(cap), kj::refcounted(kj::mv(completionPaf.fulfiller))); this->capFulfiller->fulfill(kj::mv(cap)); auto session = sent.getSession(); // We don't need to await the call itself as `session.whenResolved()` will already propagate any // errors from it. So we can drop the call promise now. // // Note that it would NOT work to use `req.sendForPipeline()` above, since pipelined capabilities // cannot resolve until the call returns, but `sendForPipeline()` explicitly inhibits the return // message. { auto drop = kj::mv(sent); } try { // Wait for `session` to resolve to a null capability. // // Note that this works even if the server is using the "old approach" where it doesn't return // a `session` at all, because in that case the return itself represents the end of the // session, and the response contains a null pointer for `session`, so this does the expected // thing: resolves `session` to null. co_await session.whenResolved().exclusiveJoin(kj::mv(completionPaf.promise)); } catch (...) { auto e = kj::getCaughtExceptionAsKj(); if (revokePaf.fulfiller->isWaiting()) { revokePaf.fulfiller->reject(e.clone()); } kj::throwFatalException(kj::mv(e)); } co_return WorkerInterface::CustomEvent::Result{.outcome = EventOutcome::OK}; } kj::Promise JsRpcSessionCustomEvent::receiveRpc(JsRpcSessionContext context, WorkerInterface& worker, kj::Own ownWorker, kj::Maybe wrapperModule) { // Client wants to start a JS RPC session, we'll dispatch to the WorkerEntrypoint // here and read the capability off the event itself. auto customEvent = kj::heap(WORKER_RPC_EVENT_TYPE, kj::mv(wrapperModule)); auto cap = customEvent->getCap(); if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_SESSION_HANDLE)) { auto promise = worker.customEvent(kj::mv(customEvent)); auto results = context.getResults(capnp::MessageSize{4, 2}); results.setTopLevel(kj::mv(cap)); // Set the returned session capability to resolve to a null capability when the event is // complete. This also neatly arranges that if the session is dropped early, the // `customEvent()` promise is canceled, thus canceling the session. results.setSession(promise.then([ownWorker = kj::mv(ownWorker)](auto outcome) { return rpc::JsRpcSession::Client(nullptr); })); } else { capnp::PipelineBuilder pipelineBuilder; pipelineBuilder.setTopLevel(cap); context.setPipeline(pipelineBuilder.build()); context.getResults().setTopLevel(kj::mv(cap)); co_await worker.customEvent(kj::mv(customEvent)); } } }; // namespace workerd::api