// Copyright (c) 2025 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include #include #include namespace workerd { using kj::byte; // Implements JsValue.ExternalPusher from worker-interface.capnp. // // ExternalPusher allows a remote peer to "push" certain kinds of objects into our address space // so that they can then be embedded in `JsValue` as `External` values. class ExternalPusherImpl: public rpc::JsValue::ExternalPusher::Server, public kj::Refcounted { public: ExternalPusherImpl(capnp::ByteStreamFactory& byteStreamFactory) : byteStreamFactory(byteStreamFactory) {} using ExternalPusher = rpc::JsValue::ExternalPusher; kj::Own unwrapStream( ExternalPusher::InputStream::Client cap, kj::LiteralStringConst debugContext); // Box which holds the reason why an AbortSignal was aborted. May be either: // - A serialized V8 value if the signal was aborted from JavaScript. // - A KJ exception if the connection from the trigger was lost. using PendingAbortReason = kj::RefcountedWrapper, kj::Exception>>; struct AbortSignal { // Resolves when `reason` has been filled in. kj::Promise signal; // The abort reason box, will be uninitialized until `signal` resolves. kj::Own reason; }; AbortSignal unwrapAbortSignal(ExternalPusher::AbortSignal::Client cap); kj::Promise pushByteStream(PushByteStreamContext context) override; kj::Promise pushAbortSignal(PushAbortSignalContext context) override; private: capnp::ByteStreamFactory& byteStreamFactory; capnp::CapabilityServerSet inputStreamSet; capnp::CapabilityServerSet abortSignalSet; kj::Promise> unwrapStreamImpl( ExternalPusher::InputStream::Client cap, kj::LiteralStringConst debugContext); kj::Promise unwrapAbortSignalImpl( ExternalPusher::AbortSignal::Client cap, kj::Own pendingReason); class InputStreamImpl; class AbortSignalImpl; }; } // namespace workerd