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