Skip to content
File

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

cpp560 lines
1// Copyright (c) 2017-2022 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 "common.h"
8 
9#include <kj/function.h>
10#include <workerd/util/state-machine.h>
11 
12namespace workerd::api {
13 
14class ReadableStreamDefaultReader;
15class ReadableStreamBYOBReader;
16 
17class ReaderImpl final {
18public:
19 ReaderImpl(ReadableStreamController::Reader& reader);
20 
21 ~ReaderImpl() noexcept(false);
22 
23 void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise);
24 
25 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason);
26 
27 void detach();
28 
29 jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed();
30 
31 void lockToStream(jsg::Lock& js, ReadableStream& stream);
32 
33 jsg::Promise<ReadResult> read(jsg::Lock& js,
34 kj::Maybe<ReadableStreamController::ByobOptions> byobOptions);
35 
36 void releaseLock(jsg::Lock& js);
37 
38 void visitForGc(jsg::GcVisitor& visitor);
39 
40 kj::StringPtr jsgGetMemoryName() const;
41 size_t jsgGetMemorySelfSize() const;
42 void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
43 
44private:
45 struct Initial {
46 static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj;
47 };
48 // While a Reader is attached to a ReadableStream, it holds a strong reference to the
49 // ReadableStream to prevent it from being GC'ed so long as the Reader is available.
50 // Once the reader is closed, released, or GC'ed the reference to the ReadableStream
51 // is cleared and the ReadableStream can be GC'ed if there are no other references to
52 // it being held anywhere. If the reader is still attached to the ReadableStream when
53 // it is destroyed, the ReadableStream's reference to the reader is cleared but the
54 // ReadableStream remains in the "reader locked" state, per the spec.
55 struct Attached {
56 static constexpr kj::StringPtr NAME KJ_UNUSED = "attached"_kj;
57 jsg::Ref<ReadableStream> stream;
58 };
59 // Released: The user explicitly called releaseLock() to detach the reader from the stream.
60 // The stream remains usable and can be locked by a new reader.
61 struct Released {
62 static constexpr kj::StringPtr NAME KJ_UNUSED = "released"_kj;
63 };
64 // Closed: The underlying stream ended (closed or errored) while the reader was attached.
65 // The stream is no longer usable.
66 struct Closed {
67 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
68 };
69 
70 // State machine for ReaderImpl:
71 // Initial -> Attached (attach() called)
72 // Attached -> Closed (detach() called when stream closes)
73 // Attached -> Released (releaseLock() called)
74 // Closed and Released are terminal states.
75 // Initial is not terminal but most methods assert if called in this state.
76 using ReaderState = StateMachine<TerminalStates<Closed, Released>,
77 ActiveState<Attached>,
78 Initial,
79 Attached,
80 Closed,
81 Released>;
82 
83 kj::Maybe<IoContext&> ioContext;
84 ReadableStreamController::Reader& reader;
85 
86 ReaderState state;
87 
88 inline void assertAttachedOrTerminal() const {
89 KJ_ASSERT(!state.is<Initial>(), "this reader was never attached");
90 }
91 kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<void>>> closedPromise;
92 
93 friend class ReadableStreamDefaultReader;
94 friend class ReadableStreamBYOBReader;
95};
96 
97class ReadableStreamDefaultReader : public jsg::Object,
98 public ReadableStreamController::Reader {
99public:
100 explicit ReadableStreamDefaultReader();
101 
102 // JavaScript API
103 
104 static jsg::Ref<ReadableStreamDefaultReader> constructor(
105 jsg::Lock& js, jsg::Ref<ReadableStream> stream);
106 
107 jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed();
108 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason);
109 jsg::Promise<ReadResult> read(jsg::Lock& js);
110 void releaseLock(jsg::Lock& js);
111 
112 JSG_RESOURCE_TYPE(ReadableStreamDefaultReader, CompatibilityFlags::Reader flags) {
113 if (flags.getJsgPropertyOnPrototypeTemplate()) {
114 JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed);
115 } else {
116 JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed);
117 }
118 JSG_METHOD(cancel);
119 JSG_METHOD(read);
120 JSG_METHOD(releaseLock);
121 
122 JSG_TS_OVERRIDE(<R = any> {
123 read(): Promise<ReadableStreamReadResult<R>>;
124 });
125 }
126 
127 // Internal API
128 
129 void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) override;
130 
131 void detach() override;
132 
133 void lockToStream(jsg::Lock& js, ReadableStream& stream);
134 
135 inline bool isByteOriented() const override { return false; }
136 
137 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
138 tracker.trackField("impl", impl);
139 }
140 
141private:
142 ReaderImpl impl;
143 
144 void visitForGc(jsg::GcVisitor& visitor);
145};
146 
147class ReadableStreamBYOBReader: public jsg::Object,
148 public ReadableStreamController::Reader {
149public:
150 explicit ReadableStreamBYOBReader();
151 
152 // JavaScript API
153 
154 static jsg::Ref<ReadableStreamBYOBReader> constructor(
155 jsg::Lock& js,
156 jsg::Ref<ReadableStream> stream);
157 
158 jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed();
159 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason);
160 
161 struct ReadableStreamBYOBReaderReadOptions {
162 jsg::Optional<int> min;
163 JSG_STRUCT(min);
164 };
165 
166 jsg::Promise<ReadResult> read(jsg::Lock& js, v8::Local<v8::ArrayBufferView> byobBuffer,
167 jsg::Optional<ReadableStreamBYOBReaderReadOptions> options = kj::none);
168 
169 // Non-standard extension so that reads can specify a minimum number of elements to read. It's a
170 // struct so that we could eventually add things like timeouts if we need to. Since there's no
171 // existing spec that's a leading contender, this is behind a different method name to avoid
172 // conflicts with any changes to `read`. Fewer than `minElements` may be returned if EOF is hit
173 // or the underlying stream is closed/errors out. In all cases the read result is either
174 // {value: theChunk, done: false} or {value: undefined, done: true} as with read.
175 // TODO(soon): Like fetch() and Cache.match(), readAtLeast() returns a promise for a V8 object.
176 jsg::Promise<ReadResult> readAtLeast(jsg::Lock& js,
177 int minElements,
178 v8::Local<v8::ArrayBufferView> byobBuffer);
179 
180 void releaseLock(jsg::Lock& js);
181 
182 JSG_RESOURCE_TYPE(ReadableStreamBYOBReader, CompatibilityFlags::Reader flags) {
183 if (flags.getJsgPropertyOnPrototypeTemplate()) {
184 JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed);
185 } else {
186 JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed);
187 }
188 JSG_METHOD(cancel);
189 JSG_METHOD(read);
190 JSG_METHOD(releaseLock);
191 
192 // Non-standard extension that should only apply to BYOB byte streams.
193 JSG_METHOD(readAtLeast);
194 
195 JSG_TS_OVERRIDE(ReadableStreamBYOBReader {
196 read<T extends ArrayBufferView>(view: T): Promise<ReadableStreamReadResult<T>>;
197 readAtLeast<T extends ArrayBufferView>(minElements: number, view: T): Promise<ReadableStreamReadResult<T>>;
198 });
199 }
200 
201 // Internal API
202 
203 void attach(
204 ReadableStreamController& controller,
205 jsg::Promise<void> closedPromise) override;
206 
207 void detach() override;
208 
209 void lockToStream(jsg::Lock& js, ReadableStream& stream);
210 
211 inline bool isByteOriented() const override { return true; }
212 
213 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
214 tracker.trackField("impl", impl);
215 }
216 
217private:
218 ReaderImpl impl;
219 
220 void visitForGc(jsg::GcVisitor& visitor);
221};
222 
223// DrainingReader is a C++ only reader (not exposed to JavaScript) that performs
224// draining reads. It locks the stream like standard readers but uses drainingRead()
225// instead of regular read() to drain all synchronously available data at once.
226// This is intended for optimized pipe operations.
227class DrainingReader: public ReadableStreamController::Reader {
228 public:
229 explicit DrainingReader();
230 
231 // Factory method to create and lock to a stream. Returns nullptr if stream is locked.
232 static kj::Maybe<kj::Own<DrainingReader>> create(jsg::Lock& js, ReadableStream& stream);
233 
234 virtual ~DrainingReader() noexcept(false);
235 
236 // Performs a draining read, returning all synchronously available data as bytes.
237 // The maxRead parameter is a soft limit - see ReadableStreamController::drainingRead.
238 jsg::Promise<DrainingReadResult> read(jsg::Lock& js, size_t maxRead = kj::maxValue);
239 
240 // Cancels the stream.
241 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason);
242 
243 // Releases the lock on the stream.
244 void releaseLock(jsg::Lock& js);
245 
246 // Returns whether this reader is still attached to a stream.
247 bool isAttached() const;
248 
249 // ReadableStreamController::Reader interface
250 void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) override;
251 void detach() override;
252 bool isByteOriented() const override { return false; }
253 
254 void visitForGc(jsg::GcVisitor& visitor);
255 
256 private:
257 struct Initial {};
258 using Attached = jsg::Ref<ReadableStream>;
259 struct Released {};
260 
261 kj::Maybe<IoContext&> ioContext;
262 kj::OneOf<Initial, Attached, StreamStates::Closed, Released> state = Initial();
263 kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<void>>> closedPromise;
264};
265 
266class ReadableStream: public jsg::Object {
267private:
268 
269 struct AsyncIteratorState {
270 kj::Maybe<IoContext&> ioContext;
271 jsg::Ref<ReadableStreamDefaultReader> reader;
272 bool preventCancel;
273 };
274 
275 static jsg::Promise<kj::Maybe<jsg::Value>> nextFunction(
276 jsg::Lock& js,
277 AsyncIteratorState& state);
278 
279 static jsg::Promise<void> returnFunction(
280 jsg::Lock& js,
281 AsyncIteratorState& state,
282 jsg::Optional<jsg::Value>& value);
283 
284public:
285 explicit ReadableStream(IoContext& ioContext,
286 kj::Own<ReadableStreamSource> source);
287 
288 explicit ReadableStream(kj::Own<ReadableStreamController> controller);
289 
290 ReadableStreamController& getController();
291 
292 jsg::Ref<ReadableStream> addRef();
293 
294 bool isDisturbed();
295 
296 // ---------------------------------------------------------------------------
297 // JS interface
298 
299 // Creates a new JS-backed ReadableStream using the provided source and strategy.
300 // We use v8::Local<v8::Object>'s here instead of jsg structs because we need
301 // to preserve the object references within the implementation.
302 static jsg::Ref<ReadableStream> constructor(
303 jsg::Lock& js,
304 jsg::Optional<UnderlyingSource> underlyingSource,
305 jsg::Optional<StreamQueuingStrategy> queuingStrategy);
306 
307 static jsg::Ref<ReadableStream> from(jsg::Lock& js, jsg::AsyncGenerator<jsg::Value> generator);
308 
309 bool isLocked();
310 
311 // Closes the stream. All present and future read requests are fulfilled with successful empty
312 // results. `reason` will be passed to the underlying source's cancel algorithm -- if this
313 // readable stream is one side of a transform stream, then its cancel algorithm causes the
314 // transform's writable side to become errored with `reason`.
315 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason);
316 
317 using Reader = kj::OneOf<jsg::Ref<ReadableStreamDefaultReader>,
318 jsg::Ref<ReadableStreamBYOBReader>>;
319 
320 struct GetReaderOptions {
321 jsg::Optional<kj::String> mode; // can be "byob" or undefined
322 
323 JSG_STRUCT(mode);
324 
325 JSG_STRUCT_TS_OVERRIDE({ mode: "byob" });
326 // Intentionally required, so we can use `GetReaderOptions` directly in the
327 // `ReadableStream#getReader()` overload.
328 };
329 
330 Reader getReader(jsg::Lock& js, jsg::Optional<GetReaderOptions> options);
331 
332 // Options specifically for the values() function.
333 struct ValuesOptions {
334 jsg::Optional<bool> preventCancel = false;
335 JSG_STRUCT(preventCancel);
336 };
337 
338 JSG_ASYNC_ITERATOR_WITH_OPTIONS(ReadableStreamAsyncIterator,
339 values,
340 jsg::Value,
341 AsyncIteratorState,
342 nextFunction,
343 returnFunction,
344 ValuesOptions);
345 struct Transform {
346 jsg::Ref<ReadableStream> readable;
347 jsg::Ref<WritableStream> writable;
348 
349 JSG_STRUCT(readable, writable);
350 JSG_STRUCT_TS_OVERRIDE(ReadableWritablePair<R = any, W = any> {
351 readable: ReadableStream<R>;
352 writable: WritableStream<W>;
353 });
354 };
355 
356 jsg::Ref<ReadableStream> pipeThrough(
357 jsg::Lock& js,
358 Transform transform,
359 jsg::Optional<PipeToOptions> options);
360 
361 jsg::Promise<void> pipeTo(
362 jsg::Lock& js,
363 jsg::Ref<WritableStream> destination,
364 jsg::Optional<PipeToOptions> options);
365 
366 // Locks the stream and returns a pair of two new ReadableStreams, each of which read the same
367 // data as this ReadableStream would.
368 kj::Array<jsg::Ref<ReadableStream>> tee(jsg::Lock& js);
369 
370 jsg::JsString inspectState(jsg::Lock& js);
371 bool inspectSupportsBYOB();
372 jsg::Optional<uint64_t> inspectLength();
373 
374 JSG_RESOURCE_TYPE(ReadableStream, CompatibilityFlags::Reader flags) {
375 if (flags.getJsgPropertyOnPrototypeTemplate()) {
376 JSG_READONLY_PROTOTYPE_PROPERTY(locked, isLocked);
377 } else {
378 JSG_READONLY_INSTANCE_PROPERTY(locked, isLocked);
379 }
380 JSG_METHOD(cancel);
381 JSG_METHOD(getReader);
382 JSG_METHOD(pipeThrough);
383 JSG_METHOD(pipeTo);
384 JSG_METHOD(tee);
385 JSG_METHOD(values);
386 JSG_STATIC_METHOD(from);
387 
388 JSG_INSPECT_PROPERTY(state, inspectState);
389 JSG_INSPECT_PROPERTY(supportsBYOB, inspectSupportsBYOB);
390 JSG_INSPECT_PROPERTY(length, inspectLength);
391 
392 JSG_ASYNC_ITERABLE(values);
393 
394 if (flags.getJsgPropertyOnPrototypeTemplate()) {
395 JSG_TS_DEFINE(interface ReadableStream<R = any> {
396 get locked(): boolean;
397 
398 cancel(reason?: any): Promise<void>;
399 
400 getReader(): ReadableStreamDefaultReader<R>;
401 getReader(options: ReadableStreamGetReaderOptions): ReadableStreamBYOBReader;
402 
403 pipeThrough<T>(transform: ReadableWritablePair<T, R>, options?: StreamPipeOptions): ReadableStream<T>;
404 pipeTo(destination: WritableStream<R>, options?: StreamPipeOptions): Promise<void>;
405 
406 tee(): [ReadableStream<R>, ReadableStream<R>];
407 
408 values(options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>;
409 [Symbol.asyncIterator](options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>;
410 });
411 } else {
412 JSG_TS_DEFINE(interface ReadableStream<R = any> {
413 readonly locked: boolean;
414 
415 cancel(reason?: any): Promise<void>;
416 
417 getReader(): ReadableStreamDefaultReader<R>;
418 getReader(options: ReadableStreamGetReaderOptions): ReadableStreamBYOBReader;
419 
420 pipeThrough<T>(transform: ReadableWritablePair<T, R>, options?: StreamPipeOptions): ReadableStream<T>;
421 pipeTo(destination: WritableStream<R>, options?: StreamPipeOptions): Promise<void>;
422 
423 tee(): [ReadableStream<R>, ReadableStream<R>];
424 
425 values(options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>;
426 [Symbol.asyncIterator](options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>;
427 });
428 }
429 // Replace ReadableStream class with an interface and const, so we can have
430 // two constructors with differing type parameters for byte-oriented and
431 // value-oriented streams.
432 JSG_TS_OVERRIDE(const ReadableStream: {
433 prototype: ReadableStream;
434 new (underlyingSource: UnderlyingByteSource, strategy?: QueuingStrategy<Uint8Array>): ReadableStream<Uint8Array>;
435 new <R = any>(underlyingSource?: UnderlyingSource<R>, strategy?: QueuingStrategy<R>): ReadableStream<R>;
436 });
437 }
438 
439 // Detaches this ReadableStream from its underlying controller state, returning a
440 // new ReadableStream instance that takes over the underlying state. This is used to
441 // support the "create a proxy" of a ReadableStream algorithm in the streams spec
442 // (see https://streams.spec.whatwg.org/#readablestream-create-a-proxy). In that
443 // algorithm, it says to create a proxy of a stream by creating a new TransformStream
444 // and piping the original through it. The readable side of the created transform
445 // becomes the proxy. That is quite inefficient so instead, we create a new
446 // ReadableStream that will take over ownership of the internal state of this one,
447 // leaving this ReadableStream locked and disturbed so that it is no longer usable.
448 // The name "detach" here is used in the sense of "detaching the internal state".
449 jsg::Ref<ReadableStream> detach(jsg::Lock& js, bool ignoreDisturbed=false);
450 
451 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding);
452 
453 // A potentially optimized version of pipe that sends this stream's data to the given
454 // sink. The entire stream is consumed. The ReadableStream will be left locked and
455 // disturbed and the DeferredProxy returned will take over ownership of the internal
456 // state of the readable.
457 kj::Promise<DeferredProxy<void>> pumpTo(jsg::Lock& js,
458 kj::Own<WritableStreamSink> sink,
459 bool end);
460 
461 // Initializes signalling mechanism for EOF detection. Returns a promise that will resolve when
462 // EOF is reached.
463 //
464 // This method should only be called once.
465 jsg::Promise<void> onEof(jsg::Lock& js);
466 
467 // Used by ReadableStreamInternalController to signal EOF being reached. Can be called even if
468 // `onEof` wasn't called.
469 void signalEof(jsg::Lock& js);
470 
471 void serialize(jsg::Lock& js, jsg::Serializer& serializer);
472 static jsg::Ref<ReadableStream> deserialize(
473 jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer);
474 
475 JSG_SERIALIZABLE(rpc::SerializationTag::READABLE_STREAM);
476 
477 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
478 
479private:
480 kj::Maybe<IoContext&> ioContext;
481 kj::Own<ReadableStreamController> controller;
482 
483 // Used to signal when this ReadableStream reads EOF. This signal is required for TCP sockets.
484 kj::Maybe<jsg::PromiseResolverPair<void>> eofResolverPair;
485 
486 void visitForGc(jsg::GcVisitor& visitor);
487};
488 
489struct QueuingStrategyInit {
490 double highWaterMark;
491 JSG_STRUCT(highWaterMark);
492};
493 
494using QueuingStrategySizeFunction =
495 jsg::Optional<uint32_t>(jsg::Optional<v8::Local<v8::Value>>);
496 
497// Utility class defined by the streams spec that uses byteLength to calculate
498// backpressure changes.
499class ByteLengthQueuingStrategy: public jsg::Object {
500public:
501 ByteLengthQueuingStrategy(QueuingStrategyInit init) : init(init) {}
502 
503 static jsg::Ref<ByteLengthQueuingStrategy> constructor(jsg::Lock& js, QueuingStrategyInit init) {
504 return js.alloc<ByteLengthQueuingStrategy>(init);
505 }
506 
507 double getHighWaterMark() const { return init.highWaterMark; }
508 
509 jsg::Function<QueuingStrategySizeFunction> getSize() const { return &size; }
510 
511 JSG_RESOURCE_TYPE(ByteLengthQueuingStrategy) {
512 JSG_READONLY_PROTOTYPE_PROPERTY(highWaterMark, getHighWaterMark);
513 JSG_READONLY_PROTOTYPE_PROPERTY(size, getSize);
514 
515 // QueuingStrategy requires the result of the size function to be defined
516 JSG_TS_OVERRIDE(implements QueuingStrategy<ArrayBufferView> {
517 get size(): (chunk?: any) => number;
518 });
519 }
520 
521private:
522 static jsg::Optional<uint32_t> size(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>>);
523 
524 QueuingStrategyInit init;
525};
526 
527// Utility class defined by the streams spec that uses a fixed value of 1 to calculate
528// backpressure change
529class CountQueuingStrategy: public jsg::Object {
530public:
531 CountQueuingStrategy(QueuingStrategyInit init) : init(init) {}
532 
533 static jsg::Ref<CountQueuingStrategy> constructor(jsg::Lock& js, QueuingStrategyInit init) {
534 return js.alloc<CountQueuingStrategy>(init);
535 }
536 
537 double getHighWaterMark() const { return init.highWaterMark; }
538 
539 jsg::Function<QueuingStrategySizeFunction> getSize() const { return &size; }
540 
541 JSG_RESOURCE_TYPE(CountQueuingStrategy) {
542 JSG_READONLY_PROTOTYPE_PROPERTY(highWaterMark, getHighWaterMark);
543 JSG_READONLY_PROTOTYPE_PROPERTY(size, getSize);
544 
545 // QueuingStrategy requires the result of the size function to be defined
546 JSG_TS_OVERRIDE(implements QueuingStrategy {
547 get size(): (chunk?: any) => number;
548 });
549 }
550 
551private:
552 static jsg::Optional<uint32_t> size(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>>) {
553 return 1;
554 }
555 
556 QueuingStrategyInit init;
557};
558 
559} // namespace workerd::api