// 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 #pragma once // Classes for calling a remote Worker/Durable Object's methods from the stub over RPC. // This file contains the generic stub object (JsRpcStub), as well as classes for sending and // delivering the RPC event. // // `JsRpcStub` specifically represents a capability that was introduced as part of some // broader RPC session. `Fetcher`, on the other hand, also supports RPC methods, where each method // call begins a new session (by dispatching a `jsRpcSession` custom event). Service bindings and // Durable Object stubs both extend from `Fetcher`, and so allow such calls. // // See worker-interface.capnp for the underlying protocol. #include #include #include #include #include #include #include namespace workerd::api { // The 32MB limit is based on the fact that Cap'n Proto's default total message size limit is 64MB, // and we want to stay clear of that. // Additionally, considering total memory of the isolate is limited to 128MB a significantly larger // memory might cause unwarrented condemnations and terminations. // Applications which need to move large amounts of data should split the data into several smaller // chunks transmitted through separate calls. constexpr size_t MAX_JS_RPC_MESSAGE_SIZE = 1u << 25; // ExternalHandler used when serializing RPC messages. Serialization functions with which to // handle RPC specially should use this. class RpcSerializerExternalHandler final: public jsg::Serializer::ExternalHandler { public: using GetStreamSinkFunc = kj::Function; using GetExternalPusherFunc = kj::Function; using GetStreamHandlerFunc = kj::OneOf; enum StubOwnership { TRANSFER, DUPLICATE }; // `getStreamSinkFunc` will be called at most once, the first time a stream is encountered in // serialization, to get the StreamSink that should be used. RpcSerializerExternalHandler( StubOwnership stubOwnership, GetStreamHandlerFunc getStreamHandlerFunc) : stubOwnership(stubOwnership), getStreamHandlerFunc(kj::mv(getStreamHandlerFunc)) {} inline StubOwnership getStubOwnership() { return stubOwnership; } using BuilderCallback = kj::Function; // Returns the ExternalPusher for the remote side. Returns kj::none if this serialization is // using the older StreamSink approach, in which case you need to call `writeStream()` instead. kj::Maybe getExternalPusher(); // Add an external. The value is a callback which will be invoked later to fill in the // JsValue::External in the Cap'n Proto structure. The external array cannot be allocated until // the number of externals are known, which is only after all calls to `add()` have completed, // hence the need for a callback. void write(BuilderCallback callback) { externals.add(kj::mv(callback)); } // Like write(), but use this when there is also a stream associated with the external, i.e. // using StreamSink. This returns a capability which will eventually resolve to the stream. // // StreamSink is being replaced by ExternalPusher. You should only call writeStream() if // getExternalPusher() returns kj::none. If ExternalPusher is available, this method will throw. capnp::Capability::Client writeStream(BuilderCallback callback); // Build the final list. capnp::Orphan> build(capnp::Orphanage orphanage); size_t size() { return externals.size(); } // Add an object that will be released once the serialized value is no longer needed to handle // pipelined calls (i.e. when we are serializing a return value). In particular, for each stub // that we found while serializing, we need to make sure its disposer is run later, so the // Own's destructor runs said disposer. // // NOTE: These are called "stub disposers" because they are most commonly used to dispose stubs // that were part of the serialized value, but other kinds of serialized objects could use // this as well. void addStubDisposer(kj::Own disposer) { stubDisposers.add(kj::mv(disposer)); } // Get the list of disposers to be attached to the pipeline kj::Vector> releaseStubDisposers() { return kj::mv(stubDisposers); } // We serialize functions by turning them into RPC stubs. void serializeFunction( jsg::Lock& js, jsg::Serializer& serializer, v8::Local func) override; // We can serialize a Proxy if it happens to wrap RpcTarget. void serializeProxy( jsg::Lock& js, jsg::Serializer& serializer, v8::Local proxy) override; private: StubOwnership stubOwnership; GetStreamHandlerFunc getStreamHandlerFunc; kj::Vector externals; kj::Vector> stubDisposers; kj::Maybe streamSink; kj::Maybe externalPusher; }; class RpcStubDisposalGroup; class StreamSinkImpl; // ExternalHandler used when deserializing RPC messages. Deserialization functions with which to // handle RPC specially should use this. class RpcDeserializerExternalHandler final: public jsg::Deserializer::ExternalHandler { public: // The `streamSink` parameter should be provided if a StreamSink already exists, e.g. when // deserializing results. If omitted, it will be constructed on-demand. RpcDeserializerExternalHandler(capnp::List::Reader externals, RpcStubDisposalGroup& disposalGroup, kj::Maybe streamSink, kj::LiteralStringConst debugContext) : externals(externals), disposalGroup(disposalGroup), streamSink(streamSink), debugContext(debugContext) {} ~RpcDeserializerExternalHandler() noexcept(false); // Read and return the next external. rpc::JsValue::External::Reader read(); // Call immediately after `read()` when reading an external that is associated with a stream. // `stream` is published back to the sender via StreamSink. void setLastStream(capnp::Capability::Client stream); // All stubs deserialized as part of a particular parameter or result set are placed in a // common disposal group so that they can be disposed together. RpcStubDisposalGroup& getDisposalGroup() { return disposalGroup; } // Call after serialization is complete to get the StreamSink that should handle streams found // while deserializing. Returns none if there were no streams. This should only be called if // a `streamSink` was NOT passed to the constructor. kj::Maybe getStreamSink() { return kj::mv(streamSinkCap); } // Return a string literal to include in deserialization errors for debugging. (In particular // this specifies if it's params or return.) kj::LiteralStringConst getDebugContext() { return debugContext; } private: capnp::List::Reader externals; uint i = 0; kj::UnwindDetector unwindDetector; RpcStubDisposalGroup& disposalGroup; kj::Maybe streamSink; kj::Maybe streamSinkCap; kj::LiteralStringConst debugContext; }; // Base class for objects which can be sent over RPC, but doing so actually sends a stub which // makes RPCs back to the original object. class JsRpcTarget: public jsg::Object { public: static jsg::Ref constructor(jsg::Lock& js) { return js.alloc(); } JSG_RESOURCE_TYPE(JsRpcTarget) {} // Serializes to JsRpcStub. void serialize(jsg::Lock& js, jsg::Serializer& serializer); JSG_ONEWAY_SERIALIZABLE(rpc::SerializationTag::JS_RPC_STUB); }; // Common superclass of JsRpcStub and Fetcher, the two types that may serve as the basis for // RPC calls. // // This class is NOT part of the JavaScript class hierarchy (it has no JSG_RESOURCE_TYPE block), // it's only a C++ class used to abstract how to get a capnp client out of the object. class JsRpcClientProvider: public jsg::Object { public: // Get a capnp client that can be used to dispatch one call. // // If this isn't the root object (i.e. this is a JsRpcProperty), the property path starting from // the root object will be appended to `path`. virtual rpc::JsRpcTarget::Client getClientForOneCall( jsg::Lock& js, kj::Vector& path) = 0; }; class JsRpcProperty; // Represents the promise returned by calling an RPC method. We don't use a regular Promise object, // but rather our own custom thenable, so that we can support pipelining on it. class JsRpcPromise: public JsRpcClientProvider { public: // A weak reference to this JsRpcPromise. Unlike the usual WeakRef pattern, though, this ref is // allocated before the promise itself is actually created, and filled in later. This is needed // to solve cyclic initialization challenges in `callImpl()`. struct WeakRef: public kj::AtomicRefcounted { // Note: The contents of `WeakRef` can only be accessed under isolate lock, but `WeakRef`'s // refcount is not protected by any lock, hence why it is AtomicRefcounted. This also implies // that it can be destroyed without a lock. kj::Maybe ref; // This is set true if the JsRpcPromise's dispose() method was explicitly called, in which // case the final result should be considered pre-disposed. bool disposed = false; }; JsRpcPromise(jsg::JsRef inner, kj::Own weakRef, IoOwn pipeline); ~JsRpcPromise() noexcept(false); void resolve(jsg::Lock& js, jsg::JsValue result); void dispose(jsg::Lock& js); rpc::JsRpcTarget::Client getClientForOneCall( jsg::Lock& js, kj::Vector& path) override; // Expect that the call is itself going to return a function... and call that. jsg::Ref call(const v8::FunctionCallbackInfo& args); // Implement standard Promise interface, especially `then()` so that this works as a custom // thenable. // // Note that we intentionally return jsg::JsValue rather than jsg::JsPromise because we actually // do not want the JSG glue to recognize we're returning a promise triggering behavior that pins // the JsRpcPromise in memory until it resolves. It's actually fine if the JsRpcPromise is GC'ed // before the inner promise resolves, because it's just a thin wrapper that delegates to the // inner promise. The inner promise will keep running until it completes, and will invoke all // the continuations then. jsg::JsValue then(jsg::Lock& js, v8::Local handler, jsg::Optional> errorHandler); jsg::JsValue catch_(jsg::Lock& js, v8::Local errorHandler); jsg::JsValue finally(jsg::Lock& js, v8::Local onFinally); // Get a nested property, using pipelining. kj::Maybe> getProperty(jsg::Lock& js, kj::String name); JSG_RESOURCE_TYPE(JsRpcPromise) { JSG_DISPOSE(dispose); JSG_CALLABLE(call); JSG_WILDCARD_PROPERTY(getProperty); JSG_METHOD(then); JSG_METHOD_NAMED(catch, catch_); JSG_METHOD(finally); } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("inner", inner); } private: jsg::JsRef inner; kj::Own weakRef; struct Pending { IoOwn pipeline; }; struct Resolved { jsg::Value result; // Dummy IoPtr to self, used only to verify that we're running in the correct context. // (Dereferencing from the wrong context would throw an exception.) // Note: Can't use IoContext::WeakRef here because it's not thread-safe (it's only intended to // be held from KJ I/O objects, but this is a JSG object). IoPtr ctxCheck; }; struct Disposed {}; // Note we don't have a "rejected" state because it works fine to just leave the state as // "Pending" -- calls to `pipeline` will rethrow the same exception, and holding the pipeline // open won't actually hold anything open on the server. kj::OneOf state; void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(inner); KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(pending, Pending) {} KJ_CASE_ONEOF(resolved, Resolved) { visitor.visit(resolved.result); } KJ_CASE_ONEOF(disposed, Disposed) {} } } }; // Represents a property -- possibly, a method -- of a remote RPC object. class JsRpcProperty: public JsRpcClientProvider { public: JsRpcProperty(jsg::Ref parent, kj::String name) : parent(kj::mv(parent)), name(kj::mv(name)) {} rpc::JsRpcTarget::Client getClientForOneCall( jsg::Lock& js, kj::Vector& path) override; // Call the property as a method. jsg::Ref call(const v8::FunctionCallbackInfo& args); // Treat the property as a promise to obtain the value. // // Note that we intentionally return jsg::JsValue rather than jsg::JsPromise because we actually // do not want the JSG glue to recognize we're returning a promise triggering behavior that pins // the JsRpcProperty in memory until it resolves. It's actually fine if the JsRpcProperty is GC'ed // before the promise resolves, since the property is just an API stub. The underlying Cap'n Proto // RPCs it starts will keep running; Cap'n Proto refcounts all the necessary resources internally. jsg::JsValue then(jsg::Lock& js, v8::Local handler, jsg::Optional> errorHandler); jsg::JsValue catch_(jsg::Lock& js, v8::Local errorHandler); jsg::JsValue finally(jsg::Lock& js, v8::Local onFinally); // Get a nested property, using pipelining. kj::Maybe> getProperty(jsg::Lock& js, kj::String name); JSG_RESOURCE_TYPE(JsRpcProperty) { // You can call the property as a function. We'll assume it is a method in this case. JSG_CALLABLE(call); // You can access further nested properties. We'll assume the property is an object in this // case. JSG_WILDCARD_PROPERTY(getProperty); // You can treat the property as a promise. This returns the value of the property. JSG_METHOD(then); JSG_METHOD_NAMED(catch, catch_); JSG_METHOD(finally); } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("parent", parent); tracker.trackField("name", name); } private: // The parent object from which this property was obtained. jsg::Ref parent; // Name of this property within its immediate parent. kj::String name; void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(parent); } }; // A JsRpcStub object forwards JS method calls to the remote Worker/Durable Object over RPC. // Since methods are not known until runtime, JsRpcStub doesn't define any JS methods. // Instead, we use JSG_WILDCARD_PROPERTY to intercept property accesses of names that are not known // at compile time. // // JsRpcStub only supports method calls. You cannot, for instance, access a property of a // Durable Object over RPC. // // The `JsRpcStub` type is used to represent capabilities passed across some previous JS RPC // call. It is NOT the type of a Durable Object stub nor a service binding. Those are instances of // `Fetcher`, which has a `getRpcMethod()` call of its own that mostly delegates to // `JsRpcStub::sendJsRpc()`. class JsRpcStub: public JsRpcClientProvider { public: // Only deserialize() needs ExternalMemoryAdjustment — it extracts a new capability from an RPC // response, creating MembraneHook + forked promise allocations. The other callers (dup() and // constructor()) just refcount existing capabilities or create local loopbacks. JsRpcStub(IoOwn capnpClient): capnpClient(kj::mv(capnpClient)) {} JsRpcStub(IoOwn capnpClient, RpcStubDisposalGroup& disposalGroup, jsg::ExternalMemoryAdjustment externalMemoryAdjustment); ~JsRpcStub() noexcept(false); rpc::JsRpcTarget::Client getClient(); rpc::JsRpcTarget::Client getClientForOneCall( jsg::Lock& js, kj::Vector& path) override; jsg::Ref dup(jsg::Lock& js); void dispose(); // Given a JsRpcTarget, make an RPC stub from it. // // Usually, applications won't use this constructor directly. Rather, they will define types // that extend `JsRpcTarget` and then they will simply return those. The serializer will // automatically handle `JsRpcTarget` by wrapping it in `JsRpcStub`. However, it can be useful // for testing to be able to construct a loopback stub. static jsg::Ref constructor(jsg::Lock& js, jsg::JsObject object); // Call the stub itself as a function. jsg::Ref call(const v8::FunctionCallbackInfo& args); kj::Maybe> getRpcMethod(jsg::Lock& js, kj::String name); JSG_RESOURCE_TYPE(JsRpcStub) { JSG_METHOD(dup); JSG_DISPOSE(dispose); JSG_CALLABLE(call); JSG_WILDCARD_PROPERTY(getRpcMethod); } void serialize(jsg::Lock& js, jsg::Serializer& serializer); static jsg::Ref deserialize( jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer); JSG_SERIALIZABLE(rpc::SerializationTag::JS_RPC_STUB); private: // Nulled out upon dispose(). kj::Maybe> capnpClient; kj::Maybe disposalGroup; kj::ListLink disposalGroupLink; kj::Maybe externalMemoryAdjustment; friend class RpcStubDisposalGroup; }; class RpcStubDisposalGroup { public: ~RpcStubDisposalGroup() noexcept(false); // Release all the stubs in the group without disposing them. They will have to be disposed // individually by calling their disposers directly. void disownAll(); // Call dispose() on every stub in the group. void disposeAll(); bool empty() { return list.empty(); } // When creating a disposal group representing an RPC response, we may also attach the // `callPipeline` from the response, to control when the server-side `dispose()` method is // invoked. This isn't part of any stub, it's just discarded upon disposal. void setCallPipeline(IoOwn value) { callPipeline = kj::mv(value); } private: kj::List list; kj::Maybe> callPipeline; friend class JsRpcStub; }; // `jsRpcSession` returns a capability that provides the client a way to call remote methods // over RPC. We drain the IncomingRequest after the capability is used to run the relevant JS. class JsRpcSessionCustomEvent final: public WorkerInterface::CustomEvent { public: JsRpcSessionCustomEvent(uint16_t typeId, kj::Maybe wrapperModule = kj::none, kj::PromiseFulfillerPair paf = kj::newPromiseAndFulfiller()) : capFulfiller(kj::mv(paf.fulfiller)), clientCap(kj::mv(paf.promise)), typeId(typeId), wrapperModule(kj::mv(wrapperModule)) {} ~JsRpcSessionCustomEvent() noexcept(false) { if (capFulfiller->isWaiting()) { capFulfiller->reject( KJ_EXCEPTION(DISCONNECTED, "JsRpcSessionCustomEvent was destroyed before completion")); } } kj::Promise run(kj::Own incomingRequest, kj::Maybe entrypointName, kj::Maybe versionInfo, Frankenvalue props, kj::TaskSet& waitUntilTasks, bool isDynamicDispatch) override; kj::Promise sendRpc(capnp::HttpOverCapnpFactory& httpOverCapnpFactory, capnp::ByteStreamFactory& byteStreamFactory, rpc::EventDispatcher::Client dispatcher) override; // Same as `EventDispatcher::Server::JsRpcSessionContext` -- but that typedef is `protected`. using JsRpcSessionContext = capnp::CallContext; // Common implemnetation of EventDispatcher::jsRpcSession(). static kj::Promise receiveRpc(JsRpcSessionContext context, kj::Own worker, kj::Maybe wrapperModule = kj::none) { auto& ref = *worker; return receiveRpc(context, ref, kj::mv(worker), kj::mv(wrapperModule)); } // Sometimes callers need to pass an `Own` owning some parent object of `worker`. static kj::Promise receiveRpc(JsRpcSessionContext context, WorkerInterface& worker, kj::Own ownWorker, kj::Maybe wrapperModule = kj::none); uint16_t getType() override { return typeId; } tracing::EventInfo getEventInfo() const override { return tracing::JsRpcEventInfo(nullptr); } rpc::JsRpcTarget::Client getCap() { auto result = kj::mv(KJ_ASSERT_NONNULL(clientCap, "can only call getCap() once")); clientCap = kj::none; return result; } kj::Promise notSupported() override { JSG_FAIL_REQUIRE(TypeError, "The receiver is not an RPC object"); } void failed(const kj::Exception& e) override { capFulfiller->reject(e.clone()); } // Event ID for jsRpcSession. // // Similar to WebSocket hibernation, we define this event ID in the internal codebase, but since // we don't create JsRpcSessionCustomEvent from our internal code, we can't pass the event // type in -- so we hardcode it here. static constexpr uint16_t WORKER_RPC_EVENT_TYPE = 9; private: kj::Own> capFulfiller; // We need to set the client/server capability on the event itself to get around CustomEvent's // limited return type. kj::Maybe clientCap; uint16_t typeId; kj::Maybe wrapperModule; class ServerTopLevelMembrane; }; #define EW_WORKER_RPC_ISOLATE_TYPES \ api::JsRpcPromise, api::JsRpcProperty, api::JsRpcStub, api::JsRpcTarget }; // namespace workerd::api