Skip to content
File

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

cpp440 lines
1#include "common.h"
2#include "readable-source.h"
3#include "readable.h"
4 
5#include <workerd/util/state-machine.h>
6 
7namespace workerd::api::streams {
8 
9// We provide two utility adapters here: ReadableStreamSourceJsAdapter and
10// ReadableSourceKjAdapter.
11//
12// ReadableStreamSourceJsAdapter adapts a ReadableStreamSource to a JavaScript-friendly
13// interface. It provides methods that return JavaScript promises and use
14// JavaScript types. It is intended to be used by JavaScript code that wants
15// to read from a kj-backed stream. It takes ownership of the ReadableStreamSource
16// and holds it with an IoOwn, ensures that all operations are performed on the
17// correct IoContext, and safely cleans up after itself if the adapter is dropped.
18//
19// ┌───────────────────────────────────────────┐
20// │ ReadableStreamSourceJsAdapter │
21// │ │
22// │ ┌─────────────────────────────────────┐ │
23// │ │ JavaScript API │ │
24// │ │ │ │
25// │ │ • read() → Promise<ReadResult> │ │
26// │ │ • readAllText() → Promise<string> │ │
27// │ │ • readAllBytes() → Promise<bytes> │ │
28// │ │ • close() → Promise<void> │ │
29// │ │ • cancel(reason) │ │
30// │ │ • tryTee() → {branch1, branch2} │ │
31// │ └─────────────────────────────────────┘ │
32// │ │ │
33// │ ▼ │
34// │ ┌─────────────────────────────────────┐ │
35// │ │ State Management │ │
36// │ │ │ │
37// │ │ Active ──► Closed │ │
38// │ │ │ │ │ │
39// │ │ │ ▼ │ │
40// │ │ └─────► Canceled/Errored │ │
41// │ └─────────────────────────────────────┘ │
42// │ │ │
43// │ ▼ │
44// │ ┌─────────────────────────────────────┐ │
45// │ │ KJ Integration │ │
46// │ │ │ │
47// │ │ IoOwn<ReadableStreamSource> │ │
48// │ │ WeakRef for safe references │ │
49// │ │ IoContext-aware operations │ │
50// │ └─────────────────────────────────────┘ │
51// └───────────────────────────────────────────┘
52// │
53// ▼
54// ┌───────────────────────────────────────────┐
55// │ ReadableStreamSource │
56// │ (KJ Native Stream) │
57// │ │
58// │ • tryRead() │
59// │ • pumpTo() │
60// │ • tryGetLength() │
61// │ • cancel() │
62// └───────────────────────────────────────────┘
63//
64// The ReadableSourceKjAdapter adapts a ReadableStream to a KJ-friendly
65// ReadableStreamSource. It holds a strong reference to the ReadableStream and
66// locks it with a ReadableStreamDefaultReader. It is intended to be used by
67// KJ code that wants to read from a JavaScript-backed stream. It ensures that
68// all operations are performed on the correct IoContext, and safely cleans up
69// after itself if the adapter is dropped.
70//
71// ┌───────────────────────────────────────────┐
72// │ ReadableSourceKjAdapter │
73// │ │
74// │ ┌─────────────────────────────────────┐ │
75// │ │ KJ Native API │ │
76// │ │ │ │
77// │ │ • tryRead(minBytes, maxBytes) │ │
78// │ │ • pumpTo(sink, end) │ │
79// │ │ • tryGetLength(encoding) │ │
80// │ │ • cancel(exception) │ │
81// │ │ • getPreferredEncoding() │ │
82// │ │ • tryTee() → none (unsupported) │ │
83// │ └─────────────────────────────────────┘ │
84// │ │ │
85// │ ▼ │
86// │ ┌─────────────────────────────────────┐ │
87// │ │ State Management │ │
88// │ │ │ │
89// │ │ Active ──► Closed │ │
90// │ │ │ │ │
91// │ │ └─────► Canceled/Errored │ │
92// │ └─────────────────────────────────────┘ │
93// │ │ │
94// │ ▼ │
95// │ ┌─────────────────────────────────────┐ │
96// │ │ JavaScript Integration │ │
97// │ │ │ │
98// │ │ ReadableStreamDefaultReader │ │
99// │ │ WeakRef for safe references │ │
100// │ │ IoContext-aware JS operations │ │
101// │ │ Promise handling & async reads │ │
102// │ └─────────────────────────────────────┘ │
103// └───────────────────────────────────────────┘
104// │
105// ▼
106// ┌───────────────────────────────────────────┐
107// │ JavaScript ReadableStream │
108// │ │
109// │ • getReader() │
110// │ • read() → Promise<{value, done}> │
111// │ • cancel(reason) │
112// │ • locked, state properties │
113// └───────────────────────────────────────────┘
114 
115// Adapts a ReadableStreamSource to a JavaScript-friendly interface.
116class ReadableStreamSourceJsAdapter final {
117 public:
118 ReadableStreamSourceJsAdapter(
119 jsg::Lock& js, IoContext& ioContext, kj::Own<ReadableSource> source);
120 KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamSourceJsAdapter);
121 ~ReadableStreamSourceJsAdapter() noexcept(false);
122 
123 // Returns true if the adapter is closed or canceled.
124 bool isClosed();
125 
126 // If the adapter is canceled, returns the exception it was
127 // canceled with. Otherwise returns null.
128 kj::Maybe<const kj::Exception&> isCanceled() KJ_LIFETIMEBOUND;
129 
130 // Cancels the underlying source if it is still active. If an
131 // exception is provided, the source will be errored with that.
132 // If no exception is provided, the source will be closed without
133 // error. All in-flight and pending read requests will be rejected.
134 // Unlike close(), the effect is immediate.
135 void cancel(kj::Exception exception);
136 
137 // Like cancel() but with the error reason provided as a JS value.
138 void cancel(jsg::Lock& js, const jsg::JsValue& reason);
139 
140 // Closes the stream immediatey without error if it is still
141 // active. All in-flight and pending read requests will be
142 // rejected with a cancelation error but the adapter will
143 // transition to the closed state rather than the errored state.
144 // If the adapter is already closed or canceled, this is a no-op.
145 void shutdown(jsg::Lock& js);
146 
147 // Causes the adapter to enter the closing state. Any pending
148 // read requests will be allowed to complete but no new read requests
149 // will be accepted. The underlying source will be closed fully
150 // once all pending reads complete. If cancel() has already been
151 // called, this is a no-op and the returned promise resolves
152 // immediately. If cancel() is called while this is pending,
153 // the returned promise will reject with the same exception
154 // and the cancel will supersede the close.
155 jsg::Promise<void> close(jsg::Lock& js);
156 
157 struct ReadOptions {
158 // The buffer to read into. The maximum number of bytes read
159 // is equal to the length of this buffer. The actual number of
160 // bytes read is indicated by the resolved value of the promise
161 // but will never exceed the length of this buffer.
162 jsg::BufferSource buffer;
163 
164 // The optional minimum number of bytes to read. If not provided,
165 // the read will complete as soon as at least the mininum number
166 // of bytes to satisfy the minimum bytes-per-element of the input
167 // buffer is available.
168 // It is often more efficient to provide a minimum number of bytes
169 // because it allows to implementation to wait until larger chunks
170 // of data are available before completing the read.
171 kj::Maybe<size_t> minBytes;
172 };
173 struct ReadResult {
174 // The buffer containing the data that was read. The length
175 // of the buffer may be less than the length of the buffer
176 // provided in ReadOptions if fewer bytes were available.
177 // The identity of the underlying ArrayBuffer will be the same
178 // but the buffer itself will be a new type array view.
179 // of the same type as that provided in ReadOptions.
180 // If the read produced no data because the stream is
181 // closed, the type array will be zero length.
182 jsg::BufferSource buffer;
183 
184 // True if the stream is now closed and no further reads
185 // are possible. If this is true, the buffer will be zero
186 // length.
187 bool done = false;
188 };
189 
190 // Submit a read request. The returned promise resolves with a
191 // BufferSource containing the data that was read.
192 jsg::Promise<ReadResult> read(jsg::Lock& js, ReadOptions options);
193 
194 // Utility function to read the entire stream as text. This is
195 // terminal in that once this is called, no further reads
196 // are possible. The entire stream will be read and concatenated
197 // and the resulting string returned. If the stream errors while
198 // reading, the promise will reject with the error.
199 // If there are pending reads when this is called, those reads
200 // will be allowed to complete first, and then the stream will
201 // be read to the end.
202 jsg::Promise<jsg::JsRef<jsg::JsString>> readAllText(jsg::Lock& js, uint64_t limit = kj::maxValue);
203 
204 // Utility function to read the entire stream as bytes. This is
205 // terminal in that once this is called, no further reads
206 // are possible. The entire stream will be read and concatenated
207 // and the resulting bytes returned as a single BufferSource.
208 // If the stream errors while reading, the promise will reject
209 // with the error.
210 // If there are pending reads when this is called, those reads
211 // will be allowed to complete first, and then the stream will
212 // be read to the end.
213 jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit = kj::maxValue);
214 
215 // If the stream is still active, tries to get the total length,
216 // if known. If the length is not known, the encoding does not
217 // match the encoding of the underlying stream, or the stream is
218 // closed or errored, returns kj::none.
219 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding);
220 
221 struct Tee {
222 kj::Own<ReadableStreamSourceJsAdapter> branch1;
223 kj::Own<ReadableStreamSourceJsAdapter> branch2;
224 };
225 // Tees the stream into two branches. The returned Tee contains
226 // two new ReadableStreamSourceJsAdapter instances that will
227 // each receive the same data as this instance. Once this is called,
228 // this instance is no longer usable and all further operations
229 // on it will fail. Each branch operates independently; closing,
230 // canceling, or erroring one branch has no effect on the other branch.
231 // If this instance is already closed or canceled, or if there are
232 // in-flight or pending reads, this will throw.
233 kj::Maybe<Tee> tryTee(jsg::Lock& js, uint64_t limit = kj::maxValue);
234 
235 private:
236 struct Active;
237 struct Closed final {
238 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
239 };
240 struct Open {
241 static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj;
242 IoOwn<Active> active;
243 };
244 
245 // State machine for tracking readable source adapter lifecycle:
246 // Open -> Closed (normal close)
247 // Open -> kj::Exception (error via cancel or read failure)
248 // Closed is terminal, kj::Exception is implicitly terminal via ErrorState.
249 using State = StateMachine<TerminalStates<Closed>,
250 ErrorState<kj::Exception>,
251 ActiveState<Open>,
252 Open,
253 Closed,
254 kj::Exception>;
255 State state;
256 
257 kj::Rc<WeakRef<ReadableStreamSourceJsAdapter>> selfRef;
258};
259 
260// ===============================================================================================
261 
262// Adapts a ReadableStream to a KJ-friendly interface.
263// The adapter fully wraps and consumes the ReadableStream instance,
264// using a ReadableStreamDefaultReader to pull data from it.
265// When the adapter is destroyed or canceled, the reader is canceled
266// and both the reader and the stream references are dropped. Critically,
267// the stream is not usable after ownership is transferred to this adapter.
268// Initializing the adapter will fail if the stream is already locked or
269// disturbed.
270//
271// If the adapter is dropped, or canceled while there are pending reads,
272// the pending reads will be rejected with the same exception as the cancel.
273// Because JavaScript promises are not cancelable, reads that are in progress
274// won't be aborted immediately but the results will be ignored when they
275// complete and a best-effort will be made to interrupt the read as soon as
276// possible. If the stream is already closed, reads will complete immediately
277// with 0 bytes read. If the stream errors, reads will reject with the same
278// exception.
279//
280// The minRead contract is enforced. The adapter will attempt to read at
281// least minBytes on each read, under the isolate lock. If the stream ends
282// before minBytes can be satisfied, the read will complete with whatever
283// bytes were available and the adapter will remember that the stream is
284// closed.
285//
286// Concurrent/overlapping reads are not allowed. If a read is already
287// pending, further read attempts will be rejected.
288//
289// While the caller is expected to follow the ReadableStreamSource contract
290// and keep the adapter and buffer alive until the read promises resolve,
291// there are some protections in place to avoid use-after-free if the caller
292// drops the adapter. There's nothing we can do if the caller drops the
293// buffer, however, so that is still a hard requirement.
294// TODO(safety): This can be made safer by having read take a kj::Array
295// as input instead of a raw pointer and size, then having the read return
296// the filled in Array after the read completes, but that's a larger refactor.
297class ReadableSourceKjAdapter final: public ReadableSource {
298 public:
299 enum class MinReadPolicy {
300 // The read will complete as soon as at least minBytes have been read,
301 // even if more bytes are available and the buffer is not full. This
302 // may result in more read calls (keeping in mind that each read needs
303 // to acquire the isolate lock) but may keep the stream flowing more.
304 IMMEDIATE,
305 // The read will attempt to fill the entire buffer until either
306 // maxBytes, the stream ends, or we determine the buffer is "full enough".
307 // This will result in fewer read calls (and thus grabbing the isolate
308 // lock less often) but may result in higher latency for each read.
309 OPPORTUNISTIC,
310 };
311 struct Options {
312 MinReadPolicy minReadPolicy;
313 };
314 
315 ReadableSourceKjAdapter(jsg::Lock& js,
316 IoContext& ioContext,
317 jsg::Ref<ReadableStream> stream,
318 Options options = {.minReadPolicy = MinReadPolicy::OPPORTUNISTIC});
319 ~ReadableSourceKjAdapter() noexcept(false);
320 
321 // Attempts to read at least minBytes and up to maxBytes into the provided
322 // buffer. The returned promise resolves with the actual number of bytes read,
323 // which may be less than minBytes if the stream is fully consumed.
324 //
325 // If the stream is already closed, the returned promise resolves
326 // immediately with 0. If the stream is canceled or errors, the returned
327 // promise rejects with the same exception.
328 //
329 // minBytes must be less than or equal to maxBytes and greater than zero.
330 // If any values outside that range are provided, minBytes will be clamped
331 // to the range [1, maxBytes].
332 //
333 // Per the contact of tryRead, it is the caller's responsibility to ensure
334 // that both the buffer and this adapter remain alive until the returned
335 // promise resolves! It is also the caller's responsibility to ensure that
336 // buffer is at least maxBytes in length. However, there are some protections
337 // implemented to avoid use-after-free if the adapter is dropped while a read
338 // is in progress.
339 //
340 // The returned promise will never resolve with more than maxBytes.
341 kj::Promise<size_t> read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes) override;
342 
343 // Reads all remaining bytes from the stream and returns them.
344 kj::Promise<kj::Array<const kj::byte>> readAllBytes(size_t limit) override;
345 
346 // Reads all remaining bytes from the stream and returns them as a string.
347 kj::Promise<kj::String> readAllText(size_t limit) override;
348 
349 // Fully consume the stream and write it to the provided WritableStreamSink.
350 // If "end" is true, the output stream will be ended once the input
351 // stream is fully consumed.
352 // Per the contract of pumpTo, it is the caller's responsibility to ensure
353 // that both the WritableStreamSink and this adapter remain alive until
354 // the returned promise resolves!
355 kj::Promise<DeferredProxy<void>> pumpTo(WritableSink& output, EndAfterPump end) override;
356 
357 // If the stream is still active, tries to get the total length,
358 // if known. If the length is not known, the encoding does not
359 // match the encoding of the underlying stream, or the stream is closed
360 // or errored, returns kj::none.
361 kj::Maybe<size_t> tryGetLength(StreamEncoding encoding) override;
362 
363 // Cancels the underlying source if it is still active.
364 void cancel(kj::Exception reason) override;
365 
366 StreamEncoding getEncoding() override {
367 // Our underlying ReadableStream produces non-encoded bytes (for now)
368 return StreamEncoding::IDENTITY;
369 };
370 
371 Tee tee(size_t limit) override;
372 
373 struct ReadContext;
374 KJ_DECLARE_NON_POLYMORPHIC(ReadContext);
375 
376 private:
377 struct Active;
378 KJ_DECLARE_NON_POLYMORPHIC(Active);
379 struct KjClosed {
380 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
381 };
382 struct KjOpen {
383 static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj;
384 kj::Own<Active> active;
385 };
386 
387 // State machine for tracking readable source adapter lifecycle:
388 // KjOpen -> KjClosed (normal close)
389 // KjOpen -> kj::Exception (error via cancel or read failure)
390 // KjClosed is terminal, kj::Exception is implicitly terminal via ErrorState.
391 using KjState = StateMachine<TerminalStates<KjClosed>,
392 ErrorState<kj::Exception>,
393 ActiveState<KjOpen>,
394 KjOpen,
395 KjClosed,
396 kj::Exception>;
397 KjState state;
398 const Options options;
399 kj::Rc<WeakRef<ReadableSourceKjAdapter>> selfRef;
400 
401 // Checks if the inner Active state is Canceling or Canceled.
402 // If so, transitions to error state and throws the exception.
403 void throwIfCancelingOrCanceled(Active& active);
404 
405 // Checks if the inner Active state is Canceling or Canceled.
406 // If so, transitions to error state and returns the exception.
407 kj::Maybe<kj::Exception> checkCancelingOrCanceled(Active& active);
408 
409 kj::Promise<size_t> readImpl(Active& active, kj::ArrayPtr<kj::byte> buffer, size_t minBytes);
410 
411 static kj::Promise<void> pumpToImpl(
412 kj::Own<Active> active, WritableSink& output, EndAfterPump end);
413 static jsg::Promise<kj::Own<ReadContext>> readInternal(
414 jsg::Lock& js, kj::Own<ReadContext> context, MinReadPolicy minReadPolicy);
415 
416 template <typename T>
417 kj::Promise<kj::Array<T>> readAllImpl(size_t limit);
418 
419 struct CancelationToken final {
420 kj::Rc<WeakRef<CancelationToken>> selfRef;
421 inline CancelationToken()
422 : selfRef(kj::rc<WeakRef<CancelationToken>>(kj::Badge<CancelationToken>{}, *this)) {}
423 inline ~CancelationToken() {
424 selfRef->invalidate();
425 }
426 inline kj::Rc<WeakRef<CancelationToken>> getWeakRef() {
427 return selfRef.addRef();
428 }
429 };
430 
431 template <typename T>
432 static jsg::Promise<kj::Array<T>> readAllReadImpl(jsg::Lock& js,
433 IoOwn<Active> active,
434 kj::Vector<T> accumulated,
435 size_t limit,
436 kj::Rc<WeakRef<CancelationToken>> cancelationToken);
437};
438 
439} // namespace workerd::api::streams