Skip to content
File

Blob: src/workerd/api/streams/readable-source.h

cpp241 lines
1#pragma once
2 
3#include <workerd/io/worker-interface.capnp.h>
4#include <workerd/util/strong-bool.h>
5 
6#include <kj/debug.h>
7 
8namespace kj {
9class AsyncInputStream;
10}
11 
12namespace workerd {
13 
14class IoContext;
15namespace api {
16template <typename T>
17struct DeferredProxy;
18class ReadableStreamSource;
19} // namespace api
20 
21namespace jsg {
22class Lock;
23} // namespace jsg
24 
25namespace api::streams {
26 
27class WritableSink;
28 
29WD_STRONG_BOOL(EndAfterPump);
30 
31// A ReadableSource is primarily intended to serve as a bridge between kj::AsyncInputStream
32// and the ReadableStream API. However, it can also be used directly by KJ-space code that needs
33// deferred proxying. While ReadableSource should probably have been a more JS-friendly
34// API, it's a bit too late to change that now. Use the ReadableSourceJsAdapter in the
35// readable-source-adapter.h file to wrap a ReadableSource for use from JavaScript.
36//
37// A ReadableSource must be treated like a KJ I/O object. Instances that are held
38// by any JS-heap objects must be held by an IoOwn.
39//
40// If the ReadableSource is canceled or dropped, all pending read() reads will be
41// canceled.
42//
43// Only one read() may be pending at a time. Attempting to initiate a second read()
44// while one is already pending will result in a rejected promise.
45//
46// Calling pumpTo initiates a sequence of read() calls until the stream is fully consumed.
47// Ownership of the underlying AsyncInputStream is transferred to the pumpTo operation and
48// the ReadableSource is put into a closed state. After calling pumpTo, no further
49// read() calls may be made directly on the ReadableSource. Dropping the returned
50// promise before it resolves will cancel the pump operation.
51//
52// It is **NOT** intended that you should implement this interface for general use.
53// It is only intended to be implemented by specific classes within workerd for the
54// purpose of bridging between kj/js streams. Streams that operate at the kj level
55// should implement the kj::Async*Stream interfaces, and streams that operate at
56// the JS level should implement to the UnderlyingSource interface. This is a
57// departure from what we've done previously but as part of the effort to simplify
58// the streams code, the goal is to reduce the number of different stream interfaces
59// that we implement to.
60class ReadableSource {
61 public:
62 // Read into the given buffer, returning a promise that resolves to the number of bytes read.
63 // The maximum number of bytes that will be read is the size of the buffer. The minimum number
64 // of bytes that will be read is minBytes. If at least minBytes cannot be read, the promise
65 // will be resolved with the number of bytes read and the stream will be closed.
66 virtual kj::Promise<size_t> read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes = 1) = 0;
67 
68 // If `end` is true, then `output.end()` will be called after pumping. Note that it's especially
69 // important to take advantage of this when using deferred proxying since calling `end()`
70 // directly might attempt to use the `IoContext` to call `registerPendingEvent()`.
71 // If the pump fails, the ReadableSource will be left in an errored state. The
72 // default implementation uses read() to read chunks of data and write them to the output
73 // using a 16KB buffer.
74 //
75 // Per the contract of pumpTo(), it is the caller's responsibility to ensure that both
76 // the WritableStreamSink and this ReadableSource remain alive until the returned
77 // promise resolves!
78 //
79 // It is the caller's responsibility to ensure that WritableStreamSink and this
80 // ReadableSource remain alive until the wrapped deferred proxy task resolves.
81 // The default implementation does arrange to make it safe to drop the source once
82 // the pump begins but that's only precautionary/defensive. It's still better/safer
83 // for the caller to keep the source alive until the pump completes.
84 virtual kj::Promise<DeferredProxy<void>> pumpTo(
85 WritableSink& output, EndAfterPump end = EndAfterPump::YES) = 0;
86 
87 // If the stream is still active, and the encoding matches an encoding that the stream
88 // can provide, gets the total length, if known. If the length is not known, or the
89 // encoding does not match the encoding of the underlying stream, or the stream is closed
90 // or errored, returns kj::none.
91 virtual kj::Maybe<size_t> tryGetLength(rpc::StreamEncoding encoding) = 0;
92 
93 // Fully consume the stream and return all of its data as a byte array. The limit
94 // parameter is the maximum number of bytes to read. If the stream contains more
95 // than this number of bytes, the promise will reject with an exception.
96 virtual kj::Promise<kj::Array<const kj::byte>> readAllBytes(size_t limit) = 0;
97 
98 // Fully consume the stream and return all of its data as a string. The limit
99 // parameter is the maximum number of bytes to read. If the stream contains more
100 // than this number of bytes, the promise will reject with an exception.
101 virtual kj::Promise<kj::String> readAllText(size_t limit) = 0;
102 
103 // Cancels the underlying source if it is still active. Must put the stream into an
104 // errored state. After calling this, all pending and future reads should fail.
105 // Dropping the ReadableSource without calling cancel() first should trigger
106 // cancel() with a generic exception.
107 virtual void cancel(kj::Exception reason) = 0;
108 
109 struct Tee {
110 kj::Own<ReadableSource> branch1;
111 kj::Own<ReadableSource> branch2;
112 };
113 
114 // Tees the stream into two branches. The returned Tee contains two new ReadableSource
115 // instances that will each receive the same data. Once this is called, this instance is no
116 // longer usable and will behave as if it has been closed.
117 // The limit parameter specifies the maximum buffer size to use when teeing.
118 virtual Tee tee(size_t limit) = 0;
119 
120 // Gets the encoding of the stream.
121 virtual rpc::StreamEncoding getEncoding() = 0;
122};
123 
124// Utility base class for ReadableSource wrappers that delegate all
125// operations to an inner ReadableSource while selectively overriding
126// some operations.
127class ReadableSourceWrapper: public ReadableSource {
128 public:
129 KJ_DISALLOW_COPY_AND_MOVE(ReadableSourceWrapper);
130 virtual ~ReadableSourceWrapper() noexcept(false) = default;
131 
132 kj::Promise<size_t> read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes = 1) override {
133 return getInner().read(buffer, minBytes);
134 }
135 
136 kj::Promise<DeferredProxy<void>> pumpTo(
137 WritableSink& output, EndAfterPump end = EndAfterPump::YES) override {
138 return getInner().pumpTo(output, end);
139 }
140 
141 kj::Promise<kj::Array<const kj::byte>> readAllBytes(size_t limit) override {
142 return getInner().readAllBytes(limit);
143 }
144 
145 kj::Promise<kj::String> readAllText(size_t limit) override {
146 return getInner().readAllText(limit);
147 }
148 
149 kj::Maybe<size_t> tryGetLength(rpc::StreamEncoding encoding) override {
150 return getInner().tryGetLength(encoding);
151 }
152 
153 void cancel(kj::Exception reason) override {
154 getInner().cancel(kj::mv(reason));
155 }
156 
157 Tee tee(size_t limit) override {
158 return getInner().tee(limit);
159 }
160 
161 rpc::StreamEncoding getEncoding() override {
162 return getInner().getEncoding();
163 }
164 
165 // Releases ownership of the inner ReadableSource. After calling this,
166 // this wrapper becomes unusable.
167 kj::Own<ReadableSource> release() {
168 auto ret = kj::mv(KJ_ASSERT_NONNULL(inner));
169 inner = kj::none;
170 return kj::mv(ret);
171 }
172 
173 protected:
174 ReadableSourceWrapper(kj::Own<ReadableSource> inner): inner(kj::mv(inner)) {}
175 
176 ReadableSource& getInner() {
177 return *KJ_ASSERT_NONNULL(inner);
178 }
179 
180 private:
181 kj::Maybe<kj::Own<ReadableSource>> inner;
182};
183 
184// Creates a ReadableSource that wraps the given kj::AsyncInputStream.
185kj::Own<ReadableSource> newReadableSource(kj::Own<kj::AsyncInputStream> inner);
186 
187// Creates a ReadableSource that is already in the errored state.
188kj::Own<ReadableSource> newErroredReadableSource(kj::Exception exception);
189 
190// Creates a ReadableSource that is already closed and will produce no data.
191kj::Own<ReadableSource> newClosedReadableSource();
192 
193// Creates a ReadableSource that produces the given bytes and then closes.
194// The backing object, if any, is held alive until the stream is closed or canceled.
195// If the backing object is not provided, the bytes are copied.
196kj::Own<ReadableSource> newReadableSourceFromBytes(
197 kj::ArrayPtr<const kj::byte> bytes, kj::Maybe<kj::Own<void>> backing = kj::none);
198 
199// Creates a ReadableSource that wraps the given source and prevents deferred proxying.
200kj::Own<ReadableSource> newIoContextWrappedReadableSource(
201 IoContext& ioctx, kj::Own<ReadableSource> inner);
202 
203// Creates a ReadableSource that calls the given producer function to produce data
204// on each read (useful primarily for testing).
205kj::Own<ReadableSource> newReadableSourceFromProducer(
206 kj::Function<kj::Promise<size_t>(kj::ArrayPtr<kj::byte>, size_t)> producer,
207 kj::Maybe<uint64_t> expectedLength = kj::none);
208 
209// Creates a ReadableSource that decodes the given stream according to the given encoding.
210kj::Own<ReadableSource> newEncodedReadableSource(
211 rpc::StreamEncoding encoding, kj::Own<kj::AsyncInputStream> inner);
212 
213// Wraps a kj::AsyncInputStream returned from a tee() call to ensure that it translates
214// errors into equivalent JS exceptions. Typically this is used when customizing tee() on
215// a ReadableSource implementation.
216kj::Own<kj::AsyncInputStream> wrapTeeBranch(kj::Own<kj::AsyncInputStream> branch);
217 
218// A ReadableStreamSource backed by in-memory data. Unlike newSystemStream() wrapping a
219// newMemoryInputStream(), this implementation does NOT support deferred proxying. This is
220// important when the backing memory has V8 heap provenance (e.g., jsg::BackingStore, Blob data,
221// kj::Array<kj::byte> with a v8::BackingStore attached, etc)
222// since the memory could be freed by GC after the IoContext completes.
223//
224// The `backing` parameter keeps the underlying memory alive for the lifetime of the stream.
225// If not provided, the bytes are copied.
226//
227// TODO(soon): Update to implement ReadableSource instead of ReadableStreamSource.
228// For now this is a ReadableStreamSource for compat with existing code. Once internal.h/c++
229// is updated to use ReadableSource, we will change this also.
230//
231// TODO(cleanup): It would be nice to eventually have some sort of stronger guarantee when
232// deferred proxying can or cannot be used with a stream. Right now it's a bit ad hoc and
233// error-prone. It requires the stream impl to keep track of whether it can be deferred-proxied
234// or not, but in this case, that may be entirely opaque behind the details of the backing memory
235// as is the case with kj::Array<kj::byte> instances that come from the type wrapper system.
236kj::Own<ReadableStreamSource> newMemorySource(
237 kj::ArrayPtr<const kj::byte> bytes, kj::Maybe<kj::Own<void>> backing = kj::none);
238 
239} // namespace api::streams
240} // namespace workerd