File
Blob: src/workerd/io/external-pusher.h
| 1 | // Copyright (c) 2025 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #pragma once |
| 6 | |
| 7 | #include <workerd/io/worker-interface.capnp.h> |
| 8 | |
| 9 | #include <capnp/compat/byte-stream.h> |
| 10 | #include <kj/async-io.h> |
| 11 | |
| 12 | namespace workerd { |
| 13 | |
| 14 | using kj::byte; |
| 15 | |
| 16 | // Implements JsValue.ExternalPusher from worker-interface.capnp. |
| 17 | // |
| 18 | // ExternalPusher allows a remote peer to "push" certain kinds of objects into our address space |
| 19 | // so that they can then be embedded in `JsValue` as `External` values. |
| 20 | class ExternalPusherImpl: public rpc::JsValue::ExternalPusher::Server, public kj::Refcounted { |
| 21 | public: |
| 22 | ExternalPusherImpl(capnp::ByteStreamFactory& byteStreamFactory) |
| 23 | : byteStreamFactory(byteStreamFactory) {} |
| 24 | |
| 25 | using ExternalPusher = rpc::JsValue::ExternalPusher; |
| 26 | |
| 27 | kj::Own<kj::AsyncInputStream> unwrapStream( |
| 28 | ExternalPusher::InputStream::Client cap, kj::LiteralStringConst debugContext); |
| 29 | |
| 30 | // Box which holds the reason why an AbortSignal was aborted. May be either: |
| 31 | // - A serialized V8 value if the signal was aborted from JavaScript. |
| 32 | // - A KJ exception if the connection from the trigger was lost. |
| 33 | using PendingAbortReason = kj::RefcountedWrapper<kj::OneOf<kj::Array<byte>, kj::Exception>>; |
| 34 | |
| 35 | struct AbortSignal { |
| 36 | // Resolves when `reason` has been filled in. |
| 37 | kj::Promise<void> signal; |
| 38 | |
| 39 | // The abort reason box, will be uninitialized until `signal` resolves. |
| 40 | kj::Own<PendingAbortReason> reason; |
| 41 | }; |
| 42 | |
| 43 | AbortSignal unwrapAbortSignal(ExternalPusher::AbortSignal::Client cap); |
| 44 | |
| 45 | kj::Promise<void> pushByteStream(PushByteStreamContext context) override; |
| 46 | kj::Promise<void> pushAbortSignal(PushAbortSignalContext context) override; |
| 47 | |
| 48 | private: |
| 49 | capnp::ByteStreamFactory& byteStreamFactory; |
| 50 | |
| 51 | capnp::CapabilityServerSet<ExternalPusher::InputStream> inputStreamSet; |
| 52 | capnp::CapabilityServerSet<ExternalPusher::AbortSignal> abortSignalSet; |
| 53 | |
| 54 | kj::Promise<kj::Own<kj::AsyncInputStream>> unwrapStreamImpl( |
| 55 | ExternalPusher::InputStream::Client cap, kj::LiteralStringConst debugContext); |
| 56 | |
| 57 | kj::Promise<void> unwrapAbortSignalImpl( |
| 58 | ExternalPusher::AbortSignal::Client cap, kj::Own<PendingAbortReason> pendingReason); |
| 59 | |
| 60 | class InputStreamImpl; |
| 61 | class AbortSignalImpl; |
| 62 | }; |
| 63 | |
| 64 | } // namespace workerd |