Skip to content
File

Blob: src/workerd/api/streams/readable.c++

31.8 KB
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#include "readable.h"
6 
7#include "internal.h"
8#include "writable.h"
9 
10#include <workerd/api/system-streams.h>
11#include <workerd/api/worker-rpc.h>
12#include <workerd/io/features.h>
13#include <workerd/jsg/jsg.h>
14 
15namespace workerd::api {
16 
17ReaderImpl::ReaderImpl(ReadableStreamController::Reader& reader)
18 : ioContext(tryGetIoContext()),
19 reader(reader),
20 state(ReaderState::create<Initial>()) {}
21 
22ReaderImpl::~ReaderImpl() noexcept(false) {
23 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
24 attached.stream->getController().releaseReader(reader, kj::none);
25 }
26}
27 
28void ReaderImpl::attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) {
29 KJ_ASSERT(state.is<Initial>());
30 state.transitionTo<Attached>(controller.addRef());
31 this->closedPromise = kj::mv(closedPromise);
32}
33 
34void ReaderImpl::detach() {
35 // Only transition from Attached to Closed.
36 // All other states (Initial, Closed, Released) are no-ops.
37 if (state.isActive()) {
38 state.transitionTo<Closed>();
39 }
40}
41 
42jsg::Promise<void> ReaderImpl::cancel(
43 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
44 assertAttachedOrTerminal();
45 if (state.is<Released>()) {
46 return js.rejectedPromise<void>(
47 js.v8TypeError("This ReadableStream reader has been released."_kj));
48 }
49 if (state.is<Closed>()) {
50 return js.resolvedPromise();
51 }
52 auto& attached = state.requireActiveUnsafe();
53 // In some edge cases, this reader is the last thing holding a strong
54 // reference to the stream. Calling cancel might cause the readers strong
55 // reference to be cleared, so let's make sure we keep a reference to
56 // the stream at least until the call to cancel completes.
57 auto ref = attached.stream.addRef();
58 return attached.stream->getController().cancel(js, maybeReason);
59}
60 
61jsg::MemoizedIdentity<jsg::Promise<void>>& ReaderImpl::getClosed() {
62 // The closed promise should always be set after the object is created so this assert
63 // should always be safe.
64 return KJ_ASSERT_NONNULL(closedPromise);
65}
66 
67void ReaderImpl::lockToStream(jsg::Lock& js, ReadableStream& stream) {
68 KJ_ASSERT(!stream.isLocked());
69 KJ_ASSERT(stream.getController().lockReader(js, reader));
70}
71 
72jsg::Promise<ReadResult> ReaderImpl::read(
73 jsg::Lock& js, kj::Maybe<ReadableStreamController::ByobOptions> byobOptions) {
74 assertAttachedOrTerminal();
75 if (state.is<Released>()) {
76 return js.rejectedPromise<ReadResult>(
77 js.v8TypeError("This ReadableStream reader has been released."_kj));
78 }
79 if (state.is<Closed>()) {
80 return js.rejectedPromise<ReadResult>(
81 js.v8TypeError("This ReadableStream has been closed."_kj));
82 }
83 auto& attached = state.requireActiveUnsafe();
84 KJ_IF_SOME(options, byobOptions) {
85 // Per the spec, we must perform these checks before disturbing the stream.
86 size_t atLeast = options.atLeast.orDefault(1);
87 
88 if (options.byteLength == 0) {
89 return js.rejectedPromise<ReadResult>(
90 js.v8TypeError("You must call read() on a \"byob\" reader with a positive-sized "
91 "TypedArray object."_kj));
92 }
93 if (atLeast == 0) {
94 return js.rejectedPromise<ReadResult>(js.v8TypeError(
95 kj::str("Requested invalid minimum number of bytes to read (", atLeast, ").")));
96 }
97 
98 // Both read() and readAtLeast() pass atLeast in element count.
99 // Convert to bytes before validation and forwarding to the controller.
100 jsg::BufferSource source(js, options.bufferView.getHandle(js));
101 auto elementSize = source.getElementSize();
102 atLeast = atLeast * elementSize;
103 
104 if (atLeast > options.byteLength) {
105 return js.rejectedPromise<ReadResult>(js.v8TypeError(kj::str("Minimum bytes to read (",
106 atLeast, ") exceeds size of buffer (", options.byteLength, ").")));
107 }
108 
109 options.atLeast = atLeast;
110 }
111 
112 return KJ_ASSERT_NONNULL(attached.stream->getController().read(js, kj::mv(byobOptions)));
113}
114 
115void ReaderImpl::releaseLock(jsg::Lock& js) {
116 // TODO(soon): Releasing the lock should cancel any pending reads. This is a recent
117 // modification to the spec that we have not yet implemented.
118 assertAttachedOrTerminal();
119 // Closed and Released states are no-ops.
120 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
121 // In some edge cases, this reader is the last thing holding a strong
122 // reference to the stream. Calling releaseLock might cause the readers strong
123 // reference to be cleared, so let's make sure we keep a reference to
124 // the stream at least until the call to releaseLock completes.
125 auto ref = attached.stream.addRef();
126 attached.stream->getController().releaseReader(reader, js);
127 state.transitionTo<Released>();
128 }
129}
130 
131void ReaderImpl::visitForGc(jsg::GcVisitor& visitor) {
132 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
133 visitor.visit(attached.stream);
134 }
135 visitor.visit(closedPromise);
136}
137 
138// ======================================================================================
139 
140ReadableStreamDefaultReader::ReadableStreamDefaultReader(): impl(*this) {}
141 
142jsg::Ref<ReadableStreamDefaultReader> ReadableStreamDefaultReader::constructor(
143 jsg::Lock& js, jsg::Ref<ReadableStream> stream) {
144 JSG_REQUIRE(
145 !stream->isLocked(), TypeError, "This ReadableStream is currently locked to a reader.");
146 auto reader = js.alloc<ReadableStreamDefaultReader>();
147 reader->lockToStream(js, *stream);
148 return kj::mv(reader);
149}
150 
151void ReadableStreamDefaultReader::attach(
152 ReadableStreamController& controller, jsg::Promise<void> closedPromise) {
153 impl.attach(controller, kj::mv(closedPromise));
154}
155 
156jsg::Promise<void> ReadableStreamDefaultReader::cancel(
157 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
158 return impl.cancel(js, kj::mv(maybeReason));
159}
160 
161void ReadableStreamDefaultReader::detach() {
162 impl.detach();
163}
164 
165jsg::MemoizedIdentity<jsg::Promise<void>>& ReadableStreamDefaultReader::getClosed() {
166 return impl.getClosed();
167}
168 
169void ReadableStreamDefaultReader::lockToStream(jsg::Lock& js, ReadableStream& stream) {
170 impl.lockToStream(js, stream);
171}
172 
173jsg::Promise<ReadResult> ReadableStreamDefaultReader::read(jsg::Lock& js) {
174 return impl.read(js, kj::none);
175}
176 
177void ReadableStreamDefaultReader::releaseLock(jsg::Lock& js) {
178 impl.releaseLock(js);
179}
180 
181void ReadableStreamDefaultReader::visitForGc(jsg::GcVisitor& visitor) {
182 visitor.visit(impl);
183}
184 
185// ======================================================================================
186 
187ReadableStreamBYOBReader::ReadableStreamBYOBReader(): impl(*this) {}
188 
189jsg::Ref<ReadableStreamBYOBReader> ReadableStreamBYOBReader::constructor(
190 jsg::Lock& js, jsg::Ref<ReadableStream> stream) {
191 JSG_REQUIRE(
192 !stream->isLocked(), TypeError, "This ReadableStream is currently locked to a reader.");
193 
194 if (!stream->getController().isClosedOrErrored()) {
195 JSG_REQUIRE(stream->getController().isByteOriented(), TypeError,
196 "This ReadableStream does not support BYOB reads.");
197 }
198 
199 auto reader = js.alloc<ReadableStreamBYOBReader>();
200 reader->lockToStream(js, *stream);
201 return kj::mv(reader);
202}
203 
204void ReadableStreamBYOBReader::attach(
205 ReadableStreamController& controller, jsg::Promise<void> closedPromise) {
206 impl.attach(controller, kj::mv(closedPromise));
207}
208 
209jsg::Promise<void> ReadableStreamBYOBReader::cancel(
210 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
211 return impl.cancel(js, kj::mv(maybeReason));
212}
213 
214void ReadableStreamBYOBReader::detach() {
215 impl.detach();
216}
217 
218jsg::MemoizedIdentity<jsg::Promise<void>>& ReadableStreamBYOBReader::getClosed() {
219 return impl.getClosed();
220}
221 
222void ReadableStreamBYOBReader::lockToStream(jsg::Lock& js, ReadableStream& stream) {
223 impl.lockToStream(js, stream);
224}
225 
226jsg::Promise<ReadResult> ReadableStreamBYOBReader::read(jsg::Lock& js,
227 v8::Local<v8::ArrayBufferView> byobBuffer,
228 jsg::Optional<ReadableStreamBYOBReaderReadOptions> maybeOptions) {
229 static const ReadableStreamBYOBReaderReadOptions defaultOptions{};
230 auto options = ReadableStreamController::ByobOptions{
231 .bufferView = js.v8Ref(byobBuffer),
232 .byteOffset = byobBuffer->ByteOffset(),
233 .byteLength = byobBuffer->ByteLength(),
234 .atLeast = maybeOptions.orDefault(defaultOptions).min.orDefault(1),
235 .detachBuffer = FeatureFlags::get(js).getStreamsByobReaderDetachesBuffer(),
236 };
237 return impl.read(js, kj::mv(options));
238}
239 
240jsg::Promise<ReadResult> ReadableStreamBYOBReader::readAtLeast(
241 jsg::Lock& js, int minElements, v8::Local<v8::ArrayBufferView> byobBuffer) {
242 auto options = ReadableStreamController::ByobOptions{
243 .bufferView = js.v8Ref(byobBuffer),
244 .byteOffset = byobBuffer->ByteOffset(),
245 .byteLength = byobBuffer->ByteLength(),
246 .atLeast = minElements,
247 .detachBuffer = true,
248 };
249 return impl.read(js, kj::mv(options));
250}
251 
252void ReadableStreamBYOBReader::releaseLock(jsg::Lock& js) {
253 impl.releaseLock(js);
254}
255 
256void ReadableStreamBYOBReader::visitForGc(jsg::GcVisitor& visitor) {
257 visitor.visit(impl);
258}
259 
260// ======================================================================================
261// DrainingReader implementation
262 
263DrainingReader::DrainingReader(): ioContext(tryGetIoContext()) {}
264 
265DrainingReader::~DrainingReader() noexcept(false) {
266 KJ_IF_SOME(stream, state.tryGet<Attached>()) {
267 stream->getController().releaseReader(*this, kj::none);
268 }
269}
270 
271kj::Maybe<kj::Own<DrainingReader>> DrainingReader::create(jsg::Lock& js, ReadableStream& stream) {
272 if (stream.isLocked()) {
273 return kj::none;
274 }
275 auto reader = kj::heap<DrainingReader>();
276 if (!stream.getController().lockReader(js, *reader)) {
277 return kj::none;
278 }
279 return kj::mv(reader);
280}
281 
282void DrainingReader::attach(
283 ReadableStreamController& controller, jsg::Promise<void> closedPromise) {
284 KJ_ASSERT(state.is<Initial>());
285 state = controller.addRef();
286 this->closedPromise = kj::mv(closedPromise);
287}
288 
289void DrainingReader::detach() {
290 KJ_SWITCH_ONEOF(state) {
291 KJ_CASE_ONEOF(i, Initial) {
292 return;
293 }
294 KJ_CASE_ONEOF(stream, Attached) {
295 state.init<StreamStates::Closed>();
296 return;
297 }
298 KJ_CASE_ONEOF(c, StreamStates::Closed) {
299 return;
300 }
301 KJ_CASE_ONEOF(r, Released) {
302 return;
303 }
304 }
305 KJ_UNREACHABLE;
306}
307 
308jsg::Promise<DrainingReadResult> DrainingReader::read(jsg::Lock& js, size_t maxRead) {
309 KJ_SWITCH_ONEOF(state) {
310 KJ_CASE_ONEOF(i, Initial) {
311 KJ_FAIL_ASSERT("this reader was never attached");
312 }
313 KJ_CASE_ONEOF(stream, Attached) {
314 auto& controller = stream->getController();
315 KJ_IF_SOME(result, controller.drainingRead(js, maxRead)) {
316 return kj::mv(result);
317 }
318 return js.rejectedPromise<DrainingReadResult>(
319 js.v8TypeError("Unable to perform draining read on this stream."_kj));
320 }
321 KJ_CASE_ONEOF(r, Released) {
322 return js.rejectedPromise<DrainingReadResult>(
323 js.v8TypeError("This ReadableStream reader has been released."_kj));
324 }
325 KJ_CASE_ONEOF(c, StreamStates::Closed) {
326 return js.resolvedPromise(DrainingReadResult{
327 .chunks = kj::Array<kj::Array<kj::byte>>(),
328 .done = true,
329 });
330 }
331 }
332 KJ_UNREACHABLE;
333}
334 
335jsg::Promise<void> DrainingReader::cancel(
336 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
337 KJ_SWITCH_ONEOF(state) {
338 KJ_CASE_ONEOF(i, Initial) {
339 KJ_FAIL_ASSERT("this reader was never attached");
340 }
341 KJ_CASE_ONEOF(stream, Attached) {
342 auto ref = stream.addRef();
343 return stream->getController().cancel(js, maybeReason);
344 }
345 KJ_CASE_ONEOF(r, Released) {
346 return js.rejectedPromise<void>(
347 js.v8TypeError("This ReadableStream reader has been released."_kj));
348 }
349 KJ_CASE_ONEOF(c, StreamStates::Closed) {
350 return js.resolvedPromise();
351 }
352 }
353 KJ_UNREACHABLE;
354}
355 
356void DrainingReader::releaseLock(jsg::Lock& js) {
357 KJ_SWITCH_ONEOF(state) {
358 KJ_CASE_ONEOF(i, Initial) {
359 KJ_FAIL_ASSERT("this reader was never attached");
360 }
361 KJ_CASE_ONEOF(stream, Attached) {
362 auto ref = stream.addRef();
363 stream->getController().releaseReader(*this, js);
364 state.init<Released>();
365 return;
366 }
367 KJ_CASE_ONEOF(c, StreamStates::Closed) {
368 return;
369 }
370 KJ_CASE_ONEOF(r, Released) {
371 return;
372 }
373 }
374 KJ_UNREACHABLE;
375}
376 
377bool DrainingReader::isAttached() const {
378 return state.is<Attached>();
379}
380 
381void DrainingReader::visitForGc(jsg::GcVisitor& visitor) {
382 KJ_IF_SOME(stream, state.tryGet<Attached>()) {
383 visitor.visit(stream);
384 }
385 visitor.visit(closedPromise);
386}
387 
388// ======================================================================================
389 
390ReadableStream::ReadableStream(IoContext& ioContext, kj::Own<ReadableStreamSource> source)
391 : ReadableStream(newReadableStreamInternalController(ioContext, kj::mv(source))) {}
392 
393ReadableStream::ReadableStream(kj::Own<ReadableStreamController> controller)
394 : ioContext(tryGetIoContext()),
395 controller(kj::mv(controller)) {
396 getController().setOwnerRef(*this);
397}
398 
399void ReadableStream::visitForGc(jsg::GcVisitor& visitor) {
400 visitor.visit(getController());
401 KJ_IF_SOME(pair, eofResolverPair) {
402 visitor.visit(pair.resolver);
403 visitor.visit(pair.promise);
404 }
405}
406 
407jsg::Ref<ReadableStream> ReadableStream::addRef() {
408 return JSG_THIS;
409}
410 
411bool ReadableStream::isDisturbed() {
412 return getController().isDisturbed();
413}
414 
415bool ReadableStream::isLocked() {
416 return getController().isLockedToReader();
417}
418 
419jsg::Promise<void> ReadableStream::onEof(jsg::Lock& js) {
420 eofResolverPair = js.newPromiseAndResolver<void>();
421 return kj::mv(KJ_ASSERT_NONNULL(eofResolverPair).promise);
422}
423 
424void ReadableStream::signalEof(jsg::Lock& js) {
425 KJ_IF_SOME(pair, eofResolverPair) {
426 pair.resolver.resolve(js);
427 }
428}
429 
430ReadableStreamController& ReadableStream::getController() {
431 return *controller;
432}
433 
434jsg::Promise<void> ReadableStream::cancel(
435 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
436 if (isLocked()) {
437 return js.rejectedPromise<void>(
438 js.v8TypeError("This ReadableStream is currently locked to a reader."_kj));
439 }
440 return getController().cancel(js, maybeReason);
441}
442 
443ReadableStream::Reader ReadableStream::getReader(
444 jsg::Lock& js, jsg::Optional<GetReaderOptions> options) {
445 JSG_REQUIRE(!isLocked(), TypeError, "This ReadableStream is currently locked to a reader.");
446 
447 bool isByob = false;
448 KJ_IF_SOME(o, options) {
449 KJ_IF_SOME(mode, o.mode) {
450 JSG_REQUIRE(
451 mode == "byob", TypeError, "mode must be undefined or 'byob' in call to getReader().");
452 // No need to check that the ReadableStream implementation is a byte stream: the first
453 // invocation of read() will do that for us and throw if necessary. Also, we should really
454 // just support reading non-byte streams with BYOB readers.
455 isByob = true;
456 }
457 }
458 
459 if (isByob) {
460 return ReadableStreamBYOBReader::constructor(js, JSG_THIS);
461 }
462 return ReadableStreamDefaultReader::constructor(js, JSG_THIS);
463}
464 
465jsg::Ref<ReadableStream::ReadableStreamAsyncIterator> ReadableStream::values(
466 jsg::Lock& js, jsg::Optional<ValuesOptions> options) {
467 static const auto defaultOptions = ValuesOptions{};
468 return js.alloc<ReadableStreamAsyncIterator>(AsyncIteratorState{.ioContext = ioContext,
469 .reader = ReadableStreamDefaultReader::constructor(js, JSG_THIS),
470 .preventCancel = options.orDefault(defaultOptions).preventCancel.orDefault(false)});
471}
472 
473jsg::Ref<ReadableStream> ReadableStream::pipeThrough(
474 jsg::Lock& js, Transform transform, jsg::Optional<PipeToOptions> maybeOptions) {
475 auto& controller = getController();
476 
477 auto& destination = transform.writable->getController();
478 JSG_REQUIRE(!isLocked(), TypeError, "This ReadableStream is currently locked to a reader.");
479 JSG_REQUIRE(!destination.isLockedToWriter(), TypeError,
480 "This WritableStream is currently locked to a writer.");
481 
482 auto options = kj::mv(maybeOptions).orDefault({});
483 options.pipeThrough = true;
484 // The lambda intentionally captures self as a visitable reference, ensuring
485 // JSG_THIS stays alive until the pipe promise resolves.
486 controller.pipeTo(js, destination, kj::mv(options))
487 .then(js,
488 JSG_VISITABLE_LAMBDA(
489 (self = JSG_THIS), (self), (jsg::Lock& js) { return js.resolvedPromise(); }))
490 .markAsHandled(js);
491 return kj::mv(transform.readable);
492}
493 
494jsg::Promise<void> ReadableStream::pipeTo(jsg::Lock& js,
495 jsg::Ref<WritableStream> destination,
496 jsg::Optional<PipeToOptions> maybeOptions) {
497 if (isLocked()) {
498 return js.rejectedPromise<void>(
499 js.v8TypeError("This ReadableStream is currently locked to a reader."_kj));
500 }
501 
502 if (destination->getController().isLockedToWriter()) {
503 return js.rejectedPromise<void>(
504 js.v8TypeError("This WritableStream is currently locked to a writer"_kj));
505 }
506 
507 auto options = kj::mv(maybeOptions).orDefault({});
508 return getController().pipeTo(js, destination->getController(), kj::mv(options));
509}
510 
511kj::Array<jsg::Ref<ReadableStream>> ReadableStream::tee(jsg::Lock& js) {
512 JSG_REQUIRE(!isLocked(), TypeError, "This ReadableStream is currently locked to a reader,");
513 auto tee = getController().tee(js);
514 return kj::arr(kj::mv(tee.branch1), kj::mv(tee.branch2));
515}
516 
517jsg::JsString ReadableStream::inspectState(jsg::Lock& js) {
518 if (controller->isClosedOrErrored()) {
519 return js.strIntern(controller->isClosed() ? "closed"_kj : "errored"_kj);
520 } else {
521 return js.strIntern("readable"_kj);
522 }
523}
524 
525bool ReadableStream::inspectSupportsBYOB() {
526 return controller->isByteOriented();
527}
528 
529jsg::Optional<uint64_t> ReadableStream::inspectLength() {
530 return tryGetLength(StreamEncoding::IDENTITY);
531}
532 
533jsg::Promise<kj::Maybe<jsg::Value>> ReadableStream::nextFunction(
534 jsg::Lock& js, AsyncIteratorState& state) {
535 return state.reader->read(js).then(
536 js, [reader = state.reader.addRef()](jsg::Lock& js, ReadResult result) mutable {
537 if (result.done) {
538 reader->releaseLock(js);
539 return js.resolvedPromise(kj::Maybe<jsg::Value>(kj::none));
540 }
541 return js.resolvedPromise<kj::Maybe<jsg::Value>>(kj::mv(result.value));
542 });
543}
544 
545jsg::Promise<void> ReadableStream::returnFunction(
546 jsg::Lock& js, AsyncIteratorState& state, jsg::Optional<jsg::Value>& value) {
547 if (state.reader.get() != nullptr) {
548 auto reader = kj::mv(state.reader);
549 if (!state.preventCancel) {
550 auto promise = reader->cancel(js, value.map([&](jsg::Value& v) { return v.getHandle(js); }));
551 reader->releaseLock(js);
552 auto result = promise.then(js,
553 JSG_VISITABLE_LAMBDA((reader = kj::mv(reader)), (reader), (jsg::Lock& js) {
554 // Ensure that the reader is not garbage collected until the cancel promise resolves.
555 return js.resolvedPromise();
556 }));
557 // When the stream is already errored, cancel() returns a rejected promise
558 // that propagates through the .then() chain. Mark it as handled so V8 does
559 // not fire unhandledrejection events during iterator teardown.
560 result.markAsHandled(js);
561 return kj::mv(result);
562 }
563 
564 reader->releaseLock(js);
565 }
566 return js.resolvedPromise();
567}
568 
569jsg::Ref<ReadableStream> ReadableStream::detach(jsg::Lock& js, bool ignoreDisturbed) {
570 JSG_REQUIRE(
571 !isDisturbed() || ignoreDisturbed, TypeError, "The ReadableStream has already been read.");
572 JSG_REQUIRE(!isLocked(), TypeError, "The ReadableStream has been locked to a reader.");
573 return js.alloc<ReadableStream>(getController().detach(js, ignoreDisturbed));
574}
575 
576kj::Maybe<uint64_t> ReadableStream::tryGetLength(StreamEncoding encoding) {
577 return getController().tryGetLength(encoding);
578}
579 
580kj::Promise<DeferredProxy<void>> ReadableStream::pumpTo(
581 jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) {
582 JSG_REQUIRE(
583 IoContext::hasCurrent(), Error, "Unable to consume this ReadableStream outside of a request");
584 JSG_REQUIRE(!isLocked(), TypeError, "The ReadableStream has been locked to a reader.");
585 return getController().pumpTo(js, kj::mv(sink), end);
586}
587 
588jsg::Ref<ReadableStream> ReadableStream::constructor(jsg::Lock& js,
589 jsg::Optional<UnderlyingSource> underlyingSource,
590 jsg::Optional<StreamQueuingStrategy> queuingStrategy) {
591 
592 JSG_REQUIRE(FeatureFlags::get(js).getStreamsJavaScriptControllers(), Error,
593 "To use the new ReadableStream() constructor, enable the "
594 "streams_enable_constructors compatibility flag. "
595 "Refer to the docs for more information: https://developers.cloudflare.com/workers/platform/compatibility-dates/#compatibility-flags");
596 // We account for the memory usage of the ReadableStream and its controller together because their
597 // lifetimes are identical and memory accounting itself has a memory overhead.
598 auto controller = newReadableStreamJsController();
599 auto stream = js.allocAccounted<ReadableStream>(
600 sizeof(ReadableStream) + controller->jsgGetMemorySelfSize(), kj::mv(controller));
601 stream->getController().setup(js, kj::mv(underlyingSource), kj::mv(queuingStrategy));
602 return kj::mv(stream);
603}
604 
605jsg::Optional<uint32_t> ByteLengthQueuingStrategy::size(
606 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeValue) {
607 KJ_IF_SOME(value, maybeValue) {
608 if ((value)->IsArrayBuffer()) {
609 auto buffer = value.As<v8::ArrayBuffer>();
610 return buffer->ByteLength();
611 } else if ((value)->IsArrayBufferView()) {
612 auto view = value.As<v8::ArrayBufferView>();
613 return view->ByteLength();
614 } else {
615 // Per the WHATWG Streams spec, ByteLengthQueuingStrategy.size should return
616 // GetV(chunk, "byteLength"), which means getting the byteLength property
617 // from any object, not just ArrayBuffer/ArrayBufferView.
618 KJ_IF_SOME(obj, jsg::JsValue(value).tryCast<jsg::JsObject>()) {
619 auto byteLength = obj.get(js, "byteLength"_kj);
620 KJ_IF_SOME(num, byteLength.tryCast<jsg::JsNumber>()) {
621 KJ_IF_SOME(val, num.value(js)) {
622 return static_cast<uint32_t>(val);
623 }
624 }
625 }
626 }
627 }
628 return kj::none;
629}
630 
631namespace {
632 
633// TODO(cleanup): These classes have been copied to external-pusher.c++. The copies here can be
634// deleted as soon as we've switched from StreamSink to ExternalPusher and can delete all the
635// StreamSink-related code. For now I'm not trying to avoid duplication.
636 
637// HACK: We need as async pipe, like kj::newOneWayPipe(), except supporting explicit end(). So we
638// wrap the two ends of the pipe in special adapters that track whether end() was called.
639class ExplicitEndOutputPipeAdapter final: public capnp::ExplicitEndOutputStream {
640 public:
641 ExplicitEndOutputPipeAdapter(
642 kj::Own<kj::AsyncOutputStream> inner, kj::Own<kj::RefcountedWrapper<bool>> ended)
643 : inner(kj::mv(inner)),
644 ended(kj::mv(ended)) {}
645 
646 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
647 return KJ_REQUIRE_NONNULL(inner)->write(buffer);
648 }
649 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
650 return KJ_REQUIRE_NONNULL(inner)->write(pieces);
651 }
652 
653 kj::Maybe<kj::Promise<uint64_t>> tryPumpFrom(
654 kj::AsyncInputStream& input, uint64_t amount) override {
655 return KJ_REQUIRE_NONNULL(inner)->tryPumpFrom(input, amount);
656 }
657 
658 kj::Promise<void> whenWriteDisconnected() override {
659 return KJ_REQUIRE_NONNULL(inner)->whenWriteDisconnected();
660 }
661 
662 kj::Promise<void> end() override {
663 // Signal to the other side that end() was actually called.
664 ended->getWrapped() = true;
665 inner = kj::none;
666 return kj::READY_NOW;
667 }
668 
669 private:
670 kj::Maybe<kj::Own<kj::AsyncOutputStream>> inner;
671 kj::Own<kj::RefcountedWrapper<bool>> ended;
672};
673 
674class ExplicitEndInputPipeAdapter final: public kj::AsyncInputStream {
675 public:
676 ExplicitEndInputPipeAdapter(kj::Own<kj::AsyncInputStream> inner,
677 kj::Own<kj::RefcountedWrapper<bool>> ended,
678 kj::Maybe<uint64_t> expectedLength)
679 : inner(kj::mv(inner)),
680 ended(kj::mv(ended)),
681 expectedLength(expectedLength) {}
682 
683 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
684 size_t result = co_await inner->tryRead(buffer, minBytes, maxBytes);
685 
686 KJ_IF_SOME(l, expectedLength) {
687 KJ_ASSERT(result <= l);
688 l -= result;
689 if (l == 0) {
690 // If we got all the bytes we expected, we treat this as a successful end, because the
691 // underlying KJ pipe is not actually going to wait for the other side to drop. This is
692 // consistent with the behavior of Content-Length in HTTP anyway.
693 ended->getWrapped() = true;
694 }
695 }
696 
697 if (result < minBytes) {
698 // Verify that end() was called.
699 if (!ended->getWrapped()) {
700 JSG_FAIL_REQUIRE(Error, "ReadableStream received over RPC disconnected prematurely.");
701 }
702 }
703 co_return result;
704 }
705 
706 kj::Maybe<uint64_t> tryGetLength() override {
707 return inner->tryGetLength();
708 }
709 
710 kj::Promise<uint64_t> pumpTo(kj::AsyncOutputStream& output, uint64_t amount) override {
711 return inner->pumpTo(output, amount);
712 }
713 
714 private:
715 kj::Own<kj::AsyncInputStream> inner;
716 kj::Own<kj::RefcountedWrapper<bool>> ended;
717 kj::Maybe<uint64_t> expectedLength;
718};
719 
720// Wrapper around ReadableStreamSource that prevents deferred proxying. We need this for RPC
721// streams because although they are "system streams", they become disconnected when the IoContext
722// is destroyed, due to the JsRpcCustomEvent being canceled.
723//
724// TODO(someday): Devise a better way for RPC streams to extend the lifetime of the RPC session
725// beyond the destruction of the IoContext, if it is being used for deferred proxying.
726class NoDeferredProxyReadableStream final: public ReadableStreamSource {
727 public:
728 NoDeferredProxyReadableStream(kj::Own<ReadableStreamSource> inner, IoContext& ioctx)
729 : inner(kj::mv(inner)),
730 ioctx(ioctx) {}
731 
732 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
733 return inner->tryRead(buffer, minBytes, maxBytes);
734 }
735 
736 kj::Promise<DeferredProxy<void>> pumpTo(WritableStreamSink& output, bool end) override {
737 // Move the deferred proxy part of the task over to the non-deferred part. To do this,
738 // we use `ioctx.waitForDeferredProxy()`, which returns a single promise covering both parts
739 // (and, importantly, registering pending events where needed). Then, we add a noop deferred
740 // proxy to the end of that.
741 return addNoopDeferredProxy(ioctx.waitForDeferredProxy(inner->pumpTo(output, end)));
742 }
743 
744 StreamEncoding getPreferredEncoding() override {
745 return inner->getPreferredEncoding();
746 }
747 
748 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override {
749 return inner->tryGetLength(encoding);
750 }
751 
752 void cancel(kj::Exception reason) override {
753 return inner->cancel(kj::mv(reason));
754 }
755 
756 kj::Maybe<Tee> tryTee(uint64_t limit) override {
757 return inner->tryTee(limit).map([&](Tee tee) {
758 return Tee{.branches = {
759 kj::heap<NoDeferredProxyReadableStream>(kj::mv(tee.branches[0]), ioctx),
760 kj::heap<NoDeferredProxyReadableStream>(kj::mv(tee.branches[1]), ioctx),
761 }};
762 });
763 }
764 
765 private:
766 kj::Own<ReadableStreamSource> inner;
767 IoContext& ioctx;
768};
769 
770} // namespace
771 
772void ReadableStream::serialize(jsg::Lock& js, jsg::Serializer& serializer) {
773 // Serialize by effectively creating a `JsRpcStub` around this object and serializing that.
774 // Except we don't actually want to do _exactly_ that, because we do not want to actually create
775 // a `JsRpcStub` locally. So do the important parts of `JsRpcStub::constructor()` followed by
776 // `JsRpcStub::serialize()`.
777 
778 auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError,
779 "ReadableStream can only be serialized for RPC.");
780 auto externalHandler = dynamic_cast<RpcSerializerExternalHandler*>(&handler);
781 JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError,
782 "ReadableStream can only be serialized for RPC.");
783 
784 // NOTE: We're counting on `pumpTo()`, below, to check that the stream is not locked or disturbed
785 // and other common checks. It's important that we don't modify the stream in any way before
786 // that call.
787 
788 IoContext& ioctx = IoContext::current();
789 
790 auto& controller = getController();
791 StreamEncoding encoding = controller.getPreferredEncoding();
792 auto expectedLength = controller.tryGetLength(encoding);
793 
794 capnp::ByteStream::Client streamCap = [&]() {
795 KJ_IF_SOME(pusher, externalHandler->getExternalPusher()) {
796 auto req = pusher.pushByteStreamRequest(capnp::MessageSize{2, 0});
797 KJ_IF_SOME(el, expectedLength) {
798 req.setLengthPlusOne(el + 1);
799 }
800 auto pipeline = req.sendForPipeline();
801 
802 externalHandler->write([encoding, expectedLength, source = pipeline.getSource()](
803 rpc::JsValue::External::Builder builder) mutable {
804 auto rs = builder.initReadableStream();
805 rs.setStream(kj::mv(source));
806 rs.setEncoding(encoding);
807 });
808 
809 return pipeline.getSink();
810 } else {
811 return externalHandler
812 ->writeStream(
813 [encoding, expectedLength](rpc::JsValue::External::Builder builder) mutable {
814 auto rs = builder.initReadableStream();
815 rs.setEncoding(encoding);
816 KJ_IF_SOME(l, expectedLength) {
817 rs.getExpectedLength().setKnown(l);
818 }
819 }).castAs<capnp::ByteStream>();
820 }
821 }();
822 
823 kj::Own<capnp::ExplicitEndOutputStream> kjStream =
824 ioctx.getByteStreamFactory().capnpToKjExplicitEnd(kj::mv(streamCap));
825 
826 auto sink = newSystemStream(kj::mv(kjStream), encoding, ioctx);
827 
828 ioctx.addTask(
829 ioctx.waitForDeferredProxy(pumpTo(js, kj::mv(sink), true)).catch_([](kj::Exception&& e) {
830 // Errors in pumpTo() are automatically propagated to the source and destination. We don't
831 // want to throw them from here since it'll cause an uncaught exception to be reported, even
832 // if the application actually does handle it!
833 }));
834}
835 
836jsg::Ref<ReadableStream> ReadableStream::deserialize(
837 jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) {
838 auto& handler = KJ_REQUIRE_NONNULL(
839 deserializer.getExternalHandler(), "got ReadableStream on non-RPC serialized object?");
840 auto externalHandler = dynamic_cast<RpcDeserializerExternalHandler*>(&handler);
841 KJ_REQUIRE(externalHandler != nullptr, "got ReadableStream on non-RPC serialized object?");
842 
843 auto reader = externalHandler->read();
844 KJ_REQUIRE(reader.isReadableStream(), "external table slot type doesn't match serialization tag");
845 
846 auto rs = reader.getReadableStream();
847 auto encoding = rs.getEncoding();
848 
849 KJ_REQUIRE(
850 static_cast<uint>(encoding) < capnp::Schema::from<StreamEncoding>().getEnumerants().size(),
851 "unknown StreamEncoding received from peer");
852 
853 auto& ioctx = IoContext::current();
854 
855 kj::Own<kj::AsyncInputStream> in;
856 if (rs.hasStream()) {
857 in =
858 ioctx.getExternalPusher()->unwrapStream(rs.getStream(), externalHandler->getDebugContext());
859 } else {
860 kj::Maybe<uint64_t> expectedLength;
861 auto el = rs.getExpectedLength();
862 if (el.isKnown()) {
863 expectedLength = el.getKnown();
864 }
865 
866 auto pipe = kj::newOneWayPipe(expectedLength);
867 
868 auto endedFlag = kj::refcounted<kj::RefcountedWrapper<bool>>(false);
869 
870 auto out = kj::heap<ExplicitEndOutputPipeAdapter>(kj::mv(pipe.out), kj::addRef(*endedFlag));
871 in = kj::heap<ExplicitEndInputPipeAdapter>(kj::mv(pipe.in), kj::mv(endedFlag), expectedLength);
872 
873 externalHandler->setLastStream(ioctx.getByteStreamFactory().kjToCapnp(kj::mv(out)));
874 }
875 
876 return js.alloc<ReadableStream>(ioctx,
877 kj::heap<NoDeferredProxyReadableStream>(newSystemStream(kj::mv(in), encoding, ioctx), ioctx));
878}
879 
880kj::StringPtr ReaderImpl::jsgGetMemoryName() const {
881 return "ReaderImpl"_kjc;
882}
883 
884size_t ReaderImpl::jsgGetMemorySelfSize() const {
885 return sizeof(ReaderImpl);
886}
887 
888void ReaderImpl::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
889 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
890 tracker.trackField("stream", attached.stream);
891 }
892 tracker.trackField("closedPromise", closedPromise);
893}
894 
895void ReadableStream::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
896 tracker.trackField("controller", controller);
897 tracker.trackField("eofResolverPair", eofResolverPair);
898}
899 
900} // namespace workerd::api