Skip to content
File

Blob: src/workerd/io/external-pusher.h

cpp65 lines
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 
12namespace workerd {
13 
14using 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.
20class 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