File
Blob: src/workerd/api/worker-rpc.h
| 1 | // Copyright (c) 2017-2023 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #pragma once |
| 6 | // Classes for calling a remote Worker/Durable Object's methods from the stub over RPC. |
| 7 | // This file contains the generic stub object (JsRpcStub), as well as classes for sending and |
| 8 | // delivering the RPC event. |
| 9 | // |
| 10 | // `JsRpcStub` specifically represents a capability that was introduced as part of some |
| 11 | // broader RPC session. `Fetcher`, on the other hand, also supports RPC methods, where each method |
| 12 | // call begins a new session (by dispatching a `jsRpcSession` custom event). Service bindings and |
| 13 | // Durable Object stubs both extend from `Fetcher`, and so allow such calls. |
| 14 | // |
| 15 | // See worker-interface.capnp for the underlying protocol. |
| 16 | |
| 17 | #include <workerd/io/io-context.h> |
| 18 | #include <workerd/io/trace.h> |
| 19 | #include <workerd/io/worker-interface.capnp.h> |
| 20 | #include <workerd/jsg/jsg.h> |
| 21 | #include <workerd/jsg/modules-new.h> |
| 22 | #include <workerd/jsg/ser.h> |
| 23 | #include <workerd/jsg/url.h> |
| 24 | |
| 25 | namespace workerd::api { |
| 26 | |
| 27 | // The 32MB limit is based on the fact that Cap'n Proto's default total message size limit is 64MB, |
| 28 | // and we want to stay clear of that. |
| 29 | // Additionally, considering total memory of the isolate is limited to 128MB a significantly larger |
| 30 | // memory might cause unwarrented condemnations and terminations. |
| 31 | // Applications which need to move large amounts of data should split the data into several smaller |
| 32 | // chunks transmitted through separate calls. |
| 33 | constexpr size_t MAX_JS_RPC_MESSAGE_SIZE = 1u << 25; |
| 34 | |
| 35 | // ExternalHandler used when serializing RPC messages. Serialization functions with which to |
| 36 | // handle RPC specially should use this. |
| 37 | class RpcSerializerExternalHandler final: public jsg::Serializer::ExternalHandler { |
| 38 | public: |
| 39 | using GetStreamSinkFunc = kj::Function<rpc::JsValue::StreamSink::Client()>; |
| 40 | using GetExternalPusherFunc = kj::Function<rpc::JsValue::ExternalPusher::Client()>; |
| 41 | using GetStreamHandlerFunc = kj::OneOf<GetStreamSinkFunc, GetExternalPusherFunc>; |
| 42 | |
| 43 | enum StubOwnership { TRANSFER, DUPLICATE }; |
| 44 | |
| 45 | // `getStreamSinkFunc` will be called at most once, the first time a stream is encountered in |
| 46 | // serialization, to get the StreamSink that should be used. |
| 47 | RpcSerializerExternalHandler( |
| 48 | StubOwnership stubOwnership, GetStreamHandlerFunc getStreamHandlerFunc) |
| 49 | : stubOwnership(stubOwnership), |
| 50 | getStreamHandlerFunc(kj::mv(getStreamHandlerFunc)) {} |
| 51 | |
| 52 | inline StubOwnership getStubOwnership() { |
| 53 | return stubOwnership; |
| 54 | } |
| 55 | |
| 56 | using BuilderCallback = kj::Function<void(rpc::JsValue::External::Builder)>; |
| 57 | |
| 58 | // Returns the ExternalPusher for the remote side. Returns kj::none if this serialization is |
| 59 | // using the older StreamSink approach, in which case you need to call `writeStream()` instead. |
| 60 | kj::Maybe<rpc::JsValue::ExternalPusher::Client> getExternalPusher(); |
| 61 | |
| 62 | // Add an external. The value is a callback which will be invoked later to fill in the |
| 63 | // JsValue::External in the Cap'n Proto structure. The external array cannot be allocated until |
| 64 | // the number of externals are known, which is only after all calls to `add()` have completed, |
| 65 | // hence the need for a callback. |
| 66 | void write(BuilderCallback callback) { |
| 67 | externals.add(kj::mv(callback)); |
| 68 | } |
| 69 | |
| 70 | // Like write(), but use this when there is also a stream associated with the external, i.e. |
| 71 | // using StreamSink. This returns a capability which will eventually resolve to the stream. |
| 72 | // |
| 73 | // StreamSink is being replaced by ExternalPusher. You should only call writeStream() if |
| 74 | // getExternalPusher() returns kj::none. If ExternalPusher is available, this method will throw. |
| 75 | capnp::Capability::Client writeStream(BuilderCallback callback); |
| 76 | |
| 77 | // Build the final list. |
| 78 | capnp::Orphan<capnp::List<rpc::JsValue::External>> build(capnp::Orphanage orphanage); |
| 79 | |
| 80 | size_t size() { |
| 81 | return externals.size(); |
| 82 | } |
| 83 | |
| 84 | // Add an object that will be released once the serialized value is no longer needed to handle |
| 85 | // pipelined calls (i.e. when we are serializing a return value). In particular, for each stub |
| 86 | // that we found while serializing, we need to make sure its disposer is run later, so the |
| 87 | // Own<void>'s destructor runs said disposer. |
| 88 | // |
| 89 | // NOTE: These are called "stub disposers" because they are most commonly used to dispose stubs |
| 90 | // that were part of the serialized value, but other kinds of serialized objects could use |
| 91 | // this as well. |
| 92 | void addStubDisposer(kj::Own<void> disposer) { |
| 93 | stubDisposers.add(kj::mv(disposer)); |
| 94 | } |
| 95 | |
| 96 | // Get the list of disposers to be attached to the pipeline |
| 97 | kj::Vector<kj::Own<void>> releaseStubDisposers() { |
| 98 | return kj::mv(stubDisposers); |
| 99 | } |
| 100 | |
| 101 | // We serialize functions by turning them into RPC stubs. |
| 102 | void serializeFunction( |
| 103 | jsg::Lock& js, jsg::Serializer& serializer, v8::Local<v8::Function> func) override; |
| 104 | |
| 105 | // We can serialize a Proxy if it happens to wrap RpcTarget. |
| 106 | void serializeProxy( |
| 107 | jsg::Lock& js, jsg::Serializer& serializer, v8::Local<v8::Proxy> proxy) override; |
| 108 | |
| 109 | private: |
| 110 | StubOwnership stubOwnership; |
| 111 | GetStreamHandlerFunc getStreamHandlerFunc; |
| 112 | |
| 113 | kj::Vector<BuilderCallback> externals; |
| 114 | kj::Vector<kj::Own<void>> stubDisposers; |
| 115 | |
| 116 | kj::Maybe<rpc::JsValue::StreamSink::Client> streamSink; |
| 117 | kj::Maybe<rpc::JsValue::ExternalPusher::Client> externalPusher; |
| 118 | }; |
| 119 | |
| 120 | class RpcStubDisposalGroup; |
| 121 | class StreamSinkImpl; |
| 122 | |
| 123 | // ExternalHandler used when deserializing RPC messages. Deserialization functions with which to |
| 124 | // handle RPC specially should use this. |
| 125 | class RpcDeserializerExternalHandler final: public jsg::Deserializer::ExternalHandler { |
| 126 | public: |
| 127 | // The `streamSink` parameter should be provided if a StreamSink already exists, e.g. when |
| 128 | // deserializing results. If omitted, it will be constructed on-demand. |
| 129 | RpcDeserializerExternalHandler(capnp::List<rpc::JsValue::External>::Reader externals, |
| 130 | RpcStubDisposalGroup& disposalGroup, |
| 131 | kj::Maybe<StreamSinkImpl&> streamSink, |
| 132 | kj::LiteralStringConst debugContext) |
| 133 | : externals(externals), |
| 134 | disposalGroup(disposalGroup), |
| 135 | streamSink(streamSink), |
| 136 | debugContext(debugContext) {} |
| 137 | ~RpcDeserializerExternalHandler() noexcept(false); |
| 138 | |
| 139 | // Read and return the next external. |
| 140 | rpc::JsValue::External::Reader read(); |
| 141 | |
| 142 | // Call immediately after `read()` when reading an external that is associated with a stream. |
| 143 | // `stream` is published back to the sender via StreamSink. |
| 144 | void setLastStream(capnp::Capability::Client stream); |
| 145 | |
| 146 | // All stubs deserialized as part of a particular parameter or result set are placed in a |
| 147 | // common disposal group so that they can be disposed together. |
| 148 | RpcStubDisposalGroup& getDisposalGroup() { |
| 149 | return disposalGroup; |
| 150 | } |
| 151 | |
| 152 | // Call after serialization is complete to get the StreamSink that should handle streams found |
| 153 | // while deserializing. Returns none if there were no streams. This should only be called if |
| 154 | // a `streamSink` was NOT passed to the constructor. |
| 155 | kj::Maybe<rpc::JsValue::StreamSink::Client> getStreamSink() { |
| 156 | return kj::mv(streamSinkCap); |
| 157 | } |
| 158 | |
| 159 | // Return a string literal to include in deserialization errors for debugging. (In particular |
| 160 | // this specifies if it's params or return.) |
| 161 | kj::LiteralStringConst getDebugContext() { |
| 162 | return debugContext; |
| 163 | } |
| 164 | |
| 165 | private: |
| 166 | capnp::List<rpc::JsValue::External>::Reader externals; |
| 167 | uint i = 0; |
| 168 | |
| 169 | kj::UnwindDetector unwindDetector; |
| 170 | RpcStubDisposalGroup& disposalGroup; |
| 171 | |
| 172 | kj::Maybe<StreamSinkImpl&> streamSink; |
| 173 | kj::Maybe<rpc::JsValue::StreamSink::Client> streamSinkCap; |
| 174 | |
| 175 | kj::LiteralStringConst debugContext; |
| 176 | }; |
| 177 | |
| 178 | // Base class for objects which can be sent over RPC, but doing so actually sends a stub which |
| 179 | // makes RPCs back to the original object. |
| 180 | class JsRpcTarget: public jsg::Object { |
| 181 | public: |
| 182 | static jsg::Ref<JsRpcTarget> constructor(jsg::Lock& js) { |
| 183 | return js.alloc<JsRpcTarget>(); |
| 184 | } |
| 185 | |
| 186 | JSG_RESOURCE_TYPE(JsRpcTarget) {} |
| 187 | |
| 188 | // Serializes to JsRpcStub. |
| 189 | void serialize(jsg::Lock& js, jsg::Serializer& serializer); |
| 190 | JSG_ONEWAY_SERIALIZABLE(rpc::SerializationTag::JS_RPC_STUB); |
| 191 | }; |
| 192 | |
| 193 | // Common superclass of JsRpcStub and Fetcher, the two types that may serve as the basis for |
| 194 | // RPC calls. |
| 195 | // |
| 196 | // This class is NOT part of the JavaScript class hierarchy (it has no JSG_RESOURCE_TYPE block), |
| 197 | // it's only a C++ class used to abstract how to get a capnp client out of the object. |
| 198 | class JsRpcClientProvider: public jsg::Object { |
| 199 | public: |
| 200 | // Get a capnp client that can be used to dispatch one call. |
| 201 | // |
| 202 | // If this isn't the root object (i.e. this is a JsRpcProperty), the property path starting from |
| 203 | // the root object will be appended to `path`. |
| 204 | virtual rpc::JsRpcTarget::Client getClientForOneCall( |
| 205 | jsg::Lock& js, kj::Vector<kj::StringPtr>& path) = 0; |
| 206 | }; |
| 207 | |
| 208 | class JsRpcProperty; |
| 209 | |
| 210 | // Represents the promise returned by calling an RPC method. We don't use a regular Promise object, |
| 211 | // but rather our own custom thenable, so that we can support pipelining on it. |
| 212 | class JsRpcPromise: public JsRpcClientProvider { |
| 213 | public: |
| 214 | // A weak reference to this JsRpcPromise. Unlike the usual WeakRef pattern, though, this ref is |
| 215 | // allocated before the promise itself is actually created, and filled in later. This is needed |
| 216 | // to solve cyclic initialization challenges in `callImpl()`. |
| 217 | struct WeakRef: public kj::AtomicRefcounted { |
| 218 | // Note: The contents of `WeakRef` can only be accessed under isolate lock, but `WeakRef`'s |
| 219 | // refcount is not protected by any lock, hence why it is AtomicRefcounted. This also implies |
| 220 | // that it can be destroyed without a lock. |
| 221 | |
| 222 | kj::Maybe<JsRpcPromise&> ref; |
| 223 | |
| 224 | // This is set true if the JsRpcPromise's dispose() method was explicitly called, in which |
| 225 | // case the final result should be considered pre-disposed. |
| 226 | bool disposed = false; |
| 227 | }; |
| 228 | |
| 229 | JsRpcPromise(jsg::JsRef<jsg::JsPromise> inner, |
| 230 | kj::Own<WeakRef> weakRef, |
| 231 | IoOwn<rpc::JsRpcTarget::CallResults::Pipeline> pipeline); |
| 232 | ~JsRpcPromise() noexcept(false); |
| 233 | |
| 234 | void resolve(jsg::Lock& js, jsg::JsValue result); |
| 235 | void dispose(jsg::Lock& js); |
| 236 | |
| 237 | rpc::JsRpcTarget::Client getClientForOneCall( |
| 238 | jsg::Lock& js, kj::Vector<kj::StringPtr>& path) override; |
| 239 | |
| 240 | // Expect that the call is itself going to return a function... and call that. |
| 241 | jsg::Ref<JsRpcPromise> call(const v8::FunctionCallbackInfo<v8::Value>& args); |
| 242 | |
| 243 | // Implement standard Promise interface, especially `then()` so that this works as a custom |
| 244 | // thenable. |
| 245 | // |
| 246 | // Note that we intentionally return jsg::JsValue rather than jsg::JsPromise because we actually |
| 247 | // do not want the JSG glue to recognize we're returning a promise triggering behavior that pins |
| 248 | // the JsRpcPromise in memory until it resolves. It's actually fine if the JsRpcPromise is GC'ed |
| 249 | // before the inner promise resolves, because it's just a thin wrapper that delegates to the |
| 250 | // inner promise. The inner promise will keep running until it completes, and will invoke all |
| 251 | // the continuations then. |
| 252 | jsg::JsValue then(jsg::Lock& js, |
| 253 | v8::Local<v8::Function> handler, |
| 254 | jsg::Optional<v8::Local<v8::Function>> errorHandler); |
| 255 | jsg::JsValue catch_(jsg::Lock& js, v8::Local<v8::Function> errorHandler); |
| 256 | jsg::JsValue finally(jsg::Lock& js, v8::Local<v8::Function> onFinally); |
| 257 | |
| 258 | // Get a nested property, using pipelining. |
| 259 | kj::Maybe<jsg::Ref<JsRpcProperty>> getProperty(jsg::Lock& js, kj::String name); |
| 260 | |
| 261 | JSG_RESOURCE_TYPE(JsRpcPromise) { |
| 262 | JSG_DISPOSE(dispose); |
| 263 | JSG_CALLABLE(call); |
| 264 | JSG_WILDCARD_PROPERTY(getProperty); |
| 265 | JSG_METHOD(then); |
| 266 | JSG_METHOD_NAMED(catch, catch_); |
| 267 | JSG_METHOD(finally); |
| 268 | } |
| 269 | |
| 270 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 271 | tracker.trackField("inner", inner); |
| 272 | } |
| 273 | |
| 274 | private: |
| 275 | jsg::JsRef<jsg::JsPromise> inner; |
| 276 | kj::Own<WeakRef> weakRef; |
| 277 | |
| 278 | struct Pending { |
| 279 | IoOwn<rpc::JsRpcTarget::CallResults::Pipeline> pipeline; |
| 280 | }; |
| 281 | struct Resolved { |
| 282 | jsg::Value result; |
| 283 | |
| 284 | // Dummy IoPtr to self, used only to verify that we're running in the correct context. |
| 285 | // (Dereferencing from the wrong context would throw an exception.) |
| 286 | // Note: Can't use IoContext::WeakRef here because it's not thread-safe (it's only intended to |
| 287 | // be held from KJ I/O objects, but this is a JSG object). |
| 288 | IoPtr<JsRpcPromise> ctxCheck; |
| 289 | }; |
| 290 | struct Disposed {}; |
| 291 | |
| 292 | // Note we don't have a "rejected" state because it works fine to just leave the state as |
| 293 | // "Pending" -- calls to `pipeline` will rethrow the same exception, and holding the pipeline |
| 294 | // open won't actually hold anything open on the server. |
| 295 | kj::OneOf<Pending, Resolved, Disposed> state; |
| 296 | |
| 297 | void visitForGc(jsg::GcVisitor& visitor) { |
| 298 | visitor.visit(inner); |
| 299 | KJ_SWITCH_ONEOF(state) { |
| 300 | KJ_CASE_ONEOF(pending, Pending) {} |
| 301 | KJ_CASE_ONEOF(resolved, Resolved) { |
| 302 | visitor.visit(resolved.result); |
| 303 | } |
| 304 | KJ_CASE_ONEOF(disposed, Disposed) {} |
| 305 | } |
| 306 | } |
| 307 | }; |
| 308 | |
| 309 | // Represents a property -- possibly, a method -- of a remote RPC object. |
| 310 | class JsRpcProperty: public JsRpcClientProvider { |
| 311 | public: |
| 312 | JsRpcProperty(jsg::Ref<JsRpcClientProvider> parent, kj::String name) |
| 313 | : parent(kj::mv(parent)), |
| 314 | name(kj::mv(name)) {} |
| 315 | |
| 316 | rpc::JsRpcTarget::Client getClientForOneCall( |
| 317 | jsg::Lock& js, kj::Vector<kj::StringPtr>& path) override; |
| 318 | |
| 319 | // Call the property as a method. |
| 320 | jsg::Ref<JsRpcPromise> call(const v8::FunctionCallbackInfo<v8::Value>& args); |
| 321 | |
| 322 | // Treat the property as a promise to obtain the value. |
| 323 | // |
| 324 | // Note that we intentionally return jsg::JsValue rather than jsg::JsPromise because we actually |
| 325 | // do not want the JSG glue to recognize we're returning a promise triggering behavior that pins |
| 326 | // the JsRpcProperty in memory until it resolves. It's actually fine if the JsRpcProperty is GC'ed |
| 327 | // before the promise resolves, since the property is just an API stub. The underlying Cap'n Proto |
| 328 | // RPCs it starts will keep running; Cap'n Proto refcounts all the necessary resources internally. |
| 329 | jsg::JsValue then(jsg::Lock& js, |
| 330 | v8::Local<v8::Function> handler, |
| 331 | jsg::Optional<v8::Local<v8::Function>> errorHandler); |
| 332 | jsg::JsValue catch_(jsg::Lock& js, v8::Local<v8::Function> errorHandler); |
| 333 | jsg::JsValue finally(jsg::Lock& js, v8::Local<v8::Function> onFinally); |
| 334 | |
| 335 | // Get a nested property, using pipelining. |
| 336 | kj::Maybe<jsg::Ref<JsRpcProperty>> getProperty(jsg::Lock& js, kj::String name); |
| 337 | |
| 338 | JSG_RESOURCE_TYPE(JsRpcProperty) { |
| 339 | // You can call the property as a function. We'll assume it is a method in this case. |
| 340 | JSG_CALLABLE(call); |
| 341 | |
| 342 | // You can access further nested properties. We'll assume the property is an object in this |
| 343 | // case. |
| 344 | JSG_WILDCARD_PROPERTY(getProperty); |
| 345 | |
| 346 | // You can treat the property as a promise. This returns the value of the property. |
| 347 | JSG_METHOD(then); |
| 348 | JSG_METHOD_NAMED(catch, catch_); |
| 349 | JSG_METHOD(finally); |
| 350 | } |
| 351 | |
| 352 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 353 | tracker.trackField("parent", parent); |
| 354 | tracker.trackField("name", name); |
| 355 | } |
| 356 | |
| 357 | private: |
| 358 | // The parent object from which this property was obtained. |
| 359 | jsg::Ref<JsRpcClientProvider> parent; |
| 360 | |
| 361 | // Name of this property within its immediate parent. |
| 362 | kj::String name; |
| 363 | |
| 364 | void visitForGc(jsg::GcVisitor& visitor) { |
| 365 | visitor.visit(parent); |
| 366 | } |
| 367 | }; |
| 368 | |
| 369 | // A JsRpcStub object forwards JS method calls to the remote Worker/Durable Object over RPC. |
| 370 | // Since methods are not known until runtime, JsRpcStub doesn't define any JS methods. |
| 371 | // Instead, we use JSG_WILDCARD_PROPERTY to intercept property accesses of names that are not known |
| 372 | // at compile time. |
| 373 | // |
| 374 | // JsRpcStub only supports method calls. You cannot, for instance, access a property of a |
| 375 | // Durable Object over RPC. |
| 376 | // |
| 377 | // The `JsRpcStub` type is used to represent capabilities passed across some previous JS RPC |
| 378 | // call. It is NOT the type of a Durable Object stub nor a service binding. Those are instances of |
| 379 | // `Fetcher`, which has a `getRpcMethod()` call of its own that mostly delegates to |
| 380 | // `JsRpcStub::sendJsRpc()`. |
| 381 | class JsRpcStub: public JsRpcClientProvider { |
| 382 | public: |
| 383 | // Only deserialize() needs ExternalMemoryAdjustment — it extracts a new capability from an RPC |
| 384 | // response, creating MembraneHook + forked promise allocations. The other callers (dup() and |
| 385 | // constructor()) just refcount existing capabilities or create local loopbacks. |
| 386 | JsRpcStub(IoOwn<rpc::JsRpcTarget::Client> capnpClient): capnpClient(kj::mv(capnpClient)) {} |
| 387 | JsRpcStub(IoOwn<rpc::JsRpcTarget::Client> capnpClient, |
| 388 | RpcStubDisposalGroup& disposalGroup, |
| 389 | jsg::ExternalMemoryAdjustment externalMemoryAdjustment); |
| 390 | ~JsRpcStub() noexcept(false); |
| 391 | |
| 392 | rpc::JsRpcTarget::Client getClient(); |
| 393 | |
| 394 | rpc::JsRpcTarget::Client getClientForOneCall( |
| 395 | jsg::Lock& js, kj::Vector<kj::StringPtr>& path) override; |
| 396 | |
| 397 | jsg::Ref<JsRpcStub> dup(jsg::Lock& js); |
| 398 | void dispose(); |
| 399 | |
| 400 | // Given a JsRpcTarget, make an RPC stub from it. |
| 401 | // |
| 402 | // Usually, applications won't use this constructor directly. Rather, they will define types |
| 403 | // that extend `JsRpcTarget` and then they will simply return those. The serializer will |
| 404 | // automatically handle `JsRpcTarget` by wrapping it in `JsRpcStub`. However, it can be useful |
| 405 | // for testing to be able to construct a loopback stub. |
| 406 | static jsg::Ref<JsRpcStub> constructor(jsg::Lock& js, jsg::JsObject object); |
| 407 | |
| 408 | // Call the stub itself as a function. |
| 409 | jsg::Ref<JsRpcPromise> call(const v8::FunctionCallbackInfo<v8::Value>& args); |
| 410 | |
| 411 | kj::Maybe<jsg::Ref<JsRpcProperty>> getRpcMethod(jsg::Lock& js, kj::String name); |
| 412 | |
| 413 | JSG_RESOURCE_TYPE(JsRpcStub) { |
| 414 | JSG_METHOD(dup); |
| 415 | JSG_DISPOSE(dispose); |
| 416 | JSG_CALLABLE(call); |
| 417 | JSG_WILDCARD_PROPERTY(getRpcMethod); |
| 418 | } |
| 419 | |
| 420 | void serialize(jsg::Lock& js, jsg::Serializer& serializer); |
| 421 | static jsg::Ref<JsRpcStub> deserialize( |
| 422 | jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer); |
| 423 | |
| 424 | JSG_SERIALIZABLE(rpc::SerializationTag::JS_RPC_STUB); |
| 425 | |
| 426 | private: |
| 427 | // Nulled out upon dispose(). |
| 428 | kj::Maybe<IoOwn<rpc::JsRpcTarget::Client>> capnpClient; |
| 429 | |
| 430 | kj::Maybe<RpcStubDisposalGroup&> disposalGroup; |
| 431 | kj::ListLink<JsRpcStub> disposalGroupLink; |
| 432 | kj::Maybe<jsg::ExternalMemoryAdjustment> externalMemoryAdjustment; |
| 433 | |
| 434 | friend class RpcStubDisposalGroup; |
| 435 | }; |
| 436 | |
| 437 | class RpcStubDisposalGroup { |
| 438 | public: |
| 439 | ~RpcStubDisposalGroup() noexcept(false); |
| 440 | |
| 441 | // Release all the stubs in the group without disposing them. They will have to be disposed |
| 442 | // individually by calling their disposers directly. |
| 443 | void disownAll(); |
| 444 | |
| 445 | // Call dispose() on every stub in the group. |
| 446 | void disposeAll(); |
| 447 | |
| 448 | bool empty() { |
| 449 | return list.empty(); |
| 450 | } |
| 451 | |
| 452 | // When creating a disposal group representing an RPC response, we may also attach the |
| 453 | // `callPipeline` from the response, to control when the server-side `dispose()` method is |
| 454 | // invoked. This isn't part of any stub, it's just discarded upon disposal. |
| 455 | void setCallPipeline(IoOwn<rpc::JsRpcTarget::Client> value) { |
| 456 | callPipeline = kj::mv(value); |
| 457 | } |
| 458 | |
| 459 | private: |
| 460 | kj::List<JsRpcStub, &JsRpcStub::disposalGroupLink> list; |
| 461 | kj::Maybe<IoOwn<rpc::JsRpcTarget::Client>> callPipeline; |
| 462 | friend class JsRpcStub; |
| 463 | }; |
| 464 | |
| 465 | // `jsRpcSession` returns a capability that provides the client a way to call remote methods |
| 466 | // over RPC. We drain the IncomingRequest after the capability is used to run the relevant JS. |
| 467 | class JsRpcSessionCustomEvent final: public WorkerInterface::CustomEvent { |
| 468 | public: |
| 469 | JsRpcSessionCustomEvent(uint16_t typeId, |
| 470 | kj::Maybe<kj::String> wrapperModule = kj::none, |
| 471 | kj::PromiseFulfillerPair<rpc::JsRpcTarget::Client> paf = |
| 472 | kj::newPromiseAndFulfiller<rpc::JsRpcTarget::Client>()) |
| 473 | : capFulfiller(kj::mv(paf.fulfiller)), |
| 474 | clientCap(kj::mv(paf.promise)), |
| 475 | typeId(typeId), |
| 476 | wrapperModule(kj::mv(wrapperModule)) {} |
| 477 | |
| 478 | ~JsRpcSessionCustomEvent() noexcept(false) { |
| 479 | if (capFulfiller->isWaiting()) { |
| 480 | capFulfiller->reject( |
| 481 | KJ_EXCEPTION(DISCONNECTED, "JsRpcSessionCustomEvent was destroyed before completion")); |
| 482 | } |
| 483 | } |
| 484 | |
| 485 | kj::Promise<Result> run(kj::Own<IoContext::IncomingRequest> incomingRequest, |
| 486 | kj::Maybe<kj::StringPtr> entrypointName, |
| 487 | kj::Maybe<Worker::VersionInfo> versionInfo, |
| 488 | Frankenvalue props, |
| 489 | kj::TaskSet& waitUntilTasks, |
| 490 | bool isDynamicDispatch) override; |
| 491 | |
| 492 | kj::Promise<Result> sendRpc(capnp::HttpOverCapnpFactory& httpOverCapnpFactory, |
| 493 | capnp::ByteStreamFactory& byteStreamFactory, |
| 494 | rpc::EventDispatcher::Client dispatcher) override; |
| 495 | |
| 496 | // Same as `EventDispatcher::Server::JsRpcSessionContext` -- but that typedef is `protected`. |
| 497 | using JsRpcSessionContext = capnp::CallContext<rpc::EventDispatcher::JsRpcSessionParams, |
| 498 | rpc::EventDispatcher::JsRpcSessionResults>; |
| 499 | |
| 500 | // Common implemnetation of EventDispatcher::jsRpcSession(). |
| 501 | static kj::Promise<void> receiveRpc(JsRpcSessionContext context, |
| 502 | kj::Own<WorkerInterface> worker, |
| 503 | kj::Maybe<kj::String> wrapperModule = kj::none) { |
| 504 | auto& ref = *worker; |
| 505 | return receiveRpc(context, ref, kj::mv(worker), kj::mv(wrapperModule)); |
| 506 | } |
| 507 | |
| 508 | // Sometimes callers need to pass an `Own<void>` owning some parent object of `worker`. |
| 509 | static kj::Promise<void> receiveRpc(JsRpcSessionContext context, |
| 510 | WorkerInterface& worker, |
| 511 | kj::Own<void> ownWorker, |
| 512 | kj::Maybe<kj::String> wrapperModule = kj::none); |
| 513 | |
| 514 | uint16_t getType() override { |
| 515 | return typeId; |
| 516 | } |
| 517 | |
| 518 | tracing::EventInfo getEventInfo() const override { |
| 519 | return tracing::JsRpcEventInfo(nullptr); |
| 520 | } |
| 521 | |
| 522 | rpc::JsRpcTarget::Client getCap() { |
| 523 | auto result = kj::mv(KJ_ASSERT_NONNULL(clientCap, "can only call getCap() once")); |
| 524 | clientCap = kj::none; |
| 525 | return result; |
| 526 | } |
| 527 | |
| 528 | kj::Promise<Result> notSupported() override { |
| 529 | JSG_FAIL_REQUIRE(TypeError, "The receiver is not an RPC object"); |
| 530 | } |
| 531 | |
| 532 | void failed(const kj::Exception& e) override { |
| 533 | capFulfiller->reject(e.clone()); |
| 534 | } |
| 535 | |
| 536 | // Event ID for jsRpcSession. |
| 537 | // |
| 538 | // Similar to WebSocket hibernation, we define this event ID in the internal codebase, but since |
| 539 | // we don't create JsRpcSessionCustomEvent from our internal code, we can't pass the event |
| 540 | // type in -- so we hardcode it here. |
| 541 | static constexpr uint16_t WORKER_RPC_EVENT_TYPE = 9; |
| 542 | |
| 543 | private: |
| 544 | kj::Own<kj::PromiseFulfiller<workerd::rpc::JsRpcTarget::Client>> capFulfiller; |
| 545 | |
| 546 | // We need to set the client/server capability on the event itself to get around CustomEvent's |
| 547 | // limited return type. |
| 548 | kj::Maybe<rpc::JsRpcTarget::Client> clientCap; |
| 549 | uint16_t typeId; |
| 550 | |
| 551 | kj::Maybe<kj::String> wrapperModule; |
| 552 | |
| 553 | class ServerTopLevelMembrane; |
| 554 | }; |
| 555 | |
| 556 | #define EW_WORKER_RPC_ISOLATE_TYPES \ |
| 557 | api::JsRpcPromise, api::JsRpcProperty, api::JsRpcStub, api::JsRpcTarget |
| 558 | |
| 559 | }; // namespace workerd::api |