Skip to content
File

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

187.3 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 "standard.h"
6 
7#include "readable.h"
8#include "writable.h"
9 
10#include <workerd/io/features.h>
11#include <workerd/jsg/jsg.h>
12#include <workerd/util/autogate.h>
13#include <workerd/util/state-machine.h>
14#include <workerd/util/weak-refs.h>
15 
16#include <kj/debug.h>
17#include <kj/vector.h>
18 
19namespace workerd::api {
20 
21using DefaultController = jsg::Ref<ReadableStreamDefaultController>;
22using ByobController = jsg::Ref<ReadableByteStreamController>;
23 
24namespace {
25struct ValueReadable;
26struct ByteReadable;
27} // namespace
28 
29// =======================================================================================
30// The Unlocked, Locked, ReaderLocked, and WriterLocked structs
31// are used to track the current lock status of JavaScript-backed streams.
32// All readable and writable streams begin in the Unlocked state. When a
33// reader or writer are attached, the streams will transition into the
34// ReaderLocked or WriterLocked state. When the reader is released, those
35// will transition back to Unlocked.
36//
37// When a readable is piped to a writable, both will enter the PipeLocked state.
38// (PipeLocked is defined within the ReadableLockImpl and WritableLockImpl classes
39// below) When the pipe completes, both will transition back to Unlocked.
40//
41// When a ReadableStreamJsController is tee()'d, it will enter the locked state.
42 
43namespace {
44 
45// A utility class used by ReadableStreamJsController
46// for implementing the reader lock in a consistent way (without duplicating any code).
47template <typename Controller>
48class ReadableLockImpl {
49 public:
50 using PipeController = ReadableStreamController::PipeController;
51 using Reader = ReadableStreamController::Reader;
52 
53 bool isLockedToReader() const {
54 return !state.template is<Unlocked>();
55 }
56 
57 bool lockReader(jsg::Lock& js, Controller& self, Reader& reader);
58 
59 // See the comment for releaseReader in common.h for details on the use of maybeJs
60 void releaseReader(Controller& self, Reader& reader, kj::Maybe<jsg::Lock&> maybeJs);
61 
62 bool lock();
63 
64 void onClose(jsg::Lock& js);
65 void onError(jsg::Lock& js, v8::Local<v8::Value> reason);
66 
67 kj::Maybe<PipeController&> tryPipeLock(Controller& self);
68 
69 void visitForGc(jsg::GcVisitor& visitor);
70 
71 kj::StringPtr jsgGetMemoryName() const {
72 return "ReadableLockImpl"_kjc;
73 }
74 size_t jsgGetMemorySelfSize() const {
75 return sizeof(ReadableLockImpl);
76 }
77 void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
78 KJ_SWITCH_ONEOF(state) {
79 KJ_CASE_ONEOF(locked, Locked) {}
80 KJ_CASE_ONEOF(unlocked, Unlocked) {}
81 KJ_CASE_ONEOF(pipeLocked, PipeLocked) {}
82 KJ_CASE_ONEOF(readerLocked, ReaderLocked) {
83 tracker.trackField("readerLocked", readerLocked);
84 }
85 }
86 }
87 
88 private:
89 class PipeLocked final: public PipeController {
90 public:
91 static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj;
92 explicit PipeLocked(Controller& inner): inner(inner) {}
93 
94 bool isClosed() override {
95 return inner.state.template is<StreamStates::Closed>();
96 }
97 
98 kj::Maybe<v8::Local<v8::Value>> tryGetErrored(jsg::Lock& js) override {
99 KJ_IF_SOME(errored, inner.state.template tryGetUnsafe<StreamStates::Errored>()) {
100 return errored.getHandle(js);
101 }
102 return kj::none;
103 }
104 
105 void cancel(jsg::Lock& js, v8::Local<v8::Value> reason) override {
106 // Cancel here returns a Promise but we do not need to propagate it.
107 // We can safely drop it on the floor here.
108 auto promise KJ_UNUSED = inner.cancel(js, reason);
109 }
110 
111 void close(jsg::Lock& js) override {
112 inner.doClose(js);
113 }
114 
115 void error(jsg::Lock& js, v8::Local<v8::Value> reason) override {
116 inner.doError(js, reason);
117 }
118 
119 void release(jsg::Lock& js, kj::Maybe<v8::Local<v8::Value>> maybeError = kj::none) override {
120 KJ_IF_SOME(error, maybeError) {
121 cancel(js, error);
122 }
123 inner.lock.state.template transitionTo<Unlocked>();
124 }
125 
126 kj::Maybe<kj::Promise<void>> tryPumpTo(WritableStreamSink& sink, bool end) override;
127 
128 jsg::Promise<ReadResult> read(jsg::Lock& js) override;
129 
130 private:
131 Controller& inner;
132 
133 friend Controller;
134 };
135 
136 // State machine for ReadableLockImpl:
137 // All states can transition to any other state (no terminal states).
138 // Unlocked -> Locked (lock() called for tee)
139 // Unlocked -> ReaderLocked (lockReader() called)
140 // Unlocked -> PipeLocked (tryPipeLock() called)
141 // ReaderLocked -> Unlocked (releaseReader() called)
142 // PipeLocked -> Unlocked (release() or onClose/onError called)
143 // Locked -> (remains until stream is done)
144 using LockState = StateMachine<Locked, PipeLocked, ReaderLocked, Unlocked>;
145 LockState state = LockState::template create<Unlocked>();
146 friend Controller;
147};
148 
149// A utility class used by WritableStreamJsController to implement the writer lock
150// mechanism. Extracted for consistency with ReadableStreamJsController and to
151// eventually allow it to be shared also with WritableStreamInternalController.
152template <typename Controller>
153class WritableLockImpl {
154 public:
155 using Writer = WritableStreamController::Writer;
156 
157 bool isLockedToWriter() const;
158 
159 bool lockWriter(jsg::Lock& js, Controller& self, Writer& writer);
160 
161 // See the comment for releaseWriter in common.h for details on the use of maybeJs
162 void releaseWriter(Controller& self, Writer& writer, kj::Maybe<jsg::Lock&> maybeJs);
163 
164 void visitForGc(jsg::GcVisitor& visitor);
165 
166 bool pipeLock(WritableStream& owner, jsg::Ref<ReadableStream> source, PipeToOptions& options);
167 void releasePipeLock();
168 
169 JSG_MEMORY_INFO(WritableLockImpl) {
170 KJ_SWITCH_ONEOF(state) {
171 KJ_CASE_ONEOF(unlocked, Unlocked) {}
172 KJ_CASE_ONEOF(locked, Locked) {}
173 KJ_CASE_ONEOF(writerLocked, WriterLocked) {
174 tracker.trackField("writerLocked", writerLocked);
175 }
176 KJ_CASE_ONEOF(pipeLocked, PipeLocked) {
177 tracker.trackField("pipeLocked", pipeLocked);
178 }
179 }
180 }
181 
182 private:
183 struct PipeLocked {
184 static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj;
185 ReadableStreamController::PipeController& source;
186 jsg::Ref<ReadableStream> readableStreamRef;
187 
188 kj::Maybe<jsg::Ref<AbortSignal>> maybeSignal;
189 
190 kj::Maybe<jsg::Promise<void>> checkSignal(jsg::Lock& js, Controller& self);
191 
192 struct Flags {
193 uint8_t preventAbort : 1 = 0;
194 uint8_t preventCancel : 1 = 0;
195 uint8_t preventClose : 1 = 0;
196 uint8_t pipeThrough : 1 = 0;
197 };
198 Flags flags{};
199 
200 JSG_MEMORY_INFO(PipeLocked) {
201 tracker.trackField("readableStreamRef", readableStreamRef);
202 tracker.trackField("signal", maybeSignal);
203 }
204 };
205 
206 // State machine for WritableLockImpl:
207 // All states can transition to any other state (no terminal states).
208 // Unlocked -> Locked (not currently used)
209 // Unlocked -> WriterLocked (lockWriter() called)
210 // Unlocked -> PipeLocked (pipeLock() called)
211 // WriterLocked -> Unlocked (releaseWriter() called)
212 // PipeLocked -> Unlocked (releasePipeLock() called)
213 using LockState = StateMachine<Unlocked, Locked, WriterLocked, PipeLocked>;
214 LockState state = LockState::template create<Unlocked>();
215 
216 inline kj::Maybe<PipeLocked&> tryGetPipe() {
217 KJ_IF_SOME(locked, state.template tryGetUnsafe<PipeLocked>()) {
218 return locked;
219 }
220 return kj::none;
221 }
222 
223 friend Controller;
224};
225 
226// ======================================================================================
227 
228template <typename Controller>
229bool ReadableLockImpl<Controller>::lock() {
230 if (isLockedToReader()) {
231 return false;
232 }
233 
234 state.template transitionTo<Locked>();
235 return true;
236}
237 
238template <typename Controller>
239bool ReadableLockImpl<Controller>::lockReader(jsg::Lock& js, Controller& self, Reader& reader) {
240 if (isLockedToReader()) {
241 return false;
242 }
243 
244 auto prp = js.newPromiseAndResolver<void>();
245 prp.promise.markAsHandled(js);
246 
247 auto lock = ReaderLocked(reader, kj::mv(prp.resolver));
248 
249 if (self.state.template is<StreamStates::Closed>()) {
250 maybeResolvePromise(js, lock.getClosedFulfiller());
251 } else KJ_IF_SOME(errored, self.state.template tryGetUnsafe<StreamStates::Errored>()) {
252 maybeRejectPromise<void>(js, lock.getClosedFulfiller(), errored.getHandle(js));
253 }
254 
255 state.template transitionTo<ReaderLocked>(kj::mv(lock));
256 reader.attach(self, kj::mv(prp.promise));
257 return true;
258}
259 
260template <typename Controller>
261void ReadableLockImpl<Controller>::releaseReader(
262 Controller& self, Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) {
263 KJ_IF_SOME(locked, state.template tryGetUnsafe<ReaderLocked>()) {
264 KJ_ASSERT(&locked.getReader() == &reader);
265 
266 KJ_IF_SOME(js, maybeJs) {
267 auto reason = js.typeError("This ReadableStream reader has been released."_kj);
268 KJ_SWITCH_ONEOF(self.state) {
269 KJ_CASE_ONEOF(initial, typename Controller::Initial) {}
270 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
271 KJ_CASE_ONEOF(errored, StreamStates::Errored) {}
272 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
273 consumer->cancelPendingReads(js, reason);
274 }
275 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
276 consumer->cancelPendingReads(js, reason);
277 }
278 }
279 maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason);
280 }
281 
282 // Keep the locked.clear() after the isolate and hasPendingReadRequests check above.
283 // Clearing will release the references and we don't want to do that if the
284 // hasPendingReadRequests check fails.
285 locked.clear();
286 
287 // When maybeJs is nullptr, that means releaseReader was called when the reader is
288 // being deconstructed and not as the result of explicitly calling releaseLock and
289 // we do not have an isolate lock. In that case, we don't want to change the lock
290 // state itself. Moving the lock above will free the lock state while keeping the
291 // ReadableStream marked as locked.
292 if (maybeJs != kj::none) {
293 state.template transitionTo<Unlocked>();
294 }
295 }
296}
297 
298template <typename Controller>
299kj::Maybe<ReadableStreamController::PipeController&> ReadableLockImpl<Controller>::tryPipeLock(
300 Controller& self) {
301 if (isLockedToReader()) {
302 return kj::none;
303 }
304 return state.template transitionTo<PipeLocked>(self);
305}
306 
307template <typename Controller>
308void ReadableLockImpl<Controller>::visitForGc(jsg::GcVisitor& visitor) {
309 KJ_SWITCH_ONEOF(state) {
310 KJ_CASE_ONEOF(locked, Locked) {}
311 KJ_CASE_ONEOF(locked, Unlocked) {}
312 KJ_CASE_ONEOF(locked, PipeLocked) {}
313 KJ_CASE_ONEOF(locked, ReaderLocked) {
314 visitor.visit(locked);
315 }
316 }
317}
318 
319template <typename Controller>
320void ReadableLockImpl<Controller>::onClose(jsg::Lock& js) {
321 KJ_IF_SOME(locked, state.template tryGetUnsafe<ReaderLocked>()) {
322 try {
323 maybeResolvePromise(js, locked.getClosedFulfiller());
324 } catch (jsg::JsExceptionThrown&) {
325 // Resolving the promise could end up throwing an exception in some cases,
326 // causing a jsg::JsExceptionThrown to be thrown. At this point, however,
327 // we are already in the process of closing the stream and an error at this
328 // point is not recoverable. Log and move on.
329 LOG_NOSENTRY(ERROR, "Error resolving ReadableStream reader closed promise");
330 };
331 } else {
332 (void)state.template transitionFromTo<PipeLocked, Unlocked>();
333 }
334}
335 
336template <typename Controller>
337void ReadableLockImpl<Controller>::onError(jsg::Lock& js, v8::Local<v8::Value> reason) {
338 KJ_IF_SOME(locked, state.template tryGetUnsafe<ReaderLocked>()) {
339 try {
340 maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason);
341 } catch (jsg::JsExceptionThrown&) {
342 // Rejecting the promise could end up throwing an exception in some cases,
343 // causing a jsg::JsExceptionThrown to be thrown. At this point, however,
344 // we are already in the process of closing the stream and an error at this
345 // point is not recoverable. Log and move on.
346 LOG_NOSENTRY(ERROR, "Error rejecting ReadableStream reader closed promise");
347 }
348 } else {
349 (void)state.template transitionFromTo<PipeLocked, Unlocked>();
350 }
351}
352 
353template <typename Controller>
354kj::Maybe<kj::Promise<void>> ReadableLockImpl<Controller>::PipeLocked::tryPumpTo(
355 WritableStreamSink& sink, bool end) {
356 // We return nullptr here because this controller does not support kj's pumpTo.
357 return kj::none;
358}
359 
360template <typename Controller>
361jsg::Promise<ReadResult> ReadableLockImpl<Controller>::PipeLocked::read(jsg::Lock& js) {
362 return KJ_ASSERT_NONNULL(inner.read(js, kj::none));
363}
364 
365// ======================================================================================
366 
367template <typename Controller>
368bool WritableLockImpl<Controller>::isLockedToWriter() const {
369 return !state.template is<Unlocked>();
370}
371 
372template <typename Controller>
373bool WritableLockImpl<Controller>::lockWriter(jsg::Lock& js, Controller& self, Writer& writer) {
374 if (isLockedToWriter()) {
375 return false;
376 }
377 
378 auto closedPrp = js.newPromiseAndResolver<void>();
379 closedPrp.promise.markAsHandled(js);
380 auto readyPrp = js.newPromiseAndResolver<void>();
381 readyPrp.promise.markAsHandled(js);
382 
383 auto lock = WriterLocked(writer, kj::mv(closedPrp.resolver), kj::mv(readyPrp.resolver));
384 
385 if (self.state.template is<StreamStates::Closed>()) {
386 maybeResolvePromise(js, lock.getClosedFulfiller());
387 maybeResolvePromise(js, lock.getReadyFulfiller());
388 } else KJ_IF_SOME(errored, self.state.template tryGetUnsafe<StreamStates::Errored>()) {
389 maybeRejectPromise<void>(js, lock.getClosedFulfiller(), errored.getHandle(js));
390 maybeRejectPromise<void>(js, lock.getReadyFulfiller(), errored.getHandle(js));
391 } else {
392 if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) {
393 // Per spec (SetUpWritableStreamDefaultWriter step 4), the ready promise
394 // is resolved when the stream is writable and not experiencing backpressure,
395 // regardless of whether the start algorithm has completed. The backpressure
396 // state is set synchronously during SetUpWritableStreamDefaultController.
397 KJ_IF_SOME(erroring, self.isErroring(js)) {
398 maybeRejectPromise<void>(js, lock.getReadyFulfiller(), erroring);
399 } else if (!self.hasBackpressure()) {
400 maybeResolvePromise(js, lock.getReadyFulfiller());
401 }
402 } else {
403 if (self.isStarted()) {
404 maybeResolvePromise(js, lock.getReadyFulfiller());
405 }
406 }
407 }
408 
409 state.template transitionTo<WriterLocked>(kj::mv(lock));
410 writer.attach(js, self, kj::mv(closedPrp.promise), kj::mv(readyPrp.promise));
411 return true;
412}
413 
414template <typename Controller>
415void WritableLockImpl<Controller>::releaseWriter(
416 Controller& self, Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) {
417 KJ_IF_SOME(locked, state.template tryGetUnsafe<WriterLocked>()) {
418 KJ_ASSERT(&locked.getWriter() == &writer);
419 KJ_IF_SOME(js, maybeJs) {
420 KJ_SWITCH_ONEOF(self.state) {
421 KJ_CASE_ONEOF(initial, typename Controller::Initial) {}
422 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
423 KJ_CASE_ONEOF(errored, StreamStates::Errored) {}
424 KJ_CASE_ONEOF(controller, jsg::Ref<WritableStreamDefaultController>) {
425 controller->cancelPendingWrites(
426 js, js.typeError("This WritableStream writer has been released."_kjc));
427 }
428 }
429 
430 // Per spec (WritableStreamDefaultWriterRelease), both the ready and closed
431 // promises must be rejected when the writer is released.
432 auto releaseReason = js.v8TypeError("This WritableStream writer has been released."_kjc);
433 if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) {
434 if (locked.getReadyFulfiller() != kj::none) {
435 maybeRejectPromise<void>(js, locked.getReadyFulfiller(), releaseReason);
436 } else {
437 // The ready fulfiller was already consumed (promise was resolved).
438 // Per spec (WritableStreamDefaultWriterEnsureReadyPromiseRejected),
439 // we must replace it with a new rejected promise.
440 auto prp = js.newPromiseAndResolver<void>();
441 prp.promise.markAsHandled(js);
442 prp.resolver.reject(js, releaseReason);
443 locked.setReadyFulfiller(js, prp);
444 }
445 } else {
446 maybeRejectPromise<void>(js, locked.getReadyFulfiller(), releaseReason);
447 }
448 maybeRejectPromise<void>(js, locked.getClosedFulfiller(), releaseReason);
449 }
450 locked.clear();
451 
452 // When maybeJs is nullptr, that means releaseWriter was called when the writer is
453 // being deconstructed and not as the result of explicitly calling releaseLock and
454 // we do not have an isolate lock. In that case, we don't want to change the lock
455 // state itself. Moving the lock above will free the lock state while keeping the
456 // WritableStream marked as locked.
457 if (maybeJs != kj::none) {
458 state.template transitionTo<Unlocked>();
459 }
460 }
461}
462 
463template <typename Controller>
464bool WritableLockImpl<Controller>::pipeLock(
465 WritableStream& owner, jsg::Ref<ReadableStream> source, PipeToOptions& options) {
466 if (isLockedToWriter()) {
467 return false;
468 }
469 
470 auto& sourceLock = KJ_ASSERT_NONNULL(source->getController().tryPipeLock());
471 
472 state.template transitionTo<PipeLocked>(PipeLocked{
473 .source = sourceLock,
474 .readableStreamRef = kj::mv(source),
475 .maybeSignal = kj::mv(options.signal),
476 .flags =
477 {
478 .preventAbort = options.preventAbort.orDefault(false),
479 .preventCancel = options.preventCancel.orDefault(false),
480 .preventClose = options.preventClose.orDefault(false),
481 .pipeThrough = options.pipeThrough,
482 },
483 });
484 return true;
485}
486 
487template <typename Controller>
488void WritableLockImpl<Controller>::releasePipeLock() {
489 if (state.template is<PipeLocked>()) {
490 state.template transitionTo<Unlocked>();
491 }
492}
493 
494template <typename Controller>
495void WritableLockImpl<Controller>::visitForGc(jsg::GcVisitor& visitor) {
496 KJ_SWITCH_ONEOF(state) {
497 KJ_CASE_ONEOF(locked, Unlocked) {}
498 KJ_CASE_ONEOF(locked, Locked) {}
499 KJ_CASE_ONEOF(locked, WriterLocked) {
500 visitor.visit(locked);
501 }
502 KJ_CASE_ONEOF(locked, PipeLocked) {
503 visitor.visit(locked.readableStreamRef);
504 KJ_IF_SOME(signal, locked.maybeSignal) {
505 visitor.visit(signal);
506 }
507 }
508 }
509}
510 
511template <typename Controller>
512kj::Maybe<jsg::Promise<void>> WritableLockImpl<Controller>::PipeLocked::checkSignal(
513 jsg::Lock& js, Controller& self) {
514 KJ_IF_SOME(signal, maybeSignal) {
515 if (signal->getAborted(js)) {
516 auto reason = signal->getReason(js);
517 if (!flags.preventCancel) {
518 source.release(js, v8::Local<v8::Value>(reason));
519 } else {
520 source.release(js);
521 }
522 if (!flags.preventAbort) {
523 return self.abort(js, reason).then(js, JSG_VISITABLE_LAMBDA((this, reason = reason.addRef(js), ref = self.addRef()), (reason, ref), (jsg::Lock& js) {
524 return rejectedMaybeHandledPromise<void>(js, reason.getHandle(js), flags.pipeThrough);
525 }));
526 }
527 return rejectedMaybeHandledPromise<void>(js, reason, flags.pipeThrough);
528 }
529 }
530 return kj::none;
531}
532 
533auto maybeAddFunctor(jsg::Lock& js, auto promise, auto onSuccess, auto onFailure) {
534 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
535 return promise.then(
536 js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure)));
537 } else {
538 return promise.then(js, kj::mv(onSuccess), kj::mv(onFailure));
539 }
540}
541 
542jsg::Promise<void> maybeRunAlgorithm(
543 jsg::Lock& js, auto& maybeAlgorithm, auto&& onSuccess, auto&& onFailure, auto&&... args) {
544 // The algorithm is a JavaScript function mapped through jsg::Function.
545 // It is expected to return a Promise mapped via jsg::Promise. If the
546 // function returns synchronously, the jsg::Promise wrapper ensures
547 // that it is properly mapped to a jsg::Promise, but if the Promise
548 // throws synchronously, we have to convert that synchronous throw
549 // into a proper rejected jsg::Promise.
550 KJ_IF_SOME(algorithm, maybeAlgorithm) {
551 // We need two layers of JSG_TRY here, unfortunately. The inner layer
552 // covers the algorithm implementation itself and is our typical error
553 // handling path. It ensures that if the algorithm throws an exception,
554 // that is properly converted in to a rejected promise that is *then*
555 // handled by the onFailure handler that is passed in. The outer JSG_TRY
556 // handles the rare and generally unexpected failure of the calls to
557 // .then() itself, which can throw JS exceptions synchronously in certain
558 // rare cases. For those we return a rejected promise but do not call the
559 // onFailure case since such errors are generally indicative of a fatal
560 // condition in the isolate (e.g. out of memory, other fatal exception, etc).
561 JSG_TRY(js) {
562 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
563 auto getInnerPromise = [&]() -> jsg::Promise<void> {
564 JSG_TRY(js) {
565 return algorithm(js, kj::fwd<decltype(args)>(args)...);
566 }
567 JSG_CATCH(exception) {
568 return js.rejectedPromise<void>(kj::mv(exception));
569 }
570 };
571 return getInnerPromise().then(
572 js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure)));
573 } else {
574 auto getInnerPromise = [&]() -> jsg::Promise<void> {
575 JSG_TRY(js) {
576 return algorithm(js, kj::fwd<decltype(args)>(args)...);
577 }
578 JSG_CATCH(exception) {
579 return js.rejectedPromise<void>(kj::mv(exception));
580 }
581 };
582 return getInnerPromise().then(js, kj::mv(onSuccess), kj::mv(onFailure));
583 }
584 }
585 JSG_CATCH(exception) {
586 return js.rejectedPromise<void>(kj::mv(exception));
587 }
588 }
589 
590 // If the algorithm does not exist, we just handle it as a success and move on.
591 onSuccess(js);
592 return js.resolvedPromise();
593}
594 
595jsg::Promise<void> maybeRunAlgorithmAsync(
596 jsg::Lock& js, auto& maybeAlgorithm, auto&& onSuccess, auto&& onFailure, auto&&... args) {
597 // The algorithm is a JavaScript function mapped through jsg::Function.
598 // It is expected to return a Promise mapped via jsg::Promise. If the
599 // function returns synchronously, the jsg::Promise wrapper ensures
600 // that it is properly mapped to a jsg::Promise, but if the Promise
601 // throws synchronously, we have to convert that synchronous throw
602 // into a proper rejected jsg::Promise.
603 KJ_IF_SOME(algorithm, maybeAlgorithm) {
604 // We need two layers of tryCatch here, unfortunately. The inner layer
605 // covers the algorithm implementation itself and is our typical error
606 // handling path. It ensures that if the algorithm throws an exception,
607 // that is properly converted in to a rejected promise that is *then*
608 // handled by the onFailure handler that is passed in. The outer tryCatch
609 // handles the rare and generally unexpected failure of the calls to
610 // .then() itself, which can throw JS exceptions synchronously in certain
611 // rare cases. For those we return a rejected promise but do not call the
612 // onFailure case since such errors are generally indicative of a fatal
613 // condition in the isolate (e.g. out of memory, other fatal exception, etc).
614 return js.tryCatch([&] {
615 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
616 return js
617 .tryCatch([&] { return algorithm(js, kj::fwd<decltype(args)>(args)...); },
618 [&](jsg::Value&& exception) { return js.rejectedPromise<void>(kj::mv(exception)); })
619 .then(js, ioContext.addFunctor(kj::mv(onSuccess)),
620 ioContext.addFunctor(kj::mv(onFailure)));
621 } else {
622 return js
623 .tryCatch([&] { return algorithm(js, kj::fwd<decltype(args)>(args)...); },
624 [&](jsg::Value&& exception) {
625 return js.rejectedPromise<void>(kj::mv(exception));
626 }).then(js, kj::mv(onSuccess), kj::mv(onFailure));
627 }
628 }, [&](jsg::Value&& exception) { return js.rejectedPromise<void>(kj::mv(exception)); });
629 }
630 
631 // If the algorithm does not exist, we handle it as a success but ensure
632 // it runs asynchronously by scheduling via a resolved promise.
633 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
634 return js.resolvedPromise().then(js, ioContext.addFunctor(kj::mv(onSuccess)));
635 } else {
636 return js.resolvedPromise().then(js, kj::mv(onSuccess));
637 }
638}
639 
640int getHighWaterMark(
641 const UnderlyingSource& underlyingSource, const StreamQueuingStrategy& queuingStrategy) {
642 bool isBytes = underlyingSource.type.map([](auto& s) { return s == "bytes"; }).orDefault(false);
643 return queuingStrategy.highWaterMark.orDefault(isBytes ? 0 : 1);
644}
645 
646} // namespace
647 
648// It is possible for the controller state to be released synchronously while
649// we are in the middle of a read. When that happens we need to defer the actual
650// close/error state change until the read call is complete. deferControllerStateChange
651// handles this for us by using the state machine's operation tracking to defer
652// pending close/error transitions until the read is complete.
653template <typename Controller>
654jsg::Promise<ReadResult> deferControllerStateChange(jsg::Lock& js,
655 Controller& controller,
656 kj::FunctionParam<jsg::Promise<ReadResult>()> readCallback) {
657 bool endOperation = true;
658 // The readCallback and the controller.doClose(..) and controller.doError(...)
659 // methods, as well as the methods can trigger JavaScript errors to be thrown
660 // synchronously in some cases. We want to make sure non-fatal errors cause the
661 // stream to error and only fatal cases bubble up.
662 return js.tryCatch([&] {
663 controller.state.beginOperation();
664 auto result = readCallback();
665 endOperation = false;
666 
667 // endOperation() will automatically apply any pending state if this was the last operation.
668 // Returns true if a pending state was applied.
669 if (controller.state.endOperation()) {
670 // A pending state was applied. Call the appropriate callback.
671 // Skip callbacks if execution is being terminated (e.g., CPU time limit) since we can't
672 // safely execute JavaScript in that state.
673 if (!js.v8Isolate->IsExecutionTerminating()) {
674 if (controller.state.template is<StreamStates::Closed>()) {
675 controller.lock.onClose(js);
676 } else if (controller.state.template is<StreamStates::Errored>()) {
677 KJ_IF_SOME(err, controller.state.template tryGetUnsafe<StreamStates::Errored>()) {
678 controller.lock.onError(js, err.getHandle(js));
679 }
680 }
681 }
682 }
683 
684 return kj::mv(result);
685 }, [&](jsg::Value exception) -> jsg::Promise<ReadResult> {
686 if (endOperation) {
687 // Clear any pending state since we're erroring
688 controller.state.clearPendingState();
689 (void)controller.state.endOperation();
690 }
691 controller.doError(js, exception.getHandle(js));
692 return js.rejectedPromise<ReadResult>(kj::mv(exception));
693 });
694}
695 
696// The ReadableStreamJsController provides the implementation of custom
697// ReadableStreams backed by a user-code provided Underlying Source. The implementation
698// is fairly complicated and defined entirely by the streams specification.
699//
700// Another important thing to understand is that there are two types of JavaScript
701// backed ReadableStreams: value-oriented, and byte-oriented.
702//
703// When user code uses the `new ReadableStream(underlyingSource)` constructor, the
704// underlyingSource argument may have a `type` property, the value of which is either
705// `undefined`, the empty string, or the string value `'bytes'`. If the underlyingSource
706// argument is not given, the default value of `type` is `undefined`. If `type` is
707// `undefined` or the empty string, the ReadableStream is value-oriented. If `type` is
708// exactly equal to `'bytes'`, the ReadableStream is byte-oriented.
709//
710// For value-oriented streams, any JavaScript value can be pushed through the stream,
711// and the stream will only support use of the ReadableStreamDefaultReader to consume
712// the stream data.
713//
714// For byte-oriented streams, only byte data (as provided by `ArrayBufferView`s) can
715// be pushed through the stream. All byte-oriented streams support using both
716// ReadableStreamDefaultReader and ReadableStreamBYOBReader to consume the stream
717// data.
718//
719// When the ReadableStreamJsController::setup() method is called the type
720// of stream is determined, and the controller will create an instance of either
721// jsg::Ref<ReadableStreamDefaultController> or jsg::Ref<ReadableByteStreamController>.
722// These are the objects that are actually passed on to the user-code's Underlying Source
723// implementation.
724class ReadableStreamJsController final: public ReadableStreamController {
725 public:
726 using ReadableLockImpl = ReadableLockImpl<ReadableStreamJsController>;
727 
728 KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamJsController);
729 
730 explicit ReadableStreamJsController();
731 explicit ReadableStreamJsController(StreamStates::Closed closed);
732 explicit ReadableStreamJsController(StreamStates::Errored errored);
733 explicit ReadableStreamJsController(jsg::Lock& js, ValueReadable& consumer);
734 explicit ReadableStreamJsController(jsg::Lock& js, ByteReadable& consumer);
735 
736 jsg::Ref<ReadableStream> addRef() override;
737 
738 void setup(jsg::Lock& js,
739 jsg::Optional<UnderlyingSource> maybeUnderlyingSource,
740 jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) override;
741 
742 // Signals that this ReadableStream is no longer interested in the underlying
743 // data source. Whether this cancels the underlying data source also depends
744 // on whether or not there are other ReadableStreams still attached to it.
745 // This operation is terminal. Once called, even while the returned Promise
746 // is still pending, the ReadableStream will be no longer usable and any
747 // data still in the queue will be dropped. Pending read requests will be
748 // rejected if a reason is given, or resolved with no data otherwise.
749 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) override;
750 
751 void doClose(jsg::Lock& js);
752 
753 void doError(jsg::Lock& js, v8::Local<v8::Value> reason);
754 
755 bool canCloseOrEnqueue();
756 bool hasBackpressure();
757 
758 bool isByteOriented() const override;
759 
760 bool isDisturbed() override;
761 
762 bool isClosedOrErrored() const override;
763 
764 bool isClosed() const override;
765 
766 bool isLockedToReader() const override;
767 
768 bool lockReader(jsg::Lock& js, Reader& reader) override;
769 
770 kj::Maybe<v8::Local<v8::Value>> isErrored(jsg::Lock& js);
771 
772 kj::Maybe<int> getDesiredSize();
773 
774 jsg::Promise<void> pipeTo(
775 jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) override;
776 
777 kj::Promise<DeferredProxy<void>> pumpTo(
778 jsg::Lock& js, kj::Own<WritableStreamSink>, bool end) override;
779 
780 kj::Maybe<jsg::Promise<ReadResult>> read(
781 jsg::Lock& js, kj::Maybe<ByobOptions> byobOptions) override;
782 
783 kj::Maybe<jsg::Promise<DrainingReadResult>> drainingRead(
784 jsg::Lock& js, size_t maxRead = kj::maxValue) override;
785 
786 // See the comment for releaseReader in common.h for details on the use of maybeJs
787 void releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) override;
788 
789 void setOwnerRef(ReadableStream& stream) override;
790 
791 Tee tee(jsg::Lock& js) override;
792 
793 kj::Maybe<PipeController&> tryPipeLock() override;
794 
795 void visitForGc(jsg::GcVisitor& visitor) override;
796 
797 kj::Maybe<kj::OneOf<DefaultController, ByobController>> getController();
798 
799 jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit) override;
800 jsg::Promise<kj::String> readAllText(jsg::Lock& js, uint64_t limit) override;
801 
802 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override;
803 
804 kj::Own<ReadableStreamController> detach(jsg::Lock& js, bool ignoreDisturbed) override;
805 
806 void setPendingClosure() override {
807 KJ_UNIMPLEMENTED("only implemented for WritableStreamInternalController");
808 }
809 
810 kj::StringPtr jsgGetMemoryName() const override;
811 size_t jsgGetMemorySelfSize() const override;
812 void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const override;
813 
814 private:
815 // If the stream was created within the scope of a request, we want to treat it as I/O
816 // and make sure it is not advanced from the scope of a different request.
817 kj::Maybe<IoContext&> ioContext;
818 kj::Maybe<ReadableStream&> owner;
819 
820 // Initial state before setup() is called.
821 struct Initial {
822 static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj;
823 };
824 
825 // State machine for ReadableStreamJsController:
826 // Initial is the default state before setup() is called
827 // ValueReadable and ByteReadable are the active states (stream has data)
828 // Closed and Errored are terminal states (stream is done)
829 // Initial -> ValueReadable or ByteReadable (setup() called)
830 // Initial -> Closed (constructed with Closed)
831 // Initial -> Errored (constructed with Errored)
832 // ValueReadable -> Closed (doClose() or cancel() called)
833 // ValueReadable -> Errored (doError() called)
834 // ByteReadable -> Closed (doClose() or cancel() called)
835 // ByteReadable -> Errored (doError() called)
836 // Note: No single ActiveState since there are two active variants.
837 // PendingStates allows Closed/Errored transitions to be deferred during reads.
838 using State = StateMachine<TerminalStates<StreamStates::Closed>,
839 ErrorState<StreamStates::Errored>,
840 PendingStates<StreamStates::Closed, StreamStates::Errored>,
841 Initial,
842 StreamStates::Closed,
843 StreamStates::Errored,
844 kj::Own<ValueReadable>,
845 kj::Own<ByteReadable>>;
846 State state = State::create<Initial>();
847 
848 kj::Maybe<uint64_t> expectedLength = kj::none;
849 bool canceling = false;
850 
851 // The lock state is separate because a closed or errored stream can still be locked.
852 ReadableLockImpl lock;
853 
854 bool disturbed = false;
855 
856 template <typename T>
857 jsg::Promise<T> readAll(jsg::Lock& js, uint64_t limit);
858 
859 friend ReadableLockImpl;
860 friend ReadableLockImpl::PipeLocked;
861 friend struct ValueReadable;
862 friend struct ByteReadable;
863 
864 template <typename Controller>
865 friend jsg::Promise<ReadResult> deferControllerStateChange(jsg::Lock& js,
866 Controller& controller,
867 kj::FunctionParam<jsg::Promise<ReadResult>()> readCallback);
868};
869 
870// The WritableStreamJsController provides the implementation of custom
871// WritableStream's backed by a user-code provided Underlying Sink. The implementation
872// is fairly complicated and defined entirely by the streams specification.
873class WritableStreamJsController final: public WritableStreamController {
874 public:
875 using WritableLockImpl = WritableLockImpl<WritableStreamJsController>;
876 
877 using Controller = jsg::Ref<WritableStreamDefaultController>;
878 
879 explicit WritableStreamJsController();
880 
881 explicit WritableStreamJsController(StreamStates::Closed closed);
882 
883 explicit WritableStreamJsController(StreamStates::Errored errored);
884 
885 ~WritableStreamJsController() noexcept(false);
886 
887 KJ_DISALLOW_COPY_AND_MOVE(WritableStreamJsController);
888 
889 jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) override;
890 
891 jsg::Ref<WritableStream> addRef() override;
892 
893 jsg::Promise<void> close(jsg::Lock& js, bool markAsHandled = false) override;
894 
895 jsg::Promise<void> flush(jsg::Lock& js, bool markAsHandled = false) override {
896 KJ_UNIMPLEMENTED("expected WritableStreamInternalController implementation to be enough");
897 }
898 
899 void doClose(jsg::Lock& js);
900 
901 void doError(jsg::Lock& js, v8::Local<v8::Value> reason);
902 
903 // Error through the underlying controller if available, going through the proper
904 // error transition (Erroring -> Errored).
905 void errorIfNeeded(jsg::Lock& js, v8::Local<v8::Value> reason);
906 
907 kj::Maybe<int> getDesiredSize() override;
908 
909 kj::Maybe<v8::Local<v8::Value>> isErroring(jsg::Lock& js) override;
910 kj::Maybe<v8::Local<v8::Value>> isErroredOrErroring(jsg::Lock& js);
911 
912 bool isLocked() const;
913 
914 bool isLockedToWriter() const override;
915 
916 bool isStarted();
917 
918 bool hasBackpressure();
919 
920 inline bool isWritable() const {
921 return state.isActive();
922 }
923 
924 bool lockWriter(jsg::Lock& js, Writer& writer) override;
925 
926 void maybeRejectReadyPromise(jsg::Lock& js, v8::Local<v8::Value> reason);
927 
928 void maybeResolveReadyPromise(jsg::Lock& js);
929 
930 // See the comment for releaseWriter in common.h for details on the use of maybeJs
931 void releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) override;
932 
933 kj::Maybe<kj::Own<WritableStreamSink>> removeSink(jsg::Lock& js) override;
934 void detach(jsg::Lock& js) override;
935 
936 void setOwnerRef(WritableStream& stream) override;
937 
938 void setup(jsg::Lock& js,
939 jsg::Optional<UnderlyingSink> maybeUnderlyingSink,
940 jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) override;
941 
942 kj::Maybe<jsg::Promise<void>> tryPipeFrom(
943 jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) override;
944 
945 void updateBackpressure(jsg::Lock& js, bool backpressure);
946 
947 jsg::Promise<void> write(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) override;
948 
949 void visitForGc(jsg::GcVisitor& visitor) override;
950 
951 bool isClosedOrClosing() override;
952 bool isErrored() override;
953 
954 inline bool isByteOriented() const override {
955 return false;
956 }
957 
958 void setPendingClosure() override {
959 KJ_UNIMPLEMENTED("only implemented for WritableStreamInternalController");
960 }
961 
962 kj::StringPtr jsgGetMemoryName() const override;
963 size_t jsgGetMemorySelfSize() const override;
964 void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override;
965 
966 private:
967 jsg::Promise<void> pipeLoop(jsg::Lock& js);
968 
969 kj::Maybe<IoContext&> ioContext;
970 kj::Maybe<WritableStream&> owner;
971 
972 // Initial state before setup() is called.
973 struct Initial {
974 static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj;
975 };
976 
977 // State machine for WritableStreamJsController:
978 // Initial is the default state before setup() is called
979 // Controller is the active state (stream is writable)
980 // Closed is terminal, Errored is implicitly terminal via ErrorState
981 using State = StateMachine<TerminalStates<StreamStates::Closed>,
982 ErrorState<StreamStates::Errored>,
983 ActiveState<Controller>,
984 Initial,
985 StreamStates::Closed,
986 StreamStates::Errored,
987 Controller>;
988 State state = State::create<Initial>();
989 
990 WritableLockImpl lock;
991 kj::Maybe<jsg::Promise<void>> maybeAbortPromise;
992 
993 friend WritableLockImpl;
994};
995 
996kj::Own<ReadableStreamController> newReadableStreamJsController() {
997 return kj::heap<ReadableStreamJsController>();
998}
999 
1000kj::Own<WritableStreamController> newWritableStreamJsController() {
1001 return kj::heap<WritableStreamJsController>();
1002}
1003 
1004template <typename Self>
1005ReadableImpl<Self>::ReadableImpl(
1006 UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy)
1007 : state(State::template create<Queue>(getHighWaterMark(underlyingSource, queuingStrategy))),
1008 algorithms(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {}
1009 
1010template <typename Self>
1011void ReadableImpl<Self>::start(jsg::Lock& js, jsg::Ref<Self> self) {
1012 KJ_ASSERT(!flags.started && !flags.starting);
1013 flags.starting = true;
1014 
1015 // Per the streams spec, the size function should be called with `undefined` as `this`,
1016 // not as a method on the strategy object.
1017 KJ_IF_SOME(sizeFunc, algorithms.size) {
1018 sizeFunc.setReceiver(jsg::Value(js.v8Isolate, js.v8Undefined()));
1019 }
1020 
1021 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) {
1022 flags.started = true;
1023 flags.starting = false;
1024 pullIfNeeded(js, kj::mv(self));
1025 });
1026 
1027 auto onFailure = JSG_VISITABLE_LAMBDA(
1028 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1029 flags.started = true;
1030 flags.starting = false;
1031 doError(js, kj::mv(reason));
1032 });
1033 
1034 maybeRunAlgorithm(js, algorithms.start, kj::mv(onSuccess), kj::mv(onFailure), kj::mv(self));
1035 algorithms.start = kj::none;
1036}
1037 
1038template <typename Self>
1039size_t ReadableImpl<Self>::consumerCount() {
1040 return state.whenActiveOr([](Queue& q) { return q.getConsumerCount(); }, size_t{0});
1041}
1042 
1043template <typename Self>
1044jsg::Promise<void> ReadableImpl<Self>::cancel(
1045 jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) {
1046 if (state.template is<StreamStates::Closed>()) {
1047 // We are already closed. There's nothing to cancel.
1048 // This shouldn't happen but we handle the case anyway, just to be safe.
1049 return js.resolvedPromise();
1050 }
1051 KJ_IF_SOME(errored, state.template tryGetUnsafe<StreamStates::Errored>()) {
1052 // We are already errored. There's nothing to cancel.
1053 // This shouldn't happen but we handle the case anyway, just to be safe.
1054 return js.rejectedPromise<void>(errored.getHandle(js));
1055 }
1056 
1057 auto& queue = state.template getUnsafe<Queue>();
1058 size_t consumerCount = queue.getConsumerCount();
1059 if (consumerCount > 1) {
1060 // If there is more than 1 consumer, then we just return here with an
1061 // immediately resolved promise. The consumer will remove itself,
1062 // canceling its interest in the underlying source but we do not yet
1063 // want to cancel the underlying source since there are still other
1064 // consumers that want data.
1065 return js.resolvedPromise();
1066 }
1067 
1068 // Otherwise, there should be exactly one consumer at this point.
1069 KJ_ASSERT(consumerCount == 1);
1070 KJ_IF_SOME(pendingCancel, maybePendingCancel) {
1071 // If we're already waiting for cancel to complete, just return the
1072 // already existing pending promise.
1073 // This shouldn't happen but we handle the case anyway, just to be safe.
1074 return pendingCancel.promise.whenResolved(js);
1075 }
1076 
1077 auto prp = js.newPromiseAndResolver<void>();
1078 maybePendingCancel = PendingCancel{
1079 .fulfiller = kj::mv(prp.resolver),
1080 .promise = kj::mv(prp.promise),
1081 };
1082 auto promise = KJ_ASSERT_NONNULL(maybePendingCancel).promise.whenResolved(js);
1083 doCancel(js, kj::mv(self), reason);
1084 return kj::mv(promise);
1085}
1086 
1087template <typename Self>
1088bool ReadableImpl<Self>::canCloseOrEnqueue() {
1089 return state.isActive();
1090}
1091 
1092// doCancel() is triggered by cancel() being called, which is an explicit signal from
1093// the ReadableStream that we don't care about the data this controller provides any
1094// more. We don't need to notify the consumers because we presume they already know
1095// that they called cancel. What we do want to do here, tho, is close the implementation
1096// and trigger the cancel algorithm.
1097template <typename Self>
1098void ReadableImpl<Self>::doCancel(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) {
1099 state.template transitionTo<StreamStates::Closed>();
1100 
1101 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) {
1102 doClose(js);
1103 KJ_IF_SOME(pendingCancel, maybePendingCancel) {
1104 maybeResolvePromise(js, pendingCancel.fulfiller);
1105 } else {
1106 // Else block to avert dangling else compiler warning.
1107 }
1108 });
1109 auto onFailure = JSG_VISITABLE_LAMBDA(
1110 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1111 // We do not call doError() here because there's really no point. Everything
1112 // that cares about the state of this controller impl has signaled that it
1113 // no longer cares and has gone away.
1114 doClose(js);
1115 KJ_IF_SOME(pendingCancel, maybePendingCancel) {
1116 maybeRejectPromise<void>(js, pendingCancel.fulfiller, reason.getHandle(js));
1117 } else {
1118 // Else block to avert dangling else compiler warning.
1119 }
1120 });
1121 
1122 maybeRunAlgorithm(js, algorithms.cancel, kj::mv(onSuccess), kj::mv(onFailure), reason);
1123}
1124 
1125template <typename Self>
1126void ReadableImpl<Self>::enqueue(jsg::Lock& js, kj::Rc<Entry> entry, jsg::Ref<Self> self) {
1127 JSG_REQUIRE(canCloseOrEnqueue(), TypeError, "This ReadableStream is closed.");
1128 KJ_DEFER(pullIfNeeded(js, kj::mv(self)));
1129 auto& queue = state.template getUnsafe<Queue>();
1130 queue.push(js, kj::mv(entry));
1131}
1132 
1133template <typename Self>
1134void ReadableImpl<Self>::close(jsg::Lock& js) {
1135 JSG_REQUIRE(canCloseOrEnqueue(), TypeError, "This ReadableStream is closed.");
1136 auto& queue = state.template getUnsafe<Queue>();
1137 
1138 if (queue.hasPartiallyFulfilledRead()) {
1139 auto error =
1140 js.v8Ref(js.v8TypeError("This ReadableStream was closed with a partial read pending."));
1141 doError(js, error.addRef(js));
1142 js.throwException(kj::mv(error));
1143 return;
1144 }
1145 
1146 queue.close(js);
1147 
1148 state.template transitionTo<StreamStates::Closed>();
1149 doClose(js);
1150}
1151 
1152template <typename Self>
1153void ReadableImpl<Self>::doClose(jsg::Lock& js) {
1154 // The state should have already been set to closed.
1155 KJ_ASSERT(state.template is<StreamStates::Closed>());
1156 algorithms.clear();
1157}
1158 
1159template <typename Self>
1160void ReadableImpl<Self>::doError(jsg::Lock& js, jsg::Value reason) {
1161 // If already closed or errored, do nothing
1162 if (state.isInactive()) {
1163 return;
1164 }
1165 
1166 auto& queue = state.template getUnsafe<Queue>();
1167 queue.error(js, reason.addRef(js));
1168 state.template transitionTo<StreamStates::Errored>(kj::mv(reason));
1169 algorithms.clear();
1170}
1171 
1172template <typename Self>
1173kj::Maybe<int> ReadableImpl<Self>::getDesiredSize() {
1174 if (state.template is<StreamStates::Closed>()) {
1175 return 0;
1176 }
1177 if (state.template is<StreamStates::Errored>()) {
1178 return kj::none;
1179 }
1180 return state.template getUnsafe<Queue>().desiredSize();
1181}
1182 
1183// We should call pull if any of the consumers known to the queue have read requests or
1184// we haven't yet signalled backpressure.
1185template <typename Self>
1186bool ReadableImpl<Self>::shouldCallPull() {
1187 return state.whenActiveOr(
1188 [this](Queue& q) { return q.wantsRead() || getDesiredSize().orDefault(0) > 0; }, false);
1189}
1190 
1191template <typename Self>
1192void ReadableImpl<Self>::pullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) {
1193 // Determining if we need to pull is fairly complicated. All of the following
1194 // must hold true:
1195 if (!shouldCallPull()) {
1196 return;
1197 }
1198 
1199 if (flags.pulling) {
1200 flags.pullAgain = true;
1201 return;
1202 }
1203 KJ_ASSERT(!flags.pullAgain);
1204 flags.pulling = true;
1205 
1206 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) {
1207 flags.pulling = false;
1208 if (flags.pullAgain) {
1209 flags.pullAgain = false;
1210 pullIfNeeded(js, kj::mv(self));
1211 }
1212 });
1213 
1214 auto onFailure = JSG_VISITABLE_LAMBDA(
1215 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1216 flags.pulling = false;
1217 doError(js, kj::mv(reason));
1218 });
1219 
1220 maybeRunAlgorithm(js, algorithms.pull, kj::mv(onSuccess), kj::mv(onFailure), self.addRef());
1221}
1222 
1223template <typename Self>
1224void ReadableImpl<Self>::forcePullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) {
1225 // Like pullIfNeeded but bypasses the shouldCallPull() check. Used for draining reads
1226 // which need to pull all available data regardless of backpressure settings.
1227 if (!canCloseOrEnqueue()) {
1228 return;
1229 }
1230 
1231 if (flags.pulling) {
1232 flags.pullAgain = true;
1233 return;
1234 }
1235 KJ_ASSERT(!flags.pullAgain);
1236 flags.pulling = true;
1237 
1238 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) {
1239 flags.pulling = false;
1240 if (flags.pullAgain) {
1241 flags.pullAgain = false;
1242 // After a force pull, we go back to normal pullIfNeeded behavior.
1243 pullIfNeeded(js, kj::mv(self));
1244 }
1245 });
1246 
1247 auto onFailure = JSG_VISITABLE_LAMBDA(
1248 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1249 flags.pulling = false;
1250 doError(js, kj::mv(reason));
1251 });
1252 
1253 maybeRunAlgorithm(js, algorithms.pull, kj::mv(onSuccess), kj::mv(onFailure), self.addRef());
1254}
1255 
1256template <typename Self>
1257void ReadableImpl<Self>::visitForGc(jsg::GcVisitor& visitor) {
1258 state.visitForGc(visitor);
1259 KJ_IF_SOME(pendingCancel, maybePendingCancel) {
1260 visitor.visit(pendingCancel.fulfiller, pendingCancel.promise);
1261 }
1262 visitor.visit(algorithms);
1263}
1264 
1265template <typename Self>
1266kj::Own<typename ReadableImpl<Self>::Consumer> ReadableImpl<Self>::getConsumer(
1267 kj::Maybe<ReadableImpl<Self>::StateListener&> listener) {
1268 auto& queue = state.template getUnsafe<Queue>();
1269 return kj::heap<typename ReadableImpl<Self>::Consumer>(queue, listener);
1270}
1271 
1272// ======================================================================================
1273 
1274template <typename Self>
1275WritableImpl<Self>::WritableImpl(
1276 jsg::Lock& js, WritableStream& owner, jsg::Ref<AbortSignal> abortSignal)
1277 : owner(owner.addWeakRef()),
1278 signal(kj::mv(abortSignal)) {
1279 flags.pedanticWpt = FeatureFlags::get(js).getPedanticWpt();
1280}
1281 
1282template <typename Self>
1283jsg::Promise<void> WritableImpl<Self>::abort(
1284 jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) {
1285 // Per the spec, the signal.reason should be a DOMException with name 'AbortError'
1286 // when no reason is provided, but the stored error should remain as the original reason.
1287 auto signalReason = [&]() -> jsg::JsValue {
1288 if (reason->IsUndefined() && FeatureFlags::get(js).getPedanticWpt()) {
1289 auto ex = js.domException(
1290 kj::str("AbortError"), kj::str("This writable stream has been aborted."), kj::none);
1291 return jsg::JsValue(KJ_ASSERT_NONNULL(ex.tryGetHandle(js)));
1292 }
1293 return jsg::JsValue(reason);
1294 }();
1295 signal->triggerAbort(js, signalReason);
1296 
1297 // We have to check this again after the AbortSignal is triggered.
1298 if (state.isTerminal()) {
1299 return js.resolvedPromise();
1300 }
1301 
1302 KJ_IF_SOME(pendingAbort, maybePendingAbort) {
1303 // Notice here that, per the spec, the reason given in this call of abort is
1304 // intentionally ignored if there is already an abort pending.
1305 return pendingAbort->whenResolved(js);
1306 }
1307 
1308 bool wasAlreadyErroring = false;
1309 if (state.template is<StreamStates::Erroring>()) {
1310 wasAlreadyErroring = true;
1311 reason = js.v8Undefined();
1312 }
1313 
1314 KJ_DEFER(if (!wasAlreadyErroring) { startErroring(js, kj::mv(self), reason); });
1315 
1316 maybePendingAbort = kj::heap<PendingAbort>(js, reason, wasAlreadyErroring);
1317 return KJ_ASSERT_NONNULL(maybePendingAbort)->whenResolved(js);
1318}
1319 
1320template <typename Self>
1321kj::Maybe<WritableStreamJsController&> WritableImpl<Self>::tryGetOwner() {
1322 KJ_IF_SOME(o, owner) {
1323 return o->tryGet().map([](WritableStream& owner) -> WritableStreamJsController& {
1324 return static_cast<WritableStreamJsController&>(owner.getController());
1325 });
1326 }
1327 return kj::none;
1328}
1329 
1330template <typename Self>
1331ssize_t WritableImpl<Self>::getDesiredSize() {
1332 return highWaterMark - amountBuffered;
1333}
1334 
1335template <typename Self>
1336void WritableImpl<Self>::advanceQueueIfNeeded(jsg::Lock& js, jsg::Ref<Self> self) {
1337 if (!flags.started || inFlightWrite != kj::none) {
1338 return;
1339 }
1340 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1341 
1342 if (state.template is<StreamStates::Erroring>()) {
1343 return finishErroring(js, kj::mv(self));
1344 }
1345 
1346 if (writeRequests.empty()) {
1347 if (closeRequest != kj::none) {
1348 KJ_ASSERT(inFlightClose == kj::none);
1349 KJ_ASSERT_NONNULL(closeRequest);
1350 inFlightClose = kj::mv(closeRequest);
1351 
1352 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self),
1353 (jsg::Lock& js) { finishInFlightClose(js, kj::mv(self)); });
1354 
1355 auto onFailure = JSG_VISITABLE_LAMBDA(
1356 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1357 finishInFlightClose(js, kj::mv(self), reason.getHandle(js));
1358 });
1359 
1360 // Per the spec, the close algorithm should always run asynchronously, even if
1361 // there's no user-provided close handler. This ensures that releaseLock() can
1362 // reject the closed promise before the close completes.
1363 // The original maybeRunAlgorithm would call the onSuccess continuation
1364 // synchronously if algorithms.close is not specified. maybeRunAlgorithmAsync
1365 // always defers to a microtask.
1366 if (FeatureFlags::get(js).getPedanticWpt()) {
1367 maybeRunAlgorithmAsync(js, algorithms.close, kj::mv(onSuccess), kj::mv(onFailure));
1368 } else {
1369 maybeRunAlgorithm(js, algorithms.close, kj::mv(onSuccess), kj::mv(onFailure));
1370 }
1371 }
1372 return;
1373 }
1374 
1375 KJ_ASSERT(inFlightWrite == kj::none);
1376 auto req = dequeueWriteRequest();
1377 auto value = req.value.addRef(js);
1378 auto size = req.size;
1379 inFlightWrite = kj::mv(req);
1380 
1381 auto onSuccess =
1382 JSG_VISITABLE_LAMBDA((this, self = self.addRef(), size), (self), (jsg::Lock& js) {
1383 amountBuffered -= size;
1384 finishInFlightWrite(js, self.addRef());
1385 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1386 if (!isCloseQueuedOrInFlight() && isWritable()) {
1387 updateBackpressure(js);
1388 }
1389 if (state.template is<StreamStates::Erroring>() || writeRequests.empty()) {
1390 // In this case, we know advanceQueueIfNeeded won't recurse further, so we can
1391 // avoid the extra microtask hop.
1392 advanceQueueIfNeeded(js, kj::mv(self));
1393 return js.resolvedPromise();
1394 }
1395 // Here, however, let's avoid potentially deep recursion by hopping to a new
1396 // microtask to continue processing the queue.
1397 return js.resolvedPromise().then(
1398 js, JSG_VISITABLE_LAMBDA((this, self = kj::mv(self)), (self), (jsg::Lock & js) mutable {
1399 if (isWritable() || state.template is<StreamStates::Erroring>()) {
1400 advanceQueueIfNeeded(js, kj::mv(self));
1401 }
1402 }));
1403 });
1404 
1405 auto onFailure = JSG_VISITABLE_LAMBDA(
1406 (this, self = self.addRef(), size), (self), (jsg::Lock& js, jsg::Value reason) {
1407 amountBuffered -= size;
1408 finishInFlightWrite(js, kj::mv(self), reason.getHandle(js));
1409 return js.resolvedPromise();
1410 });
1411 
1412 // Per the spec, the write algorithm should always run asynchronously, even if
1413 // there's no user-provided write handler. This ensures that backpressure changes
1414 // from the write don't resolve the ready promise synchronously, preserving correct
1415 // microtask ordering (e.g., ready rejects before closed on releaseLock).
1416 if (FeatureFlags::get(js).getPedanticWpt()) {
1417 maybeRunAlgorithmAsync(js, algorithms.write, kj::mv(onSuccess), kj::mv(onFailure),
1418 value.getHandle(js), self.addRef());
1419 } else {
1420 maybeRunAlgorithm(js, algorithms.write, kj::mv(onSuccess), kj::mv(onFailure),
1421 value.getHandle(js), self.addRef());
1422 }
1423}
1424 
1425template <typename Self>
1426jsg::Promise<void> WritableImpl<Self>::close(jsg::Lock& js, jsg::Ref<Self> self) {
1427 if (state.template is<StreamStates::Closed>()) {
1428 return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj));
1429 }
1430 KJ_IF_SOME(errored, state.template tryGetUnsafe<StreamStates::Errored>()) {
1431 return js.rejectedPromise<void>(errored.addRef(js));
1432 }
1433 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1434 JSG_REQUIRE(
1435 !isCloseQueuedOrInFlight(), TypeError, "Cannot close a writer that is already being closed");
1436 auto prp = js.newPromiseAndResolver<void>();
1437 closeRequest = kj::mv(prp.resolver);
1438 
1439 if (flags.backpressure && isWritable()) {
1440 KJ_IF_SOME(owner, tryGetOwner()) {
1441 owner.maybeResolveReadyPromise(js);
1442 }
1443 }
1444 
1445 advanceQueueIfNeeded(js, kj::mv(self));
1446 
1447 return kj::mv(prp.promise);
1448}
1449 
1450template <typename Self>
1451void WritableImpl<Self>::dealWithRejection(
1452 jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) {
1453 if (isWritable()) {
1454 return startErroring(js, kj::mv(self), reason);
1455 }
1456 KJ_ASSERT(state.template is<StreamStates::Erroring>());
1457 finishErroring(js, kj::mv(self));
1458}
1459 
1460template <typename Self>
1461WritableImpl<Self>::WriteRequest WritableImpl<Self>::dequeueWriteRequest() {
1462 auto write = kj::mv(writeRequests.front());
1463 writeRequests.pop_front();
1464 return kj::mv(write);
1465}
1466 
1467template <typename Self>
1468void WritableImpl<Self>::doClose(jsg::Lock& js) {
1469 KJ_ASSERT(closeRequest == kj::none);
1470 KJ_ASSERT(inFlightClose == kj::none);
1471 KJ_ASSERT(inFlightWrite == kj::none);
1472 KJ_ASSERT(maybePendingAbort == kj::none);
1473 KJ_ASSERT(writeRequests.empty());
1474 // State should have already been transitioned to Closed
1475 KJ_ASSERT(state.template is<StreamStates::Closed>());
1476 algorithms.clear();
1477 
1478 KJ_IF_SOME(owner, tryGetOwner()) {
1479 owner.doClose(js);
1480 }
1481}
1482 
1483template <typename Self>
1484void WritableImpl<Self>::doError(jsg::Lock& js, v8::Local<v8::Value> reason) {
1485 KJ_ASSERT(closeRequest == kj::none);
1486 KJ_ASSERT(inFlightClose == kj::none);
1487 KJ_ASSERT(inFlightWrite == kj::none);
1488 KJ_ASSERT(maybePendingAbort == kj::none);
1489 KJ_ASSERT(writeRequests.empty());
1490 // State should have already been transitioned to Errored
1491 KJ_ASSERT(state.template is<StreamStates::Errored>());
1492 algorithms.clear();
1493 
1494 KJ_IF_SOME(owner, tryGetOwner()) {
1495 owner.doError(js, reason);
1496 }
1497}
1498 
1499template <typename Self>
1500void WritableImpl<Self>::error(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) {
1501 if (isWritable()) {
1502 algorithms.clear();
1503 startErroring(js, kj::mv(self), reason);
1504 }
1505}
1506 
1507template <typename Self>
1508void WritableImpl<Self>::finishErroring(jsg::Lock& js, jsg::Ref<Self> self) {
1509 auto erroring = kj::mv(KJ_ASSERT_NONNULL(state.template tryGetUnsafe<StreamStates::Erroring>()));
1510 auto reason = erroring.reason.getHandle(js);
1511 KJ_ASSERT(inFlightWrite == kj::none);
1512 KJ_ASSERT(inFlightClose == kj::none);
1513 state.template transitionTo<StreamStates::Errored>(kj::mv(erroring.reason));
1514 
1515 while (!writeRequests.empty()) {
1516 dequeueWriteRequest().resolver.reject(js, reason);
1517 }
1518 KJ_ASSERT(writeRequests.empty());
1519 
1520 KJ_IF_SOME(pendingAbort, maybePendingAbort) {
1521 if (pendingAbort->reject) {
1522 pendingAbort->fail(js, reason);
1523 return rejectCloseAndClosedPromiseIfNeeded(js);
1524 }
1525 
1526 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) {
1527 auto& pendingAbort = KJ_ASSERT_NONNULL(maybePendingAbort);
1528 pendingAbort->reject = false;
1529 pendingAbort->complete(js);
1530 rejectCloseAndClosedPromiseIfNeeded(js);
1531 });
1532 
1533 auto onFailure = JSG_VISITABLE_LAMBDA(
1534 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1535 auto& pendingAbort = KJ_ASSERT_NONNULL(maybePendingAbort);
1536 pendingAbort->fail(js, reason.getHandle(js));
1537 rejectCloseAndClosedPromiseIfNeeded(js);
1538 });
1539 
1540 maybeRunAlgorithm(js, algorithms.abort, kj::mv(onSuccess), kj::mv(onFailure), reason);
1541 return;
1542 }
1543 rejectCloseAndClosedPromiseIfNeeded(js);
1544}
1545 
1546template <typename Self>
1547void WritableImpl<Self>::finishInFlightClose(
1548 jsg::Lock& js, jsg::Ref<Self> self, kj::Maybe<v8::Local<v8::Value>> maybeReason) {
1549 algorithms.clear();
1550 KJ_ASSERT_NONNULL(inFlightClose);
1551 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1552 
1553 KJ_IF_SOME(reason, maybeReason) {
1554 maybeRejectPromise<void>(js, inFlightClose, reason);
1555 
1556 KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) {
1557 pendingAbort->fail(js, reason);
1558 }
1559 
1560 return dealWithRejection(js, kj::mv(self), reason);
1561 }
1562 
1563 maybeResolvePromise(js, inFlightClose);
1564 
1565 if (state.template is<StreamStates::Erroring>()) {
1566 KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) {
1567 pendingAbort->reject = false;
1568 pendingAbort->complete(js);
1569 }
1570 }
1571 KJ_ASSERT(maybePendingAbort == kj::none);
1572 
1573 state.template transitionTo<StreamStates::Closed>();
1574 doClose(js);
1575}
1576 
1577template <typename Self>
1578void WritableImpl<Self>::finishInFlightWrite(
1579 jsg::Lock& js, jsg::Ref<Self> self, kj::Maybe<v8::Local<v8::Value>> maybeReason) {
1580 auto& write = KJ_ASSERT_NONNULL(inFlightWrite);
1581 
1582 KJ_IF_SOME(reason, maybeReason) {
1583 write.resolver.reject(js, reason);
1584 inFlightWrite = kj::none;
1585 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1586 return dealWithRejection(js, kj::mv(self), reason);
1587 }
1588 
1589 write.resolver.resolve(js);
1590 inFlightWrite = kj::none;
1591}
1592 
1593template <typename Self>
1594bool WritableImpl<Self>::isCloseQueuedOrInFlight() {
1595 return closeRequest != kj::none || inFlightClose != kj::none;
1596}
1597 
1598template <typename Self>
1599void WritableImpl<Self>::rejectCloseAndClosedPromiseIfNeeded(jsg::Lock& js) {
1600 algorithms.clear();
1601 auto reason =
1602 KJ_ASSERT_NONNULL(state.template tryGetUnsafe<StreamStates::Errored>()).getHandle(js);
1603 maybeRejectPromise<void>(js, closeRequest, reason);
1604 PendingAbort::dequeue(maybePendingAbort);
1605 doError(js, reason);
1606}
1607 
1608template <typename Self>
1609void WritableImpl<Self>::setup(jsg::Lock& js,
1610 jsg::Ref<Self> self,
1611 UnderlyingSink underlyingSink,
1612 StreamQueuingStrategy queuingStrategy) {
1613 KJ_ASSERT(!flags.started && !flags.starting);
1614 flags.starting = true;
1615 
1616 highWaterMark = queuingStrategy.highWaterMark.orDefault(1);
1617 auto startAlgorithm = kj::mv(underlyingSink.start);
1618 algorithms.write = kj::mv(underlyingSink.write);
1619 algorithms.close = kj::mv(underlyingSink.close);
1620 algorithms.abort = kj::mv(underlyingSink.abort);
1621 algorithms.size = kj::mv(queuingStrategy.size);
1622 // Per the streams spec, the size function should be called with `undefined` as `this`,
1623 // not as a method on the strategy object.
1624 KJ_IF_SOME(sizeFunc, algorithms.size) {
1625 sizeFunc.setReceiver(jsg::Value(js.v8Isolate, js.v8Undefined()));
1626 }
1627 
1628 auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) {
1629 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1630 
1631 if (isWritable()) {
1632 // Only resolve the ready promise if an abort is not pending.
1633 // It will have been rejected already.
1634 KJ_IF_SOME(owner, tryGetOwner()) {
1635 owner.maybeResolveReadyPromise(js);
1636 } else {
1637 // Else block to avert dangling else compiler warning.
1638 }
1639 }
1640 
1641 flags.started = true;
1642 flags.starting = false;
1643 advanceQueueIfNeeded(js, kj::mv(self));
1644 });
1645 
1646 auto onFailure = JSG_VISITABLE_LAMBDA(
1647 (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) {
1648 auto handle = reason.getHandle(js);
1649 KJ_ASSERT(isWritable() || state.template is<StreamStates::Erroring>());
1650 KJ_IF_SOME(owner, tryGetOwner()) {
1651 owner.maybeRejectReadyPromise(js, handle);
1652 } else {
1653 // Else block to avert dangling else compiler warning.
1654 }
1655 flags.started = true;
1656 flags.starting = false;
1657 dealWithRejection(js, kj::mv(self), handle);
1658 });
1659 
1660 flags.backpressure = getDesiredSize() <= 0;
1661 
1662 maybeRunAlgorithm(js, startAlgorithm, kj::mv(onSuccess), kj::mv(onFailure), self.addRef());
1663}
1664 
1665template <typename Self>
1666void WritableImpl<Self>::startErroring(
1667 jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason) {
1668 KJ_ASSERT(isWritable());
1669 KJ_IF_SOME(owner, tryGetOwner()) {
1670 owner.maybeRejectReadyPromise(js, reason);
1671 }
1672 state.template transitionTo<StreamStates::Erroring>(js.v8Ref(reason));
1673 if (inFlightWrite == kj::none && inFlightClose == kj::none && flags.started) {
1674 finishErroring(js, kj::mv(self));
1675 }
1676}
1677 
1678template <typename Self>
1679void WritableImpl<Self>::updateBackpressure(jsg::Lock& js) {
1680 KJ_ASSERT(isWritable());
1681 KJ_ASSERT(!isCloseQueuedOrInFlight());
1682 bool bp = getDesiredSize() <= 0;
1683 
1684 if (bp != flags.backpressure) {
1685 flags.backpressure = bp;
1686 KJ_IF_SOME(owner, tryGetOwner()) {
1687 owner.updateBackpressure(js, flags.backpressure);
1688 }
1689 }
1690}
1691 
1692template <typename Self>
1693jsg::Promise<void> WritableImpl<Self>::write(
1694 jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> value) {
1695 
1696 size_t size = 1;
1697 KJ_IF_SOME(sizeFunc, algorithms.size) {
1698 kj::Maybe<jsg::Value> failure;
1699 JSG_TRY(js) {
1700 size = sizeFunc(js, value);
1701 }
1702 JSG_CATCH(exception) {
1703 startErroring(js, self.addRef(), exception.getHandle(js));
1704 failure = kj::mv(exception);
1705 }
1706 KJ_IF_SOME(exception, failure) {
1707 return js.rejectedPromise<void>(kj::mv(exception));
1708 }
1709 }
1710 
1711 // Per spec (WritableStreamDefaultWriterWrite step 5), after calling the size
1712 // algorithm, re-check that the stream is still locked to a writer. If
1713 // releaseLock() was called from within strategy.size(), the write must be
1714 // rejected. This check must occur before any state checks, as the stream
1715 // state may still appear writable even after the writer was released.
1716 if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) {
1717 KJ_IF_SOME(owner, tryGetOwner()) {
1718 if (!owner.isLockedToWriter()) {
1719 return js.rejectedPromise<void>(
1720 js.v8TypeError("This WritableStream writer has been released."_kjc));
1721 }
1722 }
1723 }
1724 
1725 KJ_IF_SOME(error, state.template tryGetUnsafe<StreamStates::Errored>()) {
1726 return js.rejectedPromise<void>(error.addRef(js));
1727 }
1728 
1729 if (isCloseQueuedOrInFlight() || state.template is<StreamStates::Closed>()) {
1730 return js.rejectedPromise<void>(js.v8TypeError("This ReadableStream is closed."_kj));
1731 }
1732 
1733 KJ_IF_SOME(erroring, state.template tryGetUnsafe<StreamStates::Erroring>()) {
1734 return js.rejectedPromise<void>(erroring.reason.addRef(js));
1735 }
1736 
1737 KJ_ASSERT(isWritable());
1738 
1739 auto prp = js.newPromiseAndResolver<void>();
1740 writeRequests.push_back(WriteRequest{
1741 .resolver = kj::mv(prp.resolver),
1742 .value = js.v8Ref(value),
1743 .size = size,
1744 });
1745 amountBuffered += size;
1746 
1747 updateBackpressure(js);
1748 advanceQueueIfNeeded(js, kj::mv(self));
1749 return kj::mv(prp.promise);
1750}
1751 
1752template <typename Self>
1753void WritableImpl<Self>::visitForGc(jsg::GcVisitor& visitor) {
1754 state.visitForGc(visitor);
1755 visitor.visit(inFlightWrite, inFlightClose, closeRequest, algorithms, signal);
1756 KJ_IF_SOME(pendingAbort, maybePendingAbort) {
1757 visitor.visit(*pendingAbort);
1758 }
1759 visitor.visitAll(writeRequests);
1760}
1761 
1762template <typename Self>
1763bool WritableImpl<Self>::isWritable() const {
1764 return state.isActive();
1765}
1766 
1767template <typename Self>
1768void WritableImpl<Self>::cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason) {
1769 for (auto& write: writeRequests) {
1770 write.resolver.reject(js, reason);
1771 }
1772 writeRequests.clear();
1773}
1774 
1775// ======================================================================================
1776 
1777namespace {
1778template <typename Controller, typename Queue>
1779struct ReadableState {
1780 Controller controller;
1781 kj::Own<typename Queue::Consumer> consumer;
1782 ReadableStreamJsController& owner;
1783 
1784 ReadableState(Controller controller,
1785 kj::Own<typename Queue::Consumer> consumer,
1786 ReadableStreamJsController& owner)
1787 : controller(kj::mv(controller)),
1788 consumer(kj::mv(consumer)),
1789 owner(owner) {}
1790 
1791 ReadableState(Controller controller,
1792 Queue::ConsumerImpl::StateListener& listener,
1793 ReadableStreamJsController& owner)
1794 : ReadableState(controller.addRef(), controller->getConsumer(listener), owner) {}
1795 
1796 ReadableState clone(jsg::Lock& js,
1797 Queue::ConsumerImpl::StateListener& listener,
1798 ReadableStreamJsController& owner) {
1799 return ReadableState(controller.addRef(), consumer->clone(js, listener), owner);
1800 }
1801};
1802 
1803struct ValueReadable final: private api::ValueQueue::ConsumerImpl::StateListener {
1804 
1805 using State = ReadableState<DefaultController, ValueQueue>;
1806 kj::Maybe<State> state;
1807 bool reading = false;
1808 bool pendingCancel = false;
1809 
1810 JSG_MEMORY_INFO(ValueReadable) {
1811 KJ_IF_SOME(s, state) {
1812 tracker.trackField("controller", s.controller);
1813 tracker.trackField("consumer", s.consumer);
1814 }
1815 }
1816 
1817 void visitForGc(jsg::GcVisitor& visitor) {
1818 KJ_IF_SOME(s, state) {
1819 visitor.visit(s.controller, *s.consumer);
1820 }
1821 }
1822 
1823 ValueReadable(DefaultController controller, ReadableStreamJsController& owner)
1824 : state(State(kj::mv(controller), *this, owner)) {}
1825 
1826 ValueReadable(jsg::Lock& js, ReadableStreamJsController& owner, ValueReadable& other)
1827 : state(KJ_ASSERT_NONNULL(other.state).clone(js, *this, owner)) {}
1828 
1829 KJ_DISALLOW_COPY_AND_MOVE(ValueReadable);
1830 
1831 void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) {
1832 KJ_IF_SOME(s, state) {
1833 s.consumer->cancelPendingReads(js, reason);
1834 }
1835 }
1836 
1837 kj::Own<ValueReadable> clone(jsg::Lock& js, ReadableStreamJsController& owner) {
1838 // A single ReadableStreamDefaultController can have multiple consumers.
1839 // When the ValueReadable constructor is used, the new consumer is added
1840 // and starts to receive new data that becomes enqueued. When clone
1841 // is used, any state currently held by this consumer is copied to the
1842 // new consumer.
1843 return kj::heap<ValueReadable>(js, owner, *this);
1844 }
1845 
1846 jsg::Promise<ReadResult> read(jsg::Lock& js) {
1847 KJ_IF_SOME(s, state) {
1848 auto prp = js.newPromiseAndResolver<ReadResult>();
1849 reading = true;
1850 s.consumer->read(js,
1851 ValueQueue::ReadRequest{
1852 .resolver = kj::mv(prp.resolver),
1853 });
1854 reading = false;
1855 if (pendingCancel) {
1856 // If we were canceled while reading, we need to drop our state now.
1857 state = kj::none;
1858 pendingCancel = false;
1859 }
1860 return kj::mv(prp.promise);
1861 }
1862 
1863 // We are canceled! There's nothing to do.
1864 return js.resolvedPromise(ReadResult{.done = true});
1865 }
1866 
1867 jsg::Promise<DrainingReadResult> drainingRead(jsg::Lock& js, size_t maxRead) {
1868 KJ_IF_SOME(s, state) {
1869 // Note: We do NOT call beginOperation()/endOperation() here. The caller
1870 // (ReadableStreamJsController::drainingRead) manages the operation scope
1871 // around both this call and the returned promise's lifetime. If we added
1872 // our own beginOperation/endOperation here, the endOperation would fire
1873 // before the caller's wrapDrainingRead could set up its .then() callbacks,
1874 // potentially destroying the Consumer while the returned promise still has
1875 // dangling this-capturing callbacks from consumer->drainingRead().
1876 return s.consumer->drainingRead(js, maxRead);
1877 }
1878 
1879 // We are canceled! Return done with empty chunks.
1880 return js.resolvedPromise(DrainingReadResult{
1881 .chunks = kj::Array<kj::Array<kj::byte>>(),
1882 .done = true,
1883 });
1884 }
1885 
1886 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
1887 // When a ReadableStream is canceled, the expected behavior is that the underlying
1888 // controller is notified and the cancel algorithm on the underlying source is
1889 // called. When there are multiple ReadableStreams sharing consumption of a
1890 // controller, however, it should act as a shared pointer of sorts, canceling
1891 // the underlying controller only when the last reader is canceled.
1892 // Here, we rely on the controller implementing the correct behavior since it owns
1893 // the queue that knows about all of the attached consumers.
1894 if (pendingCancel) return js.resolvedPromise();
1895 KJ_IF_SOME(s, state) {
1896 // Check if there's a pending draining read before calling cancel, since cancel
1897 // will resolve the pending read and we need to know if we should defer destruction.
1898 bool hasPendingDrainingRead = s.consumer->hasPendingDrainingRead();
1899 s.consumer->cancel(js, maybeReason);
1900 auto promise = s.controller->cancel(js, kj::mv(maybeReason));
1901 // If we're currently in a read (sync or draining), we need to wait for that to
1902 // finish before dropping our state. For draining reads, the promise callbacks
1903 // capture 'this' (the Consumer) to clear hasPendingDrainingRead. If we destroy
1904 // the state now, those callbacks will UAF.
1905 if (reading || hasPendingDrainingRead) {
1906 pendingCancel = true;
1907 } else {
1908 state = kj::none;
1909 }
1910 return kj::mv(promise);
1911 }
1912 
1913 return js.resolvedPromise();
1914 }
1915 
1916 void onConsumerClose(jsg::Lock& js) override {
1917 // Called by the consumer when a state change to closed happens.
1918 // We need to notify the owner. Note that the owner may drop this
1919 // readable in doClose so it is not safe to access anything on this
1920 // after calling doClose.
1921 KJ_IF_SOME(s, state) {
1922 s.owner.doClose(js);
1923 }
1924 }
1925 
1926 void onConsumerError(jsg::Lock& js, jsg::Value reason) override {
1927 // Called by the consumer when a state change to errored happens.
1928 // We need to notify the owner. Note that the owner may drop this
1929 // readable in doClose so it is not safe to access anything on this
1930 // after calling doError.
1931 KJ_IF_SOME(s, state) {
1932 s.owner.doError(js, reason.getHandle(js));
1933 }
1934 }
1935 
1936 bool onConsumerWantsData(jsg::Lock& js) override {
1937 // Called by the consumer when it has a queued pending read and needs
1938 // data to be provided to fulfill it. We need to notify the controller
1939 // to initiate pulling to provide the data.
1940 // Returns true if the pull completed synchronously (meaning more pumping
1941 // might yield additional synchronous data), false otherwise.
1942 KJ_IF_SOME(s, state) {
1943 // Save a reference to the owner before calling pull. The pull callback
1944 // may trigger close/error which could destroy this ValueReadable. By
1945 // using beginOperation(), we ensure doClose/doError defers the
1946 // actual destruction until after we return.
1947 ReadableStreamJsController& owner = s.owner;
1948 owner.state.beginOperation();
1949 
1950 // For draining reads, use forcePull to bypass backpressure checks.
1951 // This ensures we pull all available data regardless of highWaterMark.
1952 if (s.consumer->hasPendingDrainingRead()) {
1953 s.controller->forcePull(js);
1954 } else {
1955 s.controller->pull(js);
1956 }
1957 
1958 // Check if state is still valid BEFORE calling endOperation(),
1959 // because that call may destroy this ValueReadable if close was deferred.
1960 bool result =
1961 state.map([](State& s2) { return !s2.controller->isPulling(); }).orDefault(false);
1962 
1963 // Process any deferred close/error. This may destroy this ValueReadable.
1964 if (owner.state.endOperation()) {
1965 // A pending state was applied. Call the appropriate callback.
1966 if (owner.state.template is<StreamStates::Closed>()) {
1967 owner.lock.onClose(js);
1968 } else if (owner.state.template is<StreamStates::Errored>()) {
1969 KJ_IF_SOME(err, owner.state.template tryGetUnsafe<StreamStates::Errored>()) {
1970 owner.lock.onError(js, err.getHandle(js));
1971 }
1972 }
1973 }
1974 
1975 return result;
1976 }
1977 return false;
1978 }
1979 
1980 kj::Maybe<int> getDesiredSize() {
1981 KJ_IF_SOME(s, state) {
1982 return s.controller->getDesiredSize();
1983 }
1984 return kj::none;
1985 }
1986 
1987 bool canCloseOrEnqueue() {
1988 return state.map([](State& s) { return s.controller->canCloseOrEnqueue(); }).orDefault(false);
1989 }
1990 
1991 kj::Maybe<DefaultController> getControllerRef() {
1992 return state.map([](State& s) { return s.controller.addRef(); });
1993 }
1994};
1995 
1996struct ByteReadable final: private api::ByteQueue::ConsumerImpl::StateListener {
1997 
1998 using State = ReadableState<ByobController, ByteQueue>;
1999 kj::Maybe<State> state;
2000 kj::Maybe<int> autoAllocateChunkSize;
2001 bool pendingCancel = false;
2002 
2003 JSG_MEMORY_INFO(ByteReadable) {
2004 KJ_IF_SOME(s, state) {
2005 tracker.trackField("controller", s.controller);
2006 tracker.trackField("consumer", s.consumer);
2007 }
2008 }
2009 
2010 void visitForGc(jsg::GcVisitor& visitor) {
2011 KJ_IF_SOME(s, state) {
2012 visitor.visit(s.controller, *s.consumer);
2013 }
2014 }
2015 
2016 ByteReadable(ByobController controller,
2017 ReadableStreamJsController& owner,
2018 kj::Maybe<int> autoAllocateChunkSize)
2019 : state(State(kj::mv(controller), *this, owner)),
2020 autoAllocateChunkSize(autoAllocateChunkSize) {}
2021 
2022 ByteReadable(jsg::Lock& js, ReadableStreamJsController& owner, ByteReadable& other)
2023 : state(KJ_ASSERT_NONNULL(other.state).clone(js, *this, owner)),
2024 autoAllocateChunkSize(other.autoAllocateChunkSize) {}
2025 
2026 KJ_DISALLOW_COPY_AND_MOVE(ByteReadable);
2027 
2028 void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) {
2029 KJ_IF_SOME(s, state) {
2030 s.consumer->cancelPendingReads(js, reason);
2031 }
2032 }
2033 
2034 // A single ReadableByteStreamController can have multiple consumers.
2035 // When the ByteReadable constructor is used, the new consumer is added
2036 // and starts to receive new data that becomes enqueued. When clone
2037 // is used, any state currently held by this consumer is copied to the
2038 // new consumer.
2039 kj::Own<ByteReadable> clone(jsg::Lock& js, ReadableStreamJsController& owner) {
2040 return kj::heap<ByteReadable>(js, owner, *this);
2041 }
2042 
2043 jsg::Promise<ReadResult> read(
2044 jsg::Lock& js, kj::Maybe<ReadableStreamController::ByobOptions> byobOptions) {
2045 KJ_IF_SOME(s, state) {
2046 auto prp = js.newPromiseAndResolver<ReadResult>();
2047 
2048 KJ_IF_SOME(byob, byobOptions) {
2049 jsg::BufferSource source(js, byob.bufferView.getHandle(js));
2050 // If atLeast is not given, then by default it is the element size of the view
2051 // that we were given. If atLeast is given, we make sure that it is aligned
2052 // with the element size. No matter what, atLeast cannot be less than 1.
2053 auto atLeast = kj::max(source.getElementSize(), byob.atLeast.orDefault(1));
2054 atLeast = kj::max(1, atLeast - (atLeast % source.getElementSize()));
2055 s.consumer->read(js,
2056 ByteQueue::ReadRequest(kj::mv(prp.resolver),
2057 {
2058 .store = jsg::BufferSource(js, source.detach(js)),
2059 .atLeast = atLeast,
2060 .type = ByteQueue::ReadRequest::Type::BYOB,
2061 }));
2062 } else KJ_IF_SOME(chunkSize, autoAllocateChunkSize) {
2063 // autoAllocateChunkSize is set, so we allocate a buffer and do a BYOB read.
2064 // This makes the buffer available to the underlying source via controller.byobRequest.
2065 KJ_IF_SOME(store, jsg::BufferSource::tryAlloc(js, chunkSize)) {
2066 // Ensure that the handle is created here so that the size of the buffer
2067 // is accounted for in the isolate memory tracking.
2068 s.consumer->read(js,
2069 ByteQueue::ReadRequest(kj::mv(prp.resolver),
2070 {
2071 .store = kj::mv(store),
2072 .type = ByteQueue::ReadRequest::Type::BYOB,
2073 }));
2074 } else {
2075 prp.resolver.reject(js, js.v8Error("Failed to allocate buffer for read."));
2076 }
2077 } else {
2078 // autoAllocateChunkSize is not set. Per spec, we do a DEFAULT read which means
2079 // the underlying source's pull method won't get a byobRequest. It must use
2080 // controller.enqueue() to provide data instead.
2081 constexpr size_t kDefaultReadSize = 16384; // 16KB default buffer
2082 KJ_IF_SOME(store, jsg::BufferSource::tryAlloc(js, kDefaultReadSize)) {
2083 s.consumer->read(js,
2084 ByteQueue::ReadRequest(kj::mv(prp.resolver),
2085 {
2086 .store = kj::mv(store),
2087 .type = ByteQueue::ReadRequest::Type::DEFAULT,
2088 }));
2089 } else {
2090 prp.resolver.reject(js, js.v8Error("Failed to allocate buffer for read."));
2091 }
2092 }
2093 
2094 return kj::mv(prp.promise);
2095 }
2096 
2097 // We are canceled! There's nothing else to do.
2098 KJ_IF_SOME(byob, byobOptions) {
2099 // If a BYOB buffer was given, we need to give it back wrapped in a TypedArray
2100 // whose size is set to zero.
2101 jsg::BufferSource source(js, byob.bufferView.getHandle(js));
2102 auto store = source.detach(js);
2103 store.consume(store.size());
2104 return js.resolvedPromise(ReadResult{
2105 .value = js.v8Ref(store.createHandle(js)),
2106 .done = true,
2107 });
2108 } else {
2109 return js.resolvedPromise(ReadResult{.done = true});
2110 }
2111 }
2112 
2113 jsg::Promise<DrainingReadResult> drainingRead(jsg::Lock& js, size_t maxRead) {
2114 KJ_IF_SOME(s, state) {
2115 // Note: We do NOT call beginOperation()/endOperation() here. The caller
2116 // (ReadableStreamJsController::drainingRead) manages the operation scope
2117 // around both this call and the returned promise's lifetime. See the
2118 // comment in ValueReadable::drainingRead for the detailed explanation.
2119 return s.consumer->drainingRead(js, maxRead);
2120 }
2121 
2122 // We are canceled! Return done with empty chunks.
2123 return js.resolvedPromise(DrainingReadResult{
2124 .chunks = kj::Array<kj::Array<kj::byte>>(),
2125 .done = true,
2126 });
2127 }
2128 
2129 // When a ReadableStream is canceled, the expected behavior is that the underlying
2130 // controller is notified and the cancel algorithm on the underlying source is
2131 // called. When there are multiple ReadableStreams sharing consumption of a
2132 // controller, however, it should act as a shared pointer of sorts, canceling
2133 // the underlying controller only when the last reader is canceled.
2134 // Here, we rely on the controller implementing the correct behavior since it owns
2135 // the queue that knows about all of the attached consumers.
2136 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
2137 if (pendingCancel) return js.resolvedPromise();
2138 KJ_IF_SOME(s, state) {
2139 // Check if there's a pending draining read before calling cancel, since cancel
2140 // will resolve the pending read and we need to know if we should defer destruction.
2141 bool hasPendingDrainingRead = s.consumer->hasPendingDrainingRead();
2142 s.consumer->cancel(js, maybeReason);
2143 auto promise = s.controller->cancel(js, kj::mv(maybeReason));
2144 // If there's a pending draining read, we need to wait for it to finish before
2145 // dropping our state. The draining read's promise callbacks capture 'this' (the
2146 // Consumer) to clear hasPendingDrainingRead. If we destroy the state now, those
2147 // callbacks will UAF.
2148 if (hasPendingDrainingRead) {
2149 pendingCancel = true;
2150 } else {
2151 state = kj::none;
2152 }
2153 return kj::mv(promise);
2154 }
2155 
2156 return js.resolvedPromise();
2157 }
2158 
2159 void onConsumerClose(jsg::Lock& js) override {
2160 // Note that the owner may drop this readable in doClose so it
2161 // is not safe to access anything on this after calling doClose.
2162 KJ_IF_SOME(s, state) {
2163 s.owner.doClose(js);
2164 }
2165 }
2166 
2167 void onConsumerError(jsg::Lock& js, jsg::Value reason) override {
2168 // Note that the owner may drop this readable in doClose so it
2169 // is not safe to access anything on this after calling doError.
2170 KJ_IF_SOME(s, state) {
2171 s.owner.doError(js, reason.getHandle(js));
2172 };
2173 }
2174 
2175 // Called by the consumer when it has a queued pending read and needs
2176 // data to be provided to fulfill it. We need to notify the controller
2177 // to initiate pulling to provide the data.
2178 // Returns true if the pull completed synchronously (meaning more pumping
2179 // might yield additional synchronous data), false otherwise.
2180 bool onConsumerWantsData(jsg::Lock& js) override {
2181 KJ_IF_SOME(s, state) {
2182 // Save a reference to the owner before calling pull. The pull callback
2183 // may trigger close/error which could destroy this ByteReadable. By
2184 // using beginOperation(), we ensure doClose/doError defers the
2185 // actual destruction until after we return.
2186 ReadableStreamJsController& owner = s.owner;
2187 owner.state.beginOperation();
2188 
2189 // For draining reads, use forcePull to bypass backpressure checks.
2190 // This ensures we pull all available data regardless of highWaterMark.
2191 if (s.consumer->hasPendingDrainingRead()) {
2192 s.controller->forcePull(js);
2193 } else {
2194 s.controller->pull(js);
2195 }
2196 
2197 // Check if state is still valid BEFORE calling endOperation(),
2198 // because that call may destroy this ByteReadable if close was deferred.
2199 bool result =
2200 state.map([](State& s2) { return !s2.controller->isPulling(); }).orDefault(false);
2201 
2202 // Process any deferred close/error. This may destroy this ByteReadable.
2203 if (owner.state.endOperation()) {
2204 // A pending state was applied. Call the appropriate callback.
2205 if (owner.state.template is<StreamStates::Closed>()) {
2206 owner.lock.onClose(js);
2207 } else if (owner.state.template is<StreamStates::Errored>()) {
2208 KJ_IF_SOME(err, owner.state.template tryGetUnsafe<StreamStates::Errored>()) {
2209 owner.lock.onError(js, err.getHandle(js));
2210 }
2211 }
2212 }
2213 
2214 return result;
2215 }
2216 return false;
2217 }
2218 
2219 kj::Maybe<int> getDesiredSize() {
2220 KJ_IF_SOME(s, state) {
2221 return s.controller->getDesiredSize();
2222 }
2223 return kj::none;
2224 }
2225 
2226 bool canCloseOrEnqueue() {
2227 return state.map([](State& s) { return s.controller->canCloseOrEnqueue(); }).orDefault(false);
2228 }
2229 
2230 kj::Maybe<ByobController> getControllerRef() {
2231 return state.map([](State& state) { return state.controller.addRef(); });
2232 }
2233};
2234} // namespace
2235 
2236// =======================================================================================
2237 
2238ReadableStreamDefaultController::ReadableStreamDefaultController(
2239 UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy)
2240 : ioContext(tryGetIoContext()),
2241 impl(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {}
2242 
2243kj::Maybe<StreamStates::Errored> ReadableStreamDefaultController::getMaybeErrorState(
2244 jsg::Lock& js) {
2245 KJ_IF_SOME(errored, impl.state.tryGetUnsafe<StreamStates::Errored>()) {
2246 return errored.addRef(js);
2247 }
2248 return kj::none;
2249}
2250 
2251void ReadableStreamDefaultController::start(jsg::Lock& js) {
2252 impl.start(js, JSG_THIS);
2253}
2254 
2255bool ReadableStreamDefaultController::canCloseOrEnqueue() {
2256 return impl.canCloseOrEnqueue();
2257}
2258 
2259bool ReadableStreamDefaultController::hasBackpressure() {
2260 return !impl.shouldCallPull();
2261}
2262 
2263kj::Maybe<int> ReadableStreamDefaultController::getDesiredSize() {
2264 return impl.getDesiredSize();
2265}
2266 
2267void ReadableStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) {
2268 visitor.visit(impl);
2269}
2270 
2271jsg::Promise<void> ReadableStreamDefaultController::cancel(
2272 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
2273 return impl.cancel(js, JSG_THIS, maybeReason.orDefault([&] { return js.v8Undefined(); }));
2274}
2275 
2276void ReadableStreamDefaultController::close(jsg::Lock& js) {
2277 impl.close(js);
2278}
2279 
2280void ReadableStreamDefaultController::enqueue(
2281 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> chunk) {
2282 // Hold a strong reference to prevent this controller from being freed if the
2283 // user-provided size algorithm (below) re-enters JS and errors the controller
2284 // through a side-channel (e.g. TransformStreamDefaultController::error()
2285 // dropping all external jsg::Refs to this controller).
2286 auto self = JSG_THIS;
2287 auto value = chunk.orDefault(js.undefined());
2288 
2289 JSG_REQUIRE(impl.canCloseOrEnqueue(), TypeError, "Unable to enqueue");
2290 
2291 size_t size = 1;
2292 bool errored = false;
2293 KJ_IF_SOME(sizeFunc, impl.algorithms.size) {
2294 js.tryCatch([&] { size = sizeFunc(js, value); }, [&](jsg::Value exception) {
2295 impl.doError(js, kj::mv(exception));
2296 errored = true;
2297 });
2298 }
2299 
2300 // Re-check canCloseOrEnqueue: the size callback may have errored us without
2301 // throwing (e.g. by calling transformController.error()), in which case
2302 // `errored` is still false but the impl state has transitioned to Errored.
2303 if (!errored && impl.canCloseOrEnqueue()) {
2304 impl.enqueue(js, kj::rc<ValueQueue::Entry>(js.v8Ref(value), size), kj::mv(self));
2305 }
2306}
2307 
2308void ReadableStreamDefaultController::error(jsg::Lock& js, v8::Local<v8::Value> reason) {
2309 impl.doError(js, js.v8Ref(reason));
2310}
2311 
2312// When a consumer receives a read request, but does not have the data available to
2313// fulfill the request, the consumer will call pull on the controller to pull that
2314// data if needed.
2315void ReadableStreamDefaultController::pull(jsg::Lock& js) {
2316 impl.pullIfNeeded(js, JSG_THIS);
2317}
2318 
2319void ReadableStreamDefaultController::forcePull(jsg::Lock& js) {
2320 impl.forcePullIfNeeded(js, JSG_THIS);
2321}
2322 
2323kj::Own<ValueQueue::Consumer> ReadableStreamDefaultController::getConsumer(
2324 kj::Maybe<ValueQueue::ConsumerImpl::StateListener&> stateListener) {
2325 return impl.getConsumer(stateListener);
2326}
2327 
2328// ======================================================================================
2329 
2330ReadableStreamBYOBRequest::Impl::Impl(jsg::Lock& js,
2331 kj::Own<ByteQueue::ByobRequest> readRequest,
2332 kj::Rc<WeakRef<ReadableByteStreamController>> controller)
2333 : readRequest(kj::mv(readRequest)),
2334 controller(kj::mv(controller)),
2335 view(js.v8Ref(this->readRequest->getView(js))),
2336 originalBufferByteLength(this->readRequest->getOriginalBufferByteLength(js)),
2337 originalByteOffsetPlusBytesFilled(this->readRequest->getOriginalByteOffsetPlusBytesFilled()) {
2338}
2339 
2340void ReadableStreamBYOBRequest::Impl::updateView(jsg::Lock& js) {
2341 jsg::check(view.getHandle(js)->Buffer()->Detach(v8::Local<v8::Value>()));
2342 view = js.v8Ref(readRequest->getView(js));
2343}
2344 
2345void ReadableStreamBYOBRequest::visitForGc(jsg::GcVisitor& visitor) {
2346 KJ_IF_SOME(impl, maybeImpl) {
2347 visitor.visit(impl.view);
2348 }
2349}
2350 
2351ReadableStreamBYOBRequest::ReadableStreamBYOBRequest(jsg::Lock& js,
2352 kj::Own<ByteQueue::ByobRequest> readRequest,
2353 kj::Rc<WeakRef<ReadableByteStreamController>> controller)
2354 : ioContext(tryGetIoContext()),
2355 maybeImpl(Impl(js, kj::mv(readRequest), kj::mv(controller))) {}
2356 
2357kj::Maybe<int> ReadableStreamBYOBRequest::getAtLeast() {
2358 KJ_IF_SOME(impl, maybeImpl) {
2359 return impl.readRequest->getAtLeast();
2360 }
2361 return kj::none;
2362}
2363 
2364kj::Maybe<jsg::V8Ref<v8::Uint8Array>> ReadableStreamBYOBRequest::getView(jsg::Lock& js) {
2365 KJ_IF_SOME(impl, maybeImpl) {
2366 return impl.view.addRef(js);
2367 }
2368 return kj::none;
2369}
2370 
2371void ReadableStreamBYOBRequest::invalidate(jsg::Lock& js) {
2372 KJ_IF_SOME(impl, maybeImpl) {
2373 // If the user code happened to have retained a reference to the view or
2374 // the buffer, we need to detach it so that those references cannot be used
2375 // to modify or observe modifications.
2376 jsg::check(impl.view.getHandle(js)->Buffer()->Detach(v8::Local<v8::Value>()));
2377 impl.controller->runIfAlive(
2378 [](ReadableByteStreamController& controller) { controller.maybeByobRequest = kj::none; });
2379 }
2380 maybeImpl = kj::none;
2381}
2382 
2383void ReadableStreamBYOBRequest::respond(jsg::Lock& js, int bytesWritten) {
2384 auto& impl = JSG_REQUIRE_NONNULL(
2385 maybeImpl, TypeError, "This ReadableStreamBYOBRequest has been invalidated.");
2386 JSG_REQUIRE(impl.controller->isValid(), Error, "The ReadableStreamBYOBRequest is invalid.");
2387 JSG_REQUIRE(impl.view.getHandle(js)->ByteLength() > 0, TypeError,
2388 "Cannot respond with a zero-length or detached view");
2389 impl.controller->runIfAlive([&](ReadableByteStreamController& controller) {
2390 if (!controller.canCloseOrEnqueue()) {
2391 JSG_REQUIRE(bytesWritten == 0, TypeError,
2392 "The bytesWritten must be zero after the stream is closed.");
2393 KJ_ASSERT(impl.readRequest->isInvalidated());
2394 invalidate(js);
2395 } else {
2396 bool shouldInvalidate = false;
2397 if (impl.readRequest->isInvalidated() && controller.impl.consumerCount() >= 1) {
2398 // While this particular request may be invalidated, there are still
2399 // other branches we can push the data to. Let's do so.
2400 jsg::BufferSource source(js, impl.view.getHandle(js));
2401 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, source.detach(js)));
2402 controller.impl.enqueue(js, kj::mv(entry), controller.getSelf());
2403 } else {
2404 JSG_REQUIRE(bytesWritten > 0, TypeError,
2405 "The bytesWritten must be more than zero while the stream is open.");
2406 if (impl.readRequest->respond(js, bytesWritten)) {
2407 // The read request was fulfilled, we need to invalidate.
2408 shouldInvalidate = true;
2409 } else {
2410 // The response did not fulfill the minimum requirements of the read.
2411 // We do not want to invalidate the read request and we need to update the
2412 // view so that on the next read the view will be properly adjusted.
2413 impl.updateView(js);
2414 }
2415 }
2416 controller.pull(js);
2417 if (shouldInvalidate) {
2418 invalidate(js);
2419 }
2420 }
2421 });
2422}
2423 
2424void ReadableStreamBYOBRequest::respondWithNewView(jsg::Lock& js, jsg::BufferSource view) {
2425 auto& impl = JSG_REQUIRE_NONNULL(
2426 maybeImpl, TypeError, "This ReadableStreamBYOBRequest has been invalidated.");
2427 JSG_REQUIRE(impl.controller->isValid(), Error, "The ReadableStreamBYOBRequest is invalid.");
2428 impl.controller->runIfAlive([&](ReadableByteStreamController& controller) {
2429 if (!controller.canCloseOrEnqueue()) {
2430 JSG_REQUIRE(view.size() == 0, TypeError,
2431 "The view byte length must be zero after the stream is closed.");
2432 
2433 if (FeatureFlags::get(js).getPedanticWpt()) {
2434 // Per the spec, when the stream is closed:
2435 // 1. The view byte length must be zero (TypeError if not)
2436 // 2. The underlying buffer must not be detached (TypeError)
2437 // 3. The buffer byte length must not be zero (RangeError)
2438 // 4. The buffer byte length must match the original (RangeError)
2439 auto handle = view.getHandle(js);
2440 auto buffer = handle->IsArrayBuffer() ? handle.As<v8::ArrayBuffer>()
2441 : handle.As<v8::ArrayBufferView>()->Buffer();
2442 JSG_REQUIRE(
2443 !buffer->WasDetached(), TypeError, "The underlying ArrayBuffer has been detached.");
2444 
2445 JSG_REQUIRE(view.canDetach(js), TypeError, "Unable to use non-detachable ArrayBuffer.");
2446 // Use the stored values since the ByobRequest may have been invalidated during close.
2447 auto actualBufferByteLength = buffer->ByteLength();
2448 JSG_REQUIRE(
2449 actualBufferByteLength != 0, RangeError, "The underlying ArrayBuffer is zero-length.");
2450 JSG_REQUIRE(actualBufferByteLength == impl.originalBufferByteLength, RangeError,
2451 "The underlying ArrayBuffer is not the correct length.");
2452 // The view's byte offset must match the original byte offset plus bytes filled.
2453 auto viewByteOffset =
2454 handle->IsArrayBuffer() ? 0 : handle.As<v8::ArrayBufferView>()->ByteOffset();
2455 JSG_REQUIRE(viewByteOffset == impl.originalByteOffsetPlusBytesFilled, RangeError,
2456 "The view has an invalid byte offset.");
2457 } else {
2458 KJ_ASSERT(impl.readRequest->isInvalidated());
2459 }
2460 
2461 invalidate(js);
2462 } else {
2463 bool shouldInvalidate = false;
2464 if (impl.readRequest->isInvalidated() && controller.impl.consumerCount() >= 1) {
2465 // While this particular request may be invalidated, there are still
2466 // other branches we can push the data to. Let's do so.
2467 auto entry = kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, view.detach(js)));
2468 controller.impl.enqueue(js, kj::mv(entry), controller.getSelf());
2469 } else {
2470 JSG_REQUIRE(view.size() > 0, TypeError,
2471 "The view byte length must be more than zero while the stream is open.");
2472 if (impl.readRequest->respondWithNewView(js, kj::mv(view))) {
2473 // The read request was fulfilled, we need to invalidate.
2474 shouldInvalidate = true;
2475 } else {
2476 // The response did not fulfill the minimum requirements of the read.
2477 // We do not want to invalidate the read request and we need to update the
2478 // view so that on the next read the view will be properly adjusted.
2479 impl.updateView(js);
2480 }
2481 }
2482 
2483 controller.pull(js);
2484 if (shouldInvalidate) {
2485 invalidate(js);
2486 }
2487 }
2488 });
2489}
2490 
2491bool ReadableStreamBYOBRequest::isPartiallyFulfilled() {
2492 KJ_IF_SOME(impl, maybeImpl) {
2493 return impl.readRequest->isPartiallyFulfilled();
2494 }
2495 return false;
2496}
2497 
2498// ======================================================================================
2499 
2500ReadableByteStreamController::ReadableByteStreamController(
2501 UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy)
2502 : weakSelf(kj::rc<WeakRef<ReadableByteStreamController>>(
2503 kj::Badge<ReadableByteStreamController>{}, *this)),
2504 ioContext(tryGetIoContext()),
2505 impl(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {}
2506 
2507ReadableByteStreamController::~ReadableByteStreamController() noexcept(false) {
2508 weakSelf->invalidate();
2509}
2510 
2511void ReadableByteStreamController::start(jsg::Lock& js) {
2512 impl.start(js, JSG_THIS);
2513}
2514 
2515bool ReadableByteStreamController::canCloseOrEnqueue() {
2516 return impl.canCloseOrEnqueue();
2517}
2518 
2519bool ReadableByteStreamController::hasBackpressure() {
2520 return !impl.shouldCallPull();
2521}
2522 
2523kj::Maybe<int> ReadableByteStreamController::getDesiredSize() {
2524 return impl.getDesiredSize();
2525}
2526 
2527void ReadableByteStreamController::visitForGc(jsg::GcVisitor& visitor) {
2528 visitor.visit(maybeByobRequest, impl);
2529}
2530 
2531jsg::Promise<void> ReadableByteStreamController::cancel(
2532 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
2533 KJ_IF_SOME(byobRequest, maybeByobRequest) {
2534 if (impl.consumerCount() == 1) {
2535 byobRequest->invalidate(js);
2536 }
2537 }
2538 return impl.cancel(js, JSG_THIS, maybeReason.orDefault(js.undefined()));
2539}
2540 
2541void ReadableByteStreamController::close(jsg::Lock& js) {
2542 KJ_IF_SOME(byobRequest, maybeByobRequest) {
2543 JSG_REQUIRE(!byobRequest->isPartiallyFulfilled(), TypeError,
2544 "This ReadableStream was closed with a partial read pending.");
2545 } else if (FeatureFlags::get(js).getPedanticWpt()) {
2546 // If maybeByobRequest is not set, check if there's a pending byob request.
2547 // If so, materialize it before closing so it remains accessible after
2548 // the state changes to Closed. This is required by the spec for proper
2549 // respondWithNewView() error handling in the closed state.
2550 // Only do this if the queue doesn't have a partially fulfilled read.
2551 KJ_IF_SOME(queue, impl.state.tryGetUnsafe<ByteQueue>()) {
2552 if (!queue.hasPartiallyFulfilledRead()) {
2553 getByobRequest(js);
2554 }
2555 }
2556 }
2557 impl.close(js);
2558}
2559 
2560void ReadableByteStreamController::enqueue(jsg::Lock& js, jsg::BufferSource chunk) {
2561 // Hold a strong reference up front. Operations below (invalidate, detach) touch
2562 // the JS heap and C++ argument evaluation order is unspecified, so JSG_THIS as a
2563 // function argument would not reliably precede chunk.detach(js).
2564 auto self = JSG_THIS;
2565 
2566 JSG_REQUIRE(chunk.size() > 0, TypeError, "Cannot enqueue a zero-length ArrayBuffer.");
2567 JSG_REQUIRE(chunk.canDetach(js), TypeError, "The provided ArrayBuffer must be detachable.");
2568 JSG_REQUIRE(impl.canCloseOrEnqueue(), TypeError, "This ReadableByteStreamController is closed.");
2569 
2570 KJ_IF_SOME(byobRequest, maybeByobRequest) {
2571 KJ_IF_SOME(view, byobRequest->getView(js)) {
2572 JSG_REQUIRE(view.getHandle(js)->ByteLength() > 0, TypeError,
2573 "The byobRequest.view is zero-length or was detached");
2574 }
2575 byobRequest->invalidate(js);
2576 }
2577 
2578 impl.enqueue(js, kj::rc<ByteQueue::Entry>(jsg::BufferSource(js, chunk.detach(js))), kj::mv(self));
2579}
2580 
2581void ReadableByteStreamController::error(jsg::Lock& js, v8::Local<v8::Value> reason) {
2582 impl.doError(js, js.v8Ref(reason));
2583}
2584 
2585kj::Maybe<jsg::Ref<ReadableStreamBYOBRequest>> ReadableByteStreamController::getByobRequest(
2586 jsg::Lock& js) {
2587 if (maybeByobRequest == kj::none) {
2588 KJ_IF_SOME(queue, impl.state.tryGetUnsafe<ByteQueue>()) {
2589 KJ_IF_SOME(pendingByob, queue.nextPendingByobReadRequest()) {
2590 maybeByobRequest =
2591 js.alloc<ReadableStreamBYOBRequest>(js, kj::mv(pendingByob), weakSelf.addRef());
2592 }
2593 } else {
2594 return kj::none;
2595 }
2596 }
2597 
2598 return maybeByobRequest.map(
2599 [&](jsg::Ref<ReadableStreamBYOBRequest>& req) { return req.addRef(); });
2600}
2601 
2602// When a consumer receives a read request, but does not have the data available to
2603// fulfill the request, the consumer will call pull on the controller to pull that
2604// data if needed.
2605void ReadableByteStreamController::pull(jsg::Lock& js) {
2606 impl.pullIfNeeded(js, JSG_THIS);
2607}
2608 
2609void ReadableByteStreamController::forcePull(jsg::Lock& js) {
2610 impl.forcePullIfNeeded(js, JSG_THIS);
2611}
2612 
2613kj::Own<ByteQueue::Consumer> ReadableByteStreamController::getConsumer(
2614 kj::Maybe<ByteQueue::ConsumerImpl::StateListener&> stateListener) {
2615 return impl.getConsumer(stateListener);
2616}
2617 
2618// ======================================================================================
2619 
2620ReadableStreamJsController::ReadableStreamJsController(): ioContext(tryGetIoContext()) {}
2621 
2622ReadableStreamJsController::ReadableStreamJsController(StreamStates::Closed closed)
2623 : ioContext(tryGetIoContext()) {
2624 state.transitionTo<StreamStates::Closed>();
2625}
2626 
2627ReadableStreamJsController::ReadableStreamJsController(StreamStates::Errored errored)
2628 : ioContext(tryGetIoContext()) {
2629 state.transitionTo<StreamStates::Errored>(kj::mv(errored));
2630}
2631 
2632ReadableStreamJsController::ReadableStreamJsController(jsg::Lock& js, ValueReadable& consumer)
2633 : ioContext(tryGetIoContext()) {
2634 state.transitionTo<kj::Own<ValueReadable>>(consumer.clone(js, *this));
2635}
2636 
2637ReadableStreamJsController::ReadableStreamJsController(jsg::Lock& js, ByteReadable& consumer)
2638 : ioContext(tryGetIoContext()) {
2639 state.transitionTo<kj::Own<ByteReadable>>(consumer.clone(js, *this));
2640}
2641 
2642jsg::Ref<ReadableStream> ReadableStreamJsController::addRef() {
2643 return KJ_REQUIRE_NONNULL(owner).addRef();
2644}
2645 
2646jsg::Promise<void> ReadableStreamJsController::cancel(
2647 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason) {
2648 disturbed = true;
2649 
2650 const auto doCancel = [&](auto& consumer) {
2651 auto reason = js.v8Ref(maybeReason.orDefault([&] { return js.v8Undefined(); }));
2652 KJ_DEFER(doClose(js));
2653 return consumer->cancel(js, reason.getHandle(js));
2654 };
2655 
2656 // Check for pending state first (deferred close/error during a read operation)
2657 if (state.pendingStateIs<StreamStates::Closed>()) {
2658 return js.resolvedPromise();
2659 }
2660 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
2661 return js.rejectedPromise<void>(pendingError.addRef(js));
2662 }
2663 
2664 KJ_SWITCH_ONEOF(state) {
2665 KJ_CASE_ONEOF(initial, Initial) {
2666 // Stream not yet set up, treat as closed.
2667 return js.resolvedPromise();
2668 }
2669 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
2670 return js.resolvedPromise();
2671 }
2672 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
2673 return js.rejectedPromise<void>(errored.addRef(js));
2674 }
2675 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
2676 if (canceling) return js.resolvedPromise();
2677 canceling = true;
2678 return doCancel(consumer);
2679 }
2680 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
2681 if (canceling) return js.resolvedPromise();
2682 canceling = true;
2683 return doCancel(consumer);
2684 }
2685 }
2686 
2687 KJ_UNREACHABLE;
2688}
2689 
2690// Finalizes the closed state of this ReadableStream. The connection to the underlying
2691// controller is released with no further action. Importantly, this method is triggered
2692// by the underlying controller as a result of that controller closing or being canceled.
2693// We detach ourselves from the underlying controller by releasing the ValueReadable or
2694// ByteReadable in the state and changing that to closed.
2695// We also clean up other state here.
2696void ReadableStreamJsController::doClose(jsg::Lock& js) {
2697 // If already in a terminal state, nothing to do.
2698 if (state.isTerminal()) return;
2699 
2700 // deferTransitionTo will defer if an operation is in progress, otherwise transition immediately.
2701 // Returns true if transition happened immediately.
2702 if (state.deferTransitionTo<StreamStates::Closed>()) {
2703 lock.onClose(js);
2704 }
2705 // If deferred, lock.onClose will be called when the pending state is applied
2706 // via applyPendingState in deferControllerStateChange.
2707}
2708 
2709// As with doClose(), doError() finalizes the error state of this ReadableStream.
2710// The connection to the underlying controller is released with no further action.
2711// This method is triggered by the underlying controller as a result of that controller
2712// erroring. We detach ourselves from the underlying controller by releasing the ValueReadable
2713// or ByteReadable in the state and changing that to errored.
2714// We also clean up other state here.
2715void ReadableStreamJsController::doError(jsg::Lock& js, v8::Local<v8::Value> reason) {
2716 // If already in a terminal state, nothing to do.
2717 if (state.isTerminal()) return;
2718 
2719 // deferTransitionTo will defer if an operation is in progress, otherwise transition immediately.
2720 // Returns true if transition happened immediately.
2721 if (state.deferTransitionTo<StreamStates::Errored>(js.v8Ref(reason))) {
2722 lock.onError(js, reason);
2723 }
2724 // If deferred, lock.onError will be called when the pending state is applied
2725 // via applyPendingState in deferControllerStateChange.
2726}
2727 
2728bool ReadableStreamJsController::isByteOriented() const {
2729 return state.is<kj::Own<ByteReadable>>();
2730}
2731 
2732bool ReadableStreamJsController::isClosedOrErrored() const {
2733 // Check if we're in a terminal state or have one pending
2734 return state.isTerminal() || state.hasPendingState();
2735}
2736 
2737bool ReadableStreamJsController::isClosed() const {
2738 // Check current state first, then pending state
2739 if (state.is<StreamStates::Closed>()) return true;
2740 return state.pendingStateIs<StreamStates::Closed>();
2741}
2742 
2743bool ReadableStreamJsController::isDisturbed() {
2744 return disturbed;
2745}
2746 
2747bool ReadableStreamJsController::isLockedToReader() const {
2748 return lock.isLockedToReader();
2749}
2750 
2751bool ReadableStreamJsController::lockReader(jsg::Lock& js, Reader& reader) {
2752 return lock.lockReader(js, *this, reader);
2753}
2754 
2755jsg::Promise<void> ReadableStreamJsController::pipeTo(
2756 jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) {
2757 KJ_DASSERT(!isLockedToReader());
2758 KJ_DASSERT(!destination.isLockedToWriter());
2759 
2760 disturbed = true;
2761 KJ_IF_SOME(promise, destination.tryPipeFrom(js, addRef(), kj::mv(options))) {
2762 return kj::mv(promise);
2763 }
2764 
2765 return js.rejectedPromise<void>(
2766 js.v8TypeError("This ReadableStream cannot be piped to this WritableStream"_kj));
2767}
2768 
2769kj::Maybe<jsg::Promise<ReadResult>> ReadableStreamJsController::read(
2770 jsg::Lock& js, kj::Maybe<ByobOptions> maybeByobOptions) {
2771 disturbed = true;
2772 
2773 KJ_IF_SOME(byobOptions, maybeByobOptions) {
2774 byobOptions.detachBuffer = true;
2775 auto view = byobOptions.bufferView.getHandle(js);
2776 if (!view->Buffer()->IsDetachable()) {
2777 return js.rejectedPromise<ReadResult>(
2778 js.v8TypeError("Unabled to use non-detachable ArrayBuffer."_kj));
2779 }
2780 
2781 if (view->ByteLength() == 0 || view->Buffer()->ByteLength() == 0) {
2782 return js.rejectedPromise<ReadResult>(
2783 js.v8TypeError("Unable to use a zero-length ArrayBuffer."_kj));
2784 }
2785 
2786 // Check for pending error first (deferred error during a prior read operation)
2787 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
2788 return js.rejectedPromise<ReadResult>(pendingError.addRef(js));
2789 }
2790 
2791 if (state.is<StreamStates::Closed>() || state.pendingStateIs<StreamStates::Closed>()) {
2792 // If it is a BYOB read, then the spec requires that we return an empty
2793 // view of the same type provided, that uses the same backing memory
2794 // as that provided, but with zero-length.
2795 auto source = jsg::BufferSource(js, byobOptions.bufferView.getHandle(js));
2796 auto store = source.detach(js);
2797 store.consume(store.size());
2798 return js.resolvedPromise(ReadResult{
2799 .value = js.v8Ref(store.createHandle(js)),
2800 .done = true,
2801 });
2802 }
2803 }
2804 
2805 // Check for pending state (deferred close/error during a prior read operation)
2806 if (state.pendingStateIs<StreamStates::Closed>()) {
2807 // The closed state for BYOB reads is handled in the maybeByobOptions check above.
2808 KJ_ASSERT(maybeByobOptions == kj::none);
2809 return js.resolvedPromise(ReadResult{.done = true});
2810 }
2811 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
2812 return js.rejectedPromise<ReadResult>(pendingError.addRef(js));
2813 }
2814 
2815 KJ_SWITCH_ONEOF(state) {
2816 KJ_CASE_ONEOF(initial, Initial) {
2817 // Stream not yet set up, treat as closed.
2818 KJ_ASSERT(maybeByobOptions == kj::none);
2819 return js.resolvedPromise(ReadResult{.done = true});
2820 }
2821 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
2822 // The closed state for BYOB reads is handled in the maybeByobOptions check above.
2823 KJ_ASSERT(maybeByobOptions == kj::none);
2824 return js.resolvedPromise(ReadResult{.done = true});
2825 }
2826 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
2827 return js.rejectedPromise<ReadResult>(errored.addRef(js));
2828 }
2829 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
2830 // The ReadableStreamDefaultController does not support ByobOptions.
2831 // It should never happen, but let's make sure.
2832 KJ_ASSERT(maybeByobOptions == kj::none);
2833 return deferControllerStateChange(js, *this, [&]() mutable { return consumer->read(js); });
2834 }
2835 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
2836 return deferControllerStateChange(
2837 js, *this, [&]() mutable { return consumer->read(js, kj::mv(maybeByobOptions)); });
2838 }
2839 }
2840 KJ_UNREACHABLE;
2841}
2842 
2843kj::Maybe<jsg::Promise<DrainingReadResult>> ReadableStreamJsController::drainingRead(
2844 jsg::Lock& js, size_t maxRead) {
2845 disturbed = true;
2846 
2847 // Check for pending state first (deferred close/error during a prior read operation)
2848 if (state.pendingStateIs<StreamStates::Closed>()) {
2849 return js.resolvedPromise(DrainingReadResult{
2850 .chunks = kj::Array<kj::Array<kj::byte>>(),
2851 .done = true,
2852 });
2853 }
2854 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
2855 return js.rejectedPromise<DrainingReadResult>(pendingError.addRef(js));
2856 }
2857 
2858 // Like deferControllerStateChange for regular reads, we need to prevent the controller
2859 // state from being destroyed while a draining read's promise callbacks are pending.
2860 // The drainingRead implementation captures `this` (the Consumer) in promise lambdas to
2861 // clear hasPendingDrainingRead. If the state is changed (destroying the Consumer) before
2862 // those callbacks run, we get a use-after-free.
2863 //
2864 // CRITICAL: state.beginOperation() MUST be called BEFORE consumer->drainingRead(), not
2865 // after. The consumer->drainingRead() call may trigger onConsumerWantsData -> forcePull
2866 // -> close/error, which calls deferTransitionTo. If no operation is in progress at that
2867 // point, the transition fires immediately, destroying the Consumer while we're still
2868 // inside its method and before the returned promise's .then() callbacks are set up.
2869 // The endOperation() happens in the .then() callbacks below, ensuring the deferred
2870 // state change only fires after the promise resolves/rejects and the Consumer's
2871 // this-capturing callbacks have already run.
2872 auto wrapDrainingRead =
2873 [this](jsg::Lock& js,
2874 jsg::Promise<DrainingReadResult> promise) -> jsg::Promise<DrainingReadResult> {
2875 return promise.then(js, [this](jsg::Lock& js, DrainingReadResult result) {
2876 if (state.endOperation()) {
2877 // A pending state was applied. Call the appropriate callback.
2878 if (state.template is<StreamStates::Closed>()) {
2879 lock.onClose(js);
2880 } else if (state.template is<StreamStates::Errored>()) {
2881 KJ_IF_SOME(err, state.template tryGetUnsafe<StreamStates::Errored>()) {
2882 lock.onError(js, err.getHandle(js));
2883 // The error was applied during this operation โ€” the data we collected
2884 // may be invalid. Discard it and propagate the error rather than
2885 // silently returning possibly-corrupt data.
2886 js.throwException(err.addRef(js));
2887 }
2888 }
2889 }
2890 return kj::mv(result);
2891 }, [this](jsg::Lock& js, jsg::Value exception) -> DrainingReadResult {
2892 state.clearPendingState();
2893 (void)state.endOperation();
2894 js.throwException(kj::mv(exception));
2895 });
2896 };
2897 
2898 KJ_SWITCH_ONEOF(state) {
2899 KJ_CASE_ONEOF(initial, Initial) {
2900 // Stream not yet set up, treat as closed.
2901 return js.resolvedPromise(DrainingReadResult{
2902 .chunks = kj::Array<kj::Array<kj::byte>>(),
2903 .done = true,
2904 });
2905 }
2906 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
2907 return js.resolvedPromise(DrainingReadResult{
2908 .chunks = kj::Array<kj::Array<kj::byte>>(),
2909 .done = true,
2910 });
2911 }
2912 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
2913 return js.rejectedPromise<DrainingReadResult>(errored.addRef(js));
2914 }
2915 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
2916 // beginOperation MUST be before consumer->drainingRead() โ€” see comment above.
2917 state.beginOperation();
2918 JSG_TRY(js) {
2919 return wrapDrainingRead(js, consumer->drainingRead(js, maxRead));
2920 }
2921 JSG_CATCH(exception) {
2922 state.clearPendingState();
2923 (void)state.endOperation();
2924 doError(js, exception.getHandle(js));
2925 return js.rejectedPromise<DrainingReadResult>(kj::mv(exception));
2926 };
2927 }
2928 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
2929 // beginOperation MUST be before consumer->drainingRead() โ€” see comment above.
2930 state.beginOperation();
2931 JSG_TRY(js) {
2932 return wrapDrainingRead(js, consumer->drainingRead(js, maxRead));
2933 }
2934 JSG_CATCH(exception) {
2935 state.clearPendingState();
2936 (void)state.endOperation();
2937 doError(js, exception.getHandle(js));
2938 return js.rejectedPromise<DrainingReadResult>(kj::mv(exception));
2939 };
2940 }
2941 }
2942 KJ_UNREACHABLE;
2943}
2944 
2945void ReadableStreamJsController::releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) {
2946 lock.releaseReader(*this, reader, maybeJs);
2947}
2948 
2949ReadableStreamController::Tee ReadableStreamJsController::tee(jsg::Lock& js) {
2950 JSG_REQUIRE(!isLockedToReader(), TypeError, "This ReadableStream is locked to a reader.");
2951 lock.state.transitionTo<Locked>();
2952 disturbed = true;
2953 
2954 // This will leave this stream locked, disturbed, and closed.
2955 
2956 // Check for pending state first (deferred close/error during a prior read operation)
2957 if (state.pendingStateIs<StreamStates::Closed>()) {
2958 return Tee{
2959 .branch1 =
2960 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())),
2961 .branch2 =
2962 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())),
2963 };
2964 }
2965 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
2966 return Tee{
2967 .branch1 =
2968 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(pendingError.addRef(js))),
2969 .branch2 =
2970 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(pendingError.addRef(js))),
2971 };
2972 }
2973 
2974 KJ_SWITCH_ONEOF(state) {
2975 KJ_CASE_ONEOF(initial, Initial) {
2976 // Stream not yet set up, treat as closed.
2977 return Tee{
2978 .branch1 =
2979 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())),
2980 .branch2 =
2981 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())),
2982 };
2983 }
2984 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
2985 return Tee{
2986 .branch1 =
2987 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())),
2988 .branch2 =
2989 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(StreamStates::Closed())),
2990 };
2991 }
2992 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
2993 return Tee{
2994 .branch1 =
2995 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(errored.addRef(js))),
2996 .branch2 =
2997 js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(errored.addRef(js))),
2998 };
2999 }
3000 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
3001 KJ_DEFER(state.transitionTo<StreamStates::Closed>());
3002 // We create two additional streams that clone this stream's consumer state,
3003 // then close this stream's consumer.
3004 return Tee{
3005 .branch1 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)),
3006 .branch2 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)),
3007 };
3008 }
3009 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
3010 KJ_DEFER(state.transitionTo<StreamStates::Closed>());
3011 // We create two additional streams that clone this stream's consumer state,
3012 // then close this stream's consumer.
3013 return Tee{
3014 .branch1 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)),
3015 .branch2 = js.alloc<ReadableStream>(kj::heap<ReadableStreamJsController>(js, *consumer)),
3016 };
3017 }
3018 }
3019 KJ_UNREACHABLE;
3020}
3021 
3022void ReadableStreamJsController::setOwnerRef(ReadableStream& stream) {
3023 KJ_ASSERT(owner == kj::none);
3024 owner = &stream;
3025}
3026 
3027void ReadableStreamJsController::setup(jsg::Lock& js,
3028 jsg::Optional<UnderlyingSource> maybeUnderlyingSource,
3029 jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) {
3030 auto underlyingSource = kj::mv(maybeUnderlyingSource).orDefault({});
3031 auto queuingStrategy = kj::mv(maybeQueuingStrategy).orDefault({});
3032 auto type = underlyingSource.type.map([](kj::StringPtr s) { return s; }).orDefault(""_kj);
3033 
3034 expectedLength = underlyingSource.expectedLength;
3035 
3036 if (type == "bytes") {
3037 // Per spec, autoAllocateChunkSize should only be set if the user explicitly provides it.
3038 // If not set, the underlying source's pull method won't receive a byobRequest for
3039 // non-BYOB reads and must use controller.enqueue() instead.
3040 //
3041 // However, our original implementation always defaulted to 4096, so we need a compat flag
3042 // to control this behavior. Default to legacy behavior if flags aren't available.
3043 bool useSpecCompliantBehavior = false;
3044 KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) {
3045 useSpecCompliantBehavior = flags.getNoAutoAllocateChunkSize();
3046 }
3047 
3048 kj::Maybe<int> autoAllocateChunkSize;
3049 if (useSpecCompliantBehavior) {
3050 // Spec-compliant: only set if user explicitly provides it
3051 autoAllocateChunkSize =
3052 underlyingSource.autoAllocateChunkSize.map([](int size) { return size; });
3053 } else {
3054 // Legacy behavior: apply a default autoAllocateChunkSize if not provided.
3055 auto defaultChunkSize = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE;
3056 if (util::Autogate::isEnabled(util::AutogateKey::UPDATED_AUTO_ALLOCATE_CHUNK_SIZE)) {
3057 defaultChunkSize = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2;
3058 }
3059 autoAllocateChunkSize = underlyingSource.autoAllocateChunkSize.orDefault(defaultChunkSize);
3060 }
3061 
3062 auto controller =
3063 js.alloc<ReadableByteStreamController>(kj::mv(underlyingSource), kj::mv(queuingStrategy));
3064 
3065 KJ_IF_SOME(chunkSize, autoAllocateChunkSize) {
3066 JSG_REQUIRE(chunkSize > 0, TypeError, "The autoAllocateChunkSize option cannot be zero.");
3067 }
3068 
3069 // We account for the memory usage of the ByteReadable and its controller together because
3070 // their lifetimes are identical (in practice) and memory accounting itself has a memory
3071 // overhead. The same applies to ValueReadable below.
3072 state.transitionTo<kj::Own<ByteReadable>>(
3073 kj::heap<ByteReadable>(controller.addRef(), *this, autoAllocateChunkSize)
3074 .attach(js.getExternalMemoryAdjustment(
3075 sizeof(ByteReadable) + sizeof(ReadableByteStreamController))));
3076 controller->start(js);
3077 } else {
3078 JSG_REQUIRE(
3079 type == "", TypeError, kj::str("\"", type, "\" is not a valid type of ReadableStream."));
3080 auto controller = js.alloc<ReadableStreamDefaultController>(
3081 kj::mv(underlyingSource), kj::mv(queuingStrategy));
3082 state.transitionTo<kj::Own<ValueReadable>>(
3083 kj::heap<ValueReadable>(controller.addRef(), *this)
3084 .attach(js.getExternalMemoryAdjustment(
3085 sizeof(ValueReadable) + sizeof(ReadableStreamDefaultController))));
3086 controller->start(js);
3087 }
3088}
3089 
3090kj::Maybe<ReadableStreamController::PipeController&> ReadableStreamJsController::tryPipeLock() {
3091 return lock.tryPipeLock(*this);
3092}
3093 
3094void ReadableStreamJsController::visitForGc(jsg::GcVisitor& visitor) {
3095 // Visit pending state if it's an error (Closed has no GC-traceable content)
3096 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
3097 visitor.visit(pendingError);
3098 }
3099 
3100 // Note: We cannot use state.visitForGc(visitor) here because the state machine's
3101 // visitForGc passes kj::Own<T>& to visitor.visit(), but GcVisitor expects T& for
3102 // types with visitForGc methods. We must dereference kj::Own manually.
3103 KJ_SWITCH_ONEOF(state) {
3104 KJ_CASE_ONEOF(initial, Initial) {}
3105 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
3106 KJ_CASE_ONEOF(error, StreamStates::Errored) {
3107 visitor.visit(error);
3108 }
3109 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
3110 visitor.visit(*consumer);
3111 }
3112 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
3113 visitor.visit(*consumer);
3114 }
3115 }
3116 visitor.visit(lock);
3117}
3118 
3119kj::Maybe<int> ReadableStreamJsController::getDesiredSize() {
3120 // If there's a pending state transition, return none
3121 if (state.hasPendingState()) {
3122 return kj::none;
3123 }
3124 
3125 KJ_SWITCH_ONEOF(state) {
3126 KJ_CASE_ONEOF(initial, Initial) {
3127 return kj::none;
3128 }
3129 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3130 return kj::none;
3131 }
3132 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3133 return kj::none;
3134 }
3135 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
3136 return consumer->getDesiredSize();
3137 }
3138 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
3139 return consumer->getDesiredSize();
3140 }
3141 }
3142 KJ_UNREACHABLE;
3143}
3144 
3145kj::Maybe<v8::Local<v8::Value>> ReadableStreamJsController::isErrored(jsg::Lock& js) {
3146 // Check for pending error first
3147 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
3148 return pendingError.getHandle(js);
3149 }
3150 // Pending Closed means not errored, so we can just check current state
3151 return state.tryGetUnsafe<StreamStates::Errored>().map(
3152 [&](jsg::Value& reason) { return reason.getHandle(js); });
3153}
3154 
3155bool ReadableStreamJsController::canCloseOrEnqueue() {
3156 // If there's a pending state transition, can't close or enqueue
3157 if (state.hasPendingState()) {
3158 return false;
3159 }
3160 
3161 KJ_SWITCH_ONEOF(state) {
3162 KJ_CASE_ONEOF(initial, Initial) {
3163 return false;
3164 }
3165 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3166 return false;
3167 }
3168 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3169 return false;
3170 }
3171 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
3172 return consumer->canCloseOrEnqueue();
3173 }
3174 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
3175 return consumer->canCloseOrEnqueue();
3176 }
3177 }
3178 KJ_UNREACHABLE;
3179}
3180 
3181bool ReadableStreamJsController::hasBackpressure() {
3182 KJ_IF_SOME(size, getDesiredSize()) {
3183 return size <= 0;
3184 }
3185 return false;
3186}
3187 
3188kj::Maybe<kj::OneOf<DefaultController, ByobController>> ReadableStreamJsController::
3189 getController() {
3190 // If there's a pending state transition, return none
3191 if (state.hasPendingState()) {
3192 return kj::none;
3193 }
3194 KJ_SWITCH_ONEOF(state) {
3195 KJ_CASE_ONEOF(initial, Initial) {
3196 return kj::none;
3197 }
3198 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3199 return kj::none;
3200 }
3201 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3202 return kj::none;
3203 }
3204 KJ_CASE_ONEOF(consumer, kj::Own<ValueReadable>) {
3205 return consumer->getControllerRef();
3206 }
3207 KJ_CASE_ONEOF(consumer, kj::Own<ByteReadable>) {
3208 return consumer->getControllerRef();
3209 }
3210 }
3211 KJ_UNREACHABLE;
3212}
3213 
3214namespace {
3215// Consumes all bytes from a stream, buffering in memory, with the purpose
3216// of producing either a single concatenated kj::Array<byte> or kj::String.
3217class AllReader {
3218 public:
3219 using PartList = kj::Array<kj::ArrayPtr<byte>>;
3220 
3221 AllReader(jsg::Ref<ReadableStream> stream, uint64_t limit)
3222 : state(State::create<jsg::Ref<ReadableStream>>(kj::mv(stream))),
3223 limit(limit) {}
3224 KJ_DISALLOW_COPY_AND_MOVE(AllReader);
3225 
3226 jsg::Promise<jsg::BufferSource> allBytes(jsg::Lock& js) {
3227 return loop(js).then(js, [this](auto& js, PartList&& partPtrs) -> jsg::BufferSource {
3228 auto out = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, runningTotal);
3229 copyInto(out.asArrayPtr(), partPtrs.asPtr());
3230 return jsg::BufferSource(js, kj::mv(out));
3231 });
3232 }
3233 
3234 jsg::Promise<kj::String> allText(
3235 jsg::Lock& js, ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE) {
3236 return loop(js).then(js, [this, option](auto& js, PartList&& partPtrs) {
3237 // Strip UTF-8 BOM if requested
3238 if ((option & ReadAllTextOption::STRIP_BOM) && partPtrs.size() > 0 &&
3239 hasUtf8Bom(partPtrs[0])) {
3240 partPtrs[0] = partPtrs[0].slice(UTF8_BOM_SIZE);
3241 runningTotal -= UTF8_BOM_SIZE;
3242 }
3243 
3244 JSG_REQUIRE(runningTotal <= v8::String::kMaxLength, RangeError,
3245 "String length exceeds v8::String::kMaxLength.");
3246 
3247 auto out = kj::heapArray<char>(runningTotal + 1);
3248 copyInto(out.first(out.size() - 1).asBytes(), partPtrs.asPtr());
3249 out.back() = '\0';
3250 return kj::String(kj::mv(out));
3251 });
3252 }
3253 
3254 void visitForGc(jsg::GcVisitor& visitor) {
3255 state.visitForGc(visitor);
3256 }
3257 
3258 private:
3259 // State machine for AllReader:
3260 // Closed is terminal, Errored is implicitly terminal via ErrorState.
3261 // jsg::Ref<ReadableStream> is the active state (still reading).
3262 using State = StateMachine<TerminalStates<StreamStates::Closed>,
3263 ErrorState<StreamStates::Errored>,
3264 ActiveState<jsg::Ref<ReadableStream>>,
3265 StreamStates::Closed,
3266 StreamStates::Errored,
3267 jsg::Ref<ReadableStream>>;
3268 State state;
3269 uint64_t limit;
3270 kj::Vector<jsg::BufferSource> parts;
3271 uint64_t runningTotal = 0;
3272 
3273 jsg::Promise<PartList> loop(jsg::Lock& js) {
3274 KJ_SWITCH_ONEOF(state) {
3275 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3276 return js.resolvedPromise(KJ_MAP(p, parts) { return p.asArrayPtr(); });
3277 }
3278 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3279 return js.template rejectedPromise<PartList>(errored.getHandle(js));
3280 }
3281 KJ_CASE_ONEOF(readable, jsg::Ref<ReadableStream>) {
3282 // Note that these nested lambda retain references to `this` and `readable`
3283 // and are passed into to promise returned by this method. It is the responsibility
3284 // of the caller to ensure that the AllReader instance is kept alive until the
3285 // promise is settled.
3286 auto onSuccess = JSG_VISITABLE_LAMBDA((this, readable = readable.addRef()), (readable),
3287 (jsg::Lock & js, ReadResult result) mutable->jsg::Promise<PartList> {
3288 if (result.done) {
3289 state.template transitionTo<StreamStates::Closed>();
3290 return loop(js);
3291 }
3292 
3293 // If we're not done, the result value must be interpretable as
3294 // bytes for the read to make any sense.
3295 auto handle = KJ_ASSERT_NONNULL(result.value).getHandle(js);
3296 if (!handle->IsArrayBufferView() && !handle->IsArrayBuffer()) {
3297 auto error = js.v8TypeError("This ReadableStream did not return bytes.");
3298 state.template transitionTo<StreamStates::Errored>(js.v8Ref(error));
3299 return readable->getController().cancel(js, error).then(
3300 js, [&](jsg::Lock& js) { return loop(js); });
3301 }
3302 
3303 jsg::BufferSource bufferSource(js, handle);
3304 
3305 if (bufferSource.size() == 0) {
3306 // Weird but allowed, we'll skip it.
3307 return loop(js);
3308 }
3309 
3310 if ((runningTotal + bufferSource.size()) > limit) {
3311 auto error = js.v8TypeError("Memory limit exceeded before EOF.");
3312 state.template transitionTo<StreamStates::Errored>(js.v8Ref(error));
3313 return readable->getController().cancel(js, error).then(
3314 js, [&](jsg::Lock& js) { return loop(js); });
3315 }
3316 
3317 runningTotal += bufferSource.size();
3318 parts.add(bufferSource.copy(js));
3319 return loop(js);
3320 });
3321 
3322 auto onFailure = [this](auto& js, jsg::Value exception) -> jsg::Promise<PartList> {
3323 // In this case the stream should already be errored.
3324 state.template transitionTo<StreamStates::Errored>(js.v8Ref(exception.getHandle(js)));
3325 return loop(js);
3326 };
3327 
3328 return maybeAddFunctor(js, KJ_ASSERT_NONNULL(readable->getController().read(js, kj::none)),
3329 kj::mv(onSuccess), kj::mv(onFailure));
3330 }
3331 }
3332 KJ_UNREACHABLE;
3333 }
3334 
3335 void copyInto(kj::ArrayPtr<byte> out, kj::ArrayPtr<kj::ArrayPtr<byte>> in) {
3336 for (auto& part: in) {
3337 KJ_ASSERT(part.size() <= out.size());
3338 out.first(part.size()).copyFrom(part);
3339 out = out.slice(part.size());
3340 }
3341 }
3342};
3343 
3344// PumpToReader implements the original JS promise-loop approach to pumping data from
3345// a ReadableStream to a WritableStreamSink. It reads one chunk at a time using the
3346// standard read() API, writes each chunk to the sink, and loops until done or errored.
3347// This is the fallback path used when the ENABLE_DRAINING_READ_ON_STANDARD_STREAMS
3348// autogate is not enabled.
3349class PumpToReader {
3350 public:
3351 PumpToReader(jsg::Ref<ReadableStream> stream, kj::Own<WritableStreamSink> sink, bool end)
3352 : ioContext(IoContext::current()),
3353 state(State::create<jsg::Ref<ReadableStream>>(kj::mv(stream))),
3354 sink(kj::mv(sink)),
3355 self(kj::refcounted<WeakRef<PumpToReader>>(kj::Badge<PumpToReader>{}, *this)),
3356 end(end) {}
3357 KJ_DISALLOW_COPY_AND_MOVE(PumpToReader);
3358 
3359 ~PumpToReader() noexcept(false) {
3360 self->invalidate();
3361 // Ensure that if a write promise is pending it is proactively canceled.
3362 canceler.cancel("PumpToReader was destroyed");
3363 }
3364 
3365 kj::Promise<void> pumpTo(jsg::Lock& js) {
3366 ioContext.requireCurrentOrThrowJs();
3367 KJ_SWITCH_ONEOF(state) {
3368 KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) {
3369 auto readable = stream.addRef();
3370 state.template transitionTo<Pumping>();
3371 return ioContext.awaitJs(
3372 js, pumpLoop(js, ioContext, kj::mv(readable), ioContext.addObject(self->addRef())));
3373 }
3374 KJ_CASE_ONEOF(pumping, Pumping) {
3375 return KJ_EXCEPTION(FAILED, "pumping is already in progress");
3376 }
3377 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3378 return KJ_EXCEPTION(FAILED, "stream has already been consumed");
3379 }
3380 KJ_CASE_ONEOF(errored, kj::Exception) {
3381 return errored.clone();
3382 }
3383 }
3384 KJ_UNREACHABLE;
3385 }
3386 
3387 private:
3388 struct Pumping {
3389 static constexpr kj::StringPtr NAME KJ_UNUSED = "pumping"_kj;
3390 };
3391 IoContext& ioContext;
3392 
3393 using State = StateMachine<TerminalStates<StreamStates::Closed>,
3394 ErrorState<kj::Exception>,
3395 Pumping,
3396 StreamStates::Closed,
3397 kj::Exception,
3398 jsg::Ref<ReadableStream>>;
3399 State state;
3400 kj::Own<WritableStreamSink> sink;
3401 kj::Own<WeakRef<PumpToReader>> self;
3402 kj::Canceler canceler;
3403 bool end;
3404 
3405 bool isErroredOrClosed() {
3406 return state.isTerminal();
3407 }
3408 
3409 jsg::Promise<void> pumpLoop(jsg::Lock& js,
3410 IoContext& ioContext,
3411 jsg::Ref<ReadableStream> readable,
3412 IoOwn<WeakRef<PumpToReader>> pumpToReader) {
3413 ioContext.requireCurrentOrThrowJs();
3414 
3415 KJ_SWITCH_ONEOF(state) {
3416 KJ_CASE_ONEOF(ready, jsg::Ref<ReadableStream>) {
3417 KJ_UNREACHABLE;
3418 }
3419 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3420 return end ? ioContext.awaitIoLegacy(js, sink->end().attach(kj::mv(sink)))
3421 : js.resolvedPromise();
3422 }
3423 KJ_CASE_ONEOF(errored, kj::Exception) {
3424 if (end) {
3425 sink->abort(errored.clone());
3426 }
3427 return js.rejectedPromise<void>(errored.clone());
3428 }
3429 KJ_CASE_ONEOF(pumping, Pumping) {
3430 using Result = kj::OneOf<Pumping, kj::Array<kj::byte>, StreamStates::Closed, jsg::Value>;
3431 
3432 return KJ_ASSERT_NONNULL(readable->getController().read(js, kj::none))
3433 .then(js,
3434 ioContext.addFunctor([byteStream = readable->getController().isByteOriented()](
3435 auto& js, ReadResult result) mutable -> Result {
3436 if (result.done) {
3437 return StreamStates::Closed();
3438 }
3439 
3440 auto handle = KJ_ASSERT_NONNULL(result.value).getHandle(js);
3441 if (!handle->IsArrayBufferView() && !handle->IsArrayBuffer()) {
3442 return js.v8Ref(js.v8TypeError("This ReadableStream did not return bytes."));
3443 }
3444 
3445 jsg::BufferSource bufferSource(js, handle);
3446 if (bufferSource.size() == 0) {
3447 return Pumping{};
3448 }
3449 
3450 if (byteStream) {
3451 jsg::BackingStore backing = bufferSource.detach(js);
3452 return backing.asArrayPtr().attach(kj::mv(backing));
3453 }
3454 return bufferSource.asArrayPtr().attach(kj::mv(bufferSource));
3455 }),
3456 [](auto& js, jsg::Value exception) mutable -> Result { return kj::mv(exception); })
3457 .then(js, ioContext.addFunctor( JSG_VISITABLE_LAMBDA((readable = kj::mv(readable), pumpToReader = kj::mv(pumpToReader)), (readable), (jsg::Lock & js, Result result) mutable {
3458 KJ_IF_SOME(reader, pumpToReader->tryGet()) {
3459 reader.ioContext.requireCurrentOrThrowJs();
3460 auto& ioContext = IoContext::current();
3461 KJ_SWITCH_ONEOF(result) {
3462 KJ_CASE_ONEOF(bytes, kj::Array<kj::byte>) {
3463 auto promise = reader.sink->write(bytes).attach(kj::mv(bytes));
3464 return ioContext.awaitIo(js, reader.canceler.wrap(kj::mv(promise)))
3465 .then(js,
3466 [](jsg::Lock& js) -> kj::Maybe<jsg::Value> {
3467 return kj::Maybe<jsg::Value>(kj::none);
3468 },
3469 [](jsg::Lock& js, jsg::Value exception) mutable -> kj::Maybe<jsg::Value> {
3470 return kj::mv(exception);
3471 })
3472 .then(js,
3473 ioContext.addFunctor(JSG_VISITABLE_LAMBDA(
3474 (readable = readable.addRef(), pumpToReader = kj::mv(pumpToReader)),
3475 (readable),
3476 (jsg::Lock & js, kj::Maybe<jsg::Value> maybeException) mutable {
3477 KJ_IF_SOME(reader, pumpToReader->tryGet()) {
3478 auto& ioContext = reader.ioContext;
3479 ioContext.requireCurrentOrThrowJs();
3480 KJ_IF_SOME(exception, maybeException) {
3481 if (!reader.isErroredOrClosed()) {
3482 reader.state.transitionTo<kj::Exception>(
3483 js.exceptionToKj(kj::mv(exception)));
3484 }
3485 } else {
3486 // Else block to avert dangling else compiler warning.
3487 }
3488 return reader.pumpLoop(
3489 js, ioContext, readable.addRef(), kj::mv(pumpToReader));
3490 } else {
3491 return readable->getController().cancel(js,
3492 maybeException.map(
3493 [&](jsg::Value& ex) { return ex.getHandle(js); }));
3494 }
3495 })));
3496 }
3497 KJ_CASE_ONEOF(pumping, Pumping) {}
3498 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3499 if (!reader.isErroredOrClosed()) {
3500 reader.state.transitionTo<StreamStates::Closed>();
3501 }
3502 }
3503 KJ_CASE_ONEOF(exception, jsg::Value) {
3504 if (!reader.isErroredOrClosed()) {
3505 reader.state.transitionTo<kj::Exception>(js.exceptionToKj(kj::mv(exception)));
3506 }
3507 }
3508 }
3509 return reader.pumpLoop(js, ioContext, readable.addRef(), kj::mv(pumpToReader));
3510 } else {
3511 KJ_SWITCH_ONEOF(result) {
3512 KJ_CASE_ONEOF(bytes, kj::Array<kj::byte>) {
3513 return readable->getController().cancel(js, kj::none);
3514 }
3515 KJ_CASE_ONEOF(pumping, Pumping) {
3516 return readable->getController().cancel(js, kj::none);
3517 }
3518 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3519 return js.resolvedPromise();
3520 }
3521 KJ_CASE_ONEOF(exception, jsg::Value) {
3522 return readable->getController().cancel(js, exception.getHandle(js));
3523 }
3524 }
3525 }
3526 KJ_UNREACHABLE;
3527 })));
3528 }
3529 }
3530 KJ_UNREACHABLE;
3531 }
3532};
3533 
3534// pumpToCoroutine uses a DrainingReader to efficiently pull all synchronously available
3535// data from the stream in each iteration, then writes it to the sink using vectored
3536// I/O. This minimizes isolate lock acquisitions by batching: each time the lock is
3537// held, the stream's internal queue is fully drained and the JS pull callback is
3538// pumped synchronously as many times as possible.
3539//
3540// The pump loop is a kj coroutine. Dropping the returned kj::Promise drops the
3541// coroutine frame, which destroys the DrainingReader (releasing the stream lock)
3542// and the sink. No WeakRef/IoOwn dance is needed because ownership is clear.
3543// The coroutine that implements the pump loop takes ownership of the DrainingReader
3544// and sink. The jsg::Ref<ReadableStream> is not passed into the coroutine because
3545// jsg::Ref is disallowed in coroutine parameters; instead, the DrainingReader holds
3546// a reference to the stream internally.
3547kj::Promise<void> pumpToImpl(IoContext& ioContext,
3548 kj::Own<DrainingReader> reader,
3549 kj::Own<WritableStreamSink> sink,
3550 bool end) {
3551 
3552 bool writeFailed = false;
3553 
3554 KJ_TRY {
3555 while (true) {
3556 // Perform a draining read to get all synchronously available data if possible
3557 // or fall back to a regular read if not.
3558 DrainingReadResult result = co_await ioContext.run([&reader](jsg::Lock& js) mutable {
3559 auto& ioContext = IoContext::current();
3560 // Use a 256KB limit to allow periodic yielding to the event loop,
3561 // preventing a fast producer from monopolizing the thread.
3562 constexpr size_t kMaxReadPerCycle = 256 * 1024;
3563 return ioContext.awaitJs(js, reader->read(js, kMaxReadPerCycle));
3564 });
3565 
3566 // Write all the chunks we received using vectored write for efficiency.
3567 if (result.chunks.size() > 0) {
3568 KJ_ON_SCOPE_FAILURE(writeFailed = true);
3569 auto pieces =
3570 KJ_MAP(chunk, result.chunks) -> kj::ArrayPtr<const kj::byte> { return chunk.asPtr(); };
3571 co_await sink->write(pieces);
3572 }
3573 
3574 // If the stream is done, end the output if needed and exit.
3575 if (result.done) {
3576 KJ_ON_SCOPE_FAILURE(writeFailed = true);
3577 if (end) {
3578 co_await sink->end();
3579 }
3580 co_return;
3581 }
3582 }
3583 }
3584 KJ_CATCH(exception) {
3585 if (!writeFailed) {
3586 sink->abort(exception.clone());
3587 }
3588 
3589 co_await ioContext.run([&reader, ex = exception.clone()](jsg::Lock& js) mutable {
3590 auto& ioContext = IoContext::current();
3591 auto error = js.exceptionToJsValue(kj::mv(ex));
3592 return ioContext.awaitJs(js, reader->cancel(js, error.getHandle(js)));
3593 });
3594 kj::throwFatalException(kj::mv(exception));
3595 }
3596}
3597} // namespace
3598 
3599template <typename T>
3600jsg::Promise<T> ReadableStreamJsController::readAll(jsg::Lock& js, uint64_t limit) {
3601 if (isLockedToReader()) {
3602 return js.rejectedPromise<T>(KJ_EXCEPTION(
3603 FAILED, "jsg.TypeError: This ReadableStream is currently locked to a reader."));
3604 }
3605 disturbed = true;
3606 
3607 bool stripBom = false;
3608 KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) {
3609 stripBom = flags.getStripBomInReadAllText();
3610 }
3611 
3612 // This operation leaves the stream locked and disturbed. The loop will read until
3613 // the stream is closed or errored. If the limit is reached, the loop will error.
3614 
3615 const auto readAll = [this, limit, stripBom](auto& js) -> jsg::Promise<T> {
3616 KJ_ASSERT(lock.lock());
3617 // The AllReader will hold a traceable reference to the ReadableStream.
3618 auto reader = kj::heap<AllReader>(addRef(), limit);
3619 
3620 auto promise = ([&js, &reader, stripBom]() -> jsg::Promise<T> {
3621 if constexpr (kj::isSameType<T, jsg::BufferSource>()) {
3622 (void)stripBom; // Unused in this branch.
3623 return reader->allBytes(js);
3624 } else {
3625 auto option = ReadAllTextOption::NULL_TERMINATE;
3626 if (stripBom) {
3627 option |= ReadAllTextOption::STRIP_BOM;
3628 }
3629 return reader->allText(js, option);
3630 }
3631 })();
3632 
3633 return maybeAddFunctor(js, kj::mv(promise),
3634 // reader is a GC visitable type that holds a reference to either the stream
3635 // or an error. Accordingly, we wrap it in a visitable lambda attached as a
3636 // continuation on the promise to ensure that it is GC visited and kept alive until
3637 // the promise settles.
3638 JSG_VISITABLE_LAMBDA((reader = kj::mv(reader)), (reader),
3639 (jsg::Lock & js, T result)->jsg::Promise<T> {
3640 return js.resolvedPromise(kj::mv(result));
3641 }),
3642 [](jsg::Lock& js, jsg::Value exception) -> jsg::Promise<T> {
3643 return js.rejectedPromise<T>(kj::mv(exception));
3644 });
3645 };
3646 
3647 KJ_SWITCH_ONEOF(state) {
3648 KJ_CASE_ONEOF(initial, Initial) {
3649 // Stream not yet set up, treat as closed.
3650 if constexpr (kj::isSameType<T, jsg::BufferSource>()) {
3651 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0);
3652 return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing)));
3653 } else {
3654 return js.resolvedPromise(T());
3655 }
3656 }
3657 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3658 if constexpr (kj::isSameType<T, jsg::BufferSource>()) {
3659 auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0);
3660 return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing)));
3661 } else {
3662 return js.resolvedPromise(T());
3663 }
3664 }
3665 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3666 return js.rejectedPromise<T>(errored.addRef(js));
3667 }
3668 KJ_CASE_ONEOF(valueReadable, kj::Own<ValueReadable>) {
3669 return readAll(js);
3670 }
3671 KJ_CASE_ONEOF(byteReadable, kj::Own<ByteReadable>) {
3672 return readAll(js);
3673 }
3674 }
3675 KJ_UNREACHABLE;
3676}
3677 
3678jsg::Promise<jsg::BufferSource> ReadableStreamJsController::readAllBytes(
3679 jsg::Lock& js, uint64_t limit) {
3680 return readAll<jsg::BufferSource>(js, limit);
3681}
3682 
3683jsg::Promise<kj::String> ReadableStreamJsController::readAllText(jsg::Lock& js, uint64_t limit) {
3684 return readAll<kj::String>(js, limit);
3685}
3686 
3687kj::Own<ReadableStreamController> ReadableStreamJsController::detach(
3688 jsg::Lock& js, bool ignored /* unused */) {
3689 KJ_ASSERT(!isLockedToReader());
3690 KJ_ASSERT(!isDisturbed());
3691 KJ_ASSERT(!state.hasOperationInProgress(), "Unable to detach with read pending");
3692 auto controller = kj::heap<ReadableStreamJsController>();
3693 controller->expectedLength = expectedLength;
3694 disturbed = true;
3695 
3696 // Clones this streams state into a new ReadableStreamController, leaving this stream
3697 // locked, disturbed, and closed.
3698 
3699 // The controller starts in Initial state by default, so we can use regular transitionTo.
3700 KJ_SWITCH_ONEOF(state) {
3701 KJ_CASE_ONEOF(initial, Initial) {
3702 // Still in initial state, transition to closed
3703 controller->state.transitionTo<StreamStates::Closed>();
3704 }
3705 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3706 controller->state.transitionTo<StreamStates::Closed>();
3707 }
3708 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3709 controller->state.transitionTo<StreamStates::Errored>(errored.addRef(js));
3710 }
3711 KJ_CASE_ONEOF(readable, kj::Own<ValueReadable>) {
3712 KJ_ASSERT(lock.lock());
3713 controller->state.transitionTo<kj::Own<ValueReadable>>(readable->clone(js, *controller));
3714 state.transitionTo<StreamStates::Closed>();
3715 lock.onClose(js);
3716 }
3717 KJ_CASE_ONEOF(readable, kj::Own<ByteReadable>) {
3718 KJ_ASSERT(lock.lock());
3719 controller->state.transitionTo<kj::Own<ByteReadable>>(readable->clone(js, *controller));
3720 state.transitionTo<StreamStates::Closed>();
3721 lock.onClose(js);
3722 }
3723 }
3724 
3725 return kj::mv(controller);
3726}
3727 
3728kj::Maybe<uint64_t> ReadableStreamJsController::tryGetLength(StreamEncoding encoding) {
3729 return expectedLength;
3730}
3731 
3732kj::Promise<DeferredProxy<void>> ReadableStreamJsController::pumpTo(
3733 jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) {
3734 KJ_ASSERT(IoContext::hasCurrent(), "Unable to consume this ReadableStream outside of a request");
3735 KJ_REQUIRE(!isLockedToReader(), "This ReadableStream is currently locked to a reader.");
3736 disturbed = true;
3737 
3738 // This operation will leave the ReadableStream locked and disturbed. It will consume
3739 // the stream until it either closed or errors.
3740 //
3741 // When the ENABLE_DRAINING_READ_ON_STANDARD_STREAMS autogate is enabled, uses the new
3742 // pumpToImpl coroutine with DrainingReader for batched reads and vectored writes.
3743 // Otherwise, falls back to the original PumpToReader JS promise loop that reads one
3744 // chunk at a time.
3745 
3746 const auto handlePump = [&] {
3747 if (util::Autogate::isEnabled(util::AutogateKey::ENABLE_DRAINING_READ_ON_STANDARD_STREAMS)) {
3748 auto reader = KJ_ASSERT_NONNULL(DrainingReader::create(js, *this->addRef()),
3749 "Failed to create DrainingReader โ€” stream should not be locked");
3750 auto& ioContext = IoContext::current();
3751 return addNoopDeferredProxy(pumpToImpl(ioContext, kj::mv(reader), kj::mv(sink), end));
3752 } else {
3753 KJ_ASSERT(lock.lock());
3754 auto reader = kj::heap<PumpToReader>(addRef(), kj::mv(sink), end);
3755 return addNoopDeferredProxy(reader->pumpTo(js).attach(kj::mv(reader)));
3756 }
3757 };
3758 
3759 KJ_SWITCH_ONEOF(state) {
3760 KJ_CASE_ONEOF(initial, Initial) {
3761 // Stream not yet set up, treat as closed.
3762 return addNoopDeferredProxy(sink->end().attach(kj::mv(sink)));
3763 }
3764 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3765 return addNoopDeferredProxy(sink->end().attach(kj::mv(sink)));
3766 }
3767 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3768 return js.exceptionToKj(errored.addRef(js));
3769 }
3770 KJ_CASE_ONEOF(readable, kj::Own<ValueReadable>) {
3771 return handlePump();
3772 }
3773 KJ_CASE_ONEOF(readable, kj::Own<ByteReadable>) {
3774 return handlePump();
3775 }
3776 }
3777 
3778 KJ_UNREACHABLE;
3779}
3780 
3781// ======================================================================================
3782 
3783WritableStreamDefaultController::WritableStreamDefaultController(
3784 jsg::Lock& js, WritableStream& owner, jsg::Ref<AbortSignal> abortSignal)
3785 : ioContext(tryGetIoContext()),
3786 impl(js, owner, kj::mv(abortSignal)) {}
3787 
3788jsg::Promise<void> WritableStreamDefaultController::abort(
3789 jsg::Lock& js, v8::Local<v8::Value> reason) {
3790 return impl.abort(js, JSG_THIS, reason);
3791}
3792 
3793void WritableStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) {
3794 visitor.visit(impl);
3795}
3796 
3797jsg::Promise<void> WritableStreamDefaultController::close(jsg::Lock& js) {
3798 return impl.close(js, JSG_THIS);
3799}
3800 
3801void WritableStreamDefaultController::error(
3802 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) {
3803 impl.error(js, JSG_THIS, reason.orDefault(js.undefined()));
3804}
3805 
3806kj::Maybe<ssize_t> WritableStreamDefaultController::getDesiredSize() {
3807 // Per the spec, desiredSize should be null when the stream is erroring.
3808 if (impl.flags.pedanticWpt && isErroring()) {
3809 return kj::none;
3810 }
3811 return impl.getDesiredSize();
3812}
3813 
3814jsg::Ref<AbortSignal> WritableStreamDefaultController::getSignal() {
3815 return impl.signal.addRef();
3816}
3817 
3818kj::Maybe<v8::Local<v8::Value>> WritableStreamDefaultController::isErroring(jsg::Lock& js) {
3819 KJ_IF_SOME(erroring, impl.state.tryGetUnsafe<StreamStates::Erroring>()) {
3820 return erroring.reason.getHandle(js);
3821 }
3822 return kj::none;
3823}
3824 
3825void WritableStreamDefaultController::setup(
3826 jsg::Lock& js, UnderlyingSink underlyingSink, StreamQueuingStrategy queuingStrategy) {
3827 impl.setup(js, JSG_THIS, kj::mv(underlyingSink), kj::mv(queuingStrategy));
3828}
3829 
3830jsg::Promise<void> WritableStreamDefaultController::write(
3831 jsg::Lock& js, v8::Local<v8::Value> value) {
3832 return impl.write(js, JSG_THIS, value);
3833}
3834 
3835void WritableStreamDefaultController::cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason) {
3836 impl.cancelPendingWrites(js, reason);
3837}
3838 
3839void WritableStreamDefaultController::clearAlgorithms() {
3840 impl.algorithms.clear();
3841}
3842 
3843WritableStreamDefaultController::~WritableStreamDefaultController() noexcept(false) {
3844 // Clear algorithms in destructor to break circular references
3845 clearAlgorithms();
3846}
3847 
3848// ======================================================================================
3849WritableStreamJsController::WritableStreamJsController(): ioContext(tryGetIoContext()) {}
3850 
3851WritableStreamJsController::~WritableStreamJsController() noexcept(false) {
3852 // Clear algorithms to break circular references during destruction
3853 KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) {
3854 controller->clearAlgorithms();
3855 }
3856 // Clear the state to break the circular reference to the controller.
3857 // During destruction, we force the transition since the current state doesn't matter.
3858 state.forceTransitionTo<StreamStates::Closed>();
3859 // Clear owner reference
3860 owner = kj::none;
3861 // Clear any pending abort promise
3862 maybeAbortPromise = kj::none;
3863}
3864 
3865WritableStreamJsController::WritableStreamJsController(StreamStates::Closed closed)
3866 : ioContext(tryGetIoContext()) {
3867 state.transitionTo<StreamStates::Closed>();
3868}
3869 
3870WritableStreamJsController::WritableStreamJsController(StreamStates::Errored errored)
3871 : ioContext(tryGetIoContext()) {
3872 state.transitionTo<StreamStates::Errored>(kj::mv(errored));
3873}
3874 
3875jsg::Promise<void> WritableStreamJsController::abort(
3876 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) {
3877 // The spec requires that if abort is called multiple times, it is supposed to return the same
3878 // promise each time. That's a bit cumbersome here with jsg::Promise so we intentionally just
3879 // return a continuation branch off the same promise.
3880 KJ_IF_SOME(abortPromise, maybeAbortPromise) {
3881 return abortPromise.whenResolved(js);
3882 }
3883 KJ_SWITCH_ONEOF(state) {
3884 KJ_CASE_ONEOF(initial, Initial) {
3885 // Stream hasn't been set up yet - treat like closed for abort purposes
3886 maybeAbortPromise = js.resolvedPromise();
3887 return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js);
3888 }
3889 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3890 maybeAbortPromise = js.resolvedPromise();
3891 return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js);
3892 }
3893 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3894 // Per the spec, if the stream is errored, we are to return a resolved promise.
3895 maybeAbortPromise = js.resolvedPromise();
3896 return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js);
3897 }
3898 KJ_CASE_ONEOF(controller, Controller) {
3899 maybeAbortPromise = controller->abort(js, reason.orDefault(js.undefined()));
3900 return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js);
3901 }
3902 }
3903 KJ_UNREACHABLE;
3904}
3905 
3906jsg::Ref<WritableStream> WritableStreamJsController::addRef() {
3907 return KJ_ASSERT_NONNULL(owner).addRef();
3908}
3909 
3910bool WritableStreamJsController::isClosedOrClosing() {
3911 return state.is<StreamStates::Closed>();
3912}
3913 
3914bool WritableStreamJsController::isErrored() {
3915 return state.isErrored();
3916}
3917 
3918jsg::Promise<void> WritableStreamJsController::close(jsg::Lock& js, bool markAsHandled) {
3919 KJ_SWITCH_ONEOF(state) {
3920 KJ_CASE_ONEOF(initial, Initial) {
3921 return rejectedMaybeHandledPromise<void>(
3922 js, js.v8TypeError("This WritableStream has been closed."_kj), markAsHandled);
3923 }
3924 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
3925 return rejectedMaybeHandledPromise<void>(
3926 js, js.v8TypeError("This WritableStream has been closed."_kj), markAsHandled);
3927 }
3928 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
3929 if (FeatureFlags::get(js).getPedanticWpt()) {
3930 return rejectedMaybeHandledPromise<void>(
3931 js, js.v8TypeError("This WritableStream has been errored."_kj), markAsHandled);
3932 }
3933 return rejectedMaybeHandledPromise<void>(js, errored.getHandle(js), markAsHandled);
3934 }
3935 KJ_CASE_ONEOF(controller, Controller) {
3936 return controller->close(js);
3937 }
3938 }
3939 KJ_UNREACHABLE;
3940}
3941 
3942void WritableStreamJsController::doClose(jsg::Lock& js) {
3943 // If already in a terminal state, nothing to do.
3944 if (state.isTerminal()) return;
3945 
3946 // Clear algorithms to break circular references before changing state
3947 KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) {
3948 controller->clearAlgorithms();
3949 }
3950 
3951 state.transitionTo<StreamStates::Closed>();
3952 KJ_IF_SOME(locked, lock.state.tryGetUnsafe<WriterLocked>()) {
3953 maybeResolvePromise(js, locked.getClosedFulfiller());
3954 maybeResolvePromise(js, locked.getReadyFulfiller());
3955 } else {
3956 (void)lock.state.transitionFromTo<WritableLockImpl::PipeLocked, Unlocked>();
3957 }
3958}
3959 
3960void WritableStreamJsController::doError(jsg::Lock& js, v8::Local<v8::Value> reason) {
3961 // If already in a terminal state, nothing to do.
3962 if (state.isTerminal()) return;
3963 
3964 // Clear algorithms to break circular references before changing state
3965 KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) {
3966 controller->clearAlgorithms();
3967 }
3968 
3969 state.transitionTo<StreamStates::Errored>(js.v8Ref(reason));
3970 KJ_IF_SOME(locked, lock.state.tryGetUnsafe<WriterLocked>()) {
3971 maybeRejectPromise<void>(js, locked.getClosedFulfiller(), reason);
3972 maybeResolvePromise(js, locked.getReadyFulfiller());
3973 } else KJ_IF_SOME(pipeLocked, lock.state.tryGetUnsafe<WritableLockImpl::PipeLocked>()) {
3974 // When the writable side of a pipe errors, we need to release the source stream.
3975 // The pipeLoop may be waiting on a read from the source that will never complete,
3976 // so we need to proactively release the source here.
3977 if (!pipeLocked.flags.preventCancel) {
3978 pipeLocked.source.release(js, reason);
3979 } else {
3980 pipeLocked.source.release(js);
3981 }
3982 lock.state.transitionTo<Unlocked>();
3983 }
3984}
3985 
3986void WritableStreamJsController::errorIfNeeded(jsg::Lock& js, v8::Local<v8::Value> reason) {
3987 // Error through the underlying controller if available, which goes through the proper
3988 // error transition (Erroring -> Errored). This allows close() to be called while the
3989 // stream is "erroring" and reject with the stored error.
3990 KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) {
3991 controller->error(js, reason);
3992 }
3993 // If state is not Controller (already Closed or Errored), this is a no-op.
3994}
3995 
3996kj::Maybe<int> WritableStreamJsController::getDesiredSize() {
3997 KJ_SWITCH_ONEOF(state) {
3998 KJ_CASE_ONEOF(initial, Initial) {
3999 return 0;
4000 }
4001 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
4002 return 0;
4003 }
4004 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
4005 return kj::none;
4006 }
4007 KJ_CASE_ONEOF(controller, Controller) {
4008 return controller->getDesiredSize().map([](ssize_t size) -> int { return size; });
4009 }
4010 }
4011 KJ_UNREACHABLE;
4012}
4013 
4014kj::Maybe<v8::Local<v8::Value>> WritableStreamJsController::isErroring(jsg::Lock& js) {
4015 KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) {
4016 return controller->isErroring(js);
4017 }
4018 return kj::none;
4019}
4020 
4021bool WritableStreamDefaultController::isErroring() const {
4022 return impl.state.is<StreamStates::Erroring>();
4023}
4024 
4025kj::Maybe<v8::Local<v8::Value>> WritableStreamJsController::isErroredOrErroring(jsg::Lock& js) {
4026 KJ_IF_SOME(err, state.tryGetErrorUnsafe()) {
4027 return err.getHandle(js);
4028 }
4029 return isErroring(js);
4030}
4031 
4032bool WritableStreamJsController::isStarted() {
4033 KJ_SWITCH_ONEOF(state) {
4034 KJ_CASE_ONEOF(initial, Initial) {
4035 return false;
4036 }
4037 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
4038 return true;
4039 }
4040 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
4041 return true;
4042 }
4043 KJ_CASE_ONEOF(controller, Controller) {
4044 return controller->isStarted();
4045 }
4046 }
4047 KJ_UNREACHABLE;
4048}
4049 
4050bool WritableStreamJsController::hasBackpressure() {
4051 KJ_IF_SOME(controller, state.tryGetUnsafe<Controller>()) {
4052 return controller->hasBackpressure();
4053 }
4054 return false;
4055}
4056 
4057bool WritableStreamJsController::isLocked() const {
4058 return isLockedToWriter();
4059}
4060 
4061bool WritableStreamJsController::isLockedToWriter() const {
4062 return !lock.state.is<Unlocked>();
4063}
4064 
4065bool WritableStreamJsController::lockWriter(jsg::Lock& js, Writer& writer) {
4066 return lock.lockWriter(js, *this, writer);
4067}
4068 
4069void WritableStreamJsController::maybeRejectReadyPromise(
4070 jsg::Lock& js, v8::Local<v8::Value> reason) {
4071 KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe<WriterLocked>()) {
4072 if (writerLock.getReadyFulfiller() != kj::none) {
4073 maybeRejectPromise<void>(js, writerLock.getReadyFulfiller(), reason);
4074 } else {
4075 auto prp = js.newPromiseAndResolver<void>();
4076 prp.promise.markAsHandled(js);
4077 prp.resolver.reject(js, reason);
4078 writerLock.setReadyFulfiller(js, prp);
4079 }
4080 }
4081}
4082 
4083void WritableStreamJsController::maybeResolveReadyPromise(jsg::Lock& js) {
4084 KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe<WriterLocked>()) {
4085 maybeResolvePromise(js, writerLock.getReadyFulfiller());
4086 }
4087}
4088 
4089void WritableStreamJsController::releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) {
4090 lock.releaseWriter(*this, writer, maybeJs);
4091}
4092 
4093kj::Maybe<kj::Own<WritableStreamSink>> WritableStreamJsController::removeSink(jsg::Lock& js) {
4094 return kj::none;
4095}
4096void WritableStreamJsController::detach(jsg::Lock& js) {
4097 KJ_UNIMPLEMENTED("WritableStreamJsController::detach is not implemented");
4098}
4099 
4100void WritableStreamJsController::setOwnerRef(WritableStream& stream) {
4101 owner = stream;
4102}
4103 
4104void WritableStreamJsController::setup(jsg::Lock& js,
4105 jsg::Optional<UnderlyingSink> maybeUnderlyingSink,
4106 jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) {
4107 auto underlyingSink = kj::mv(maybeUnderlyingSink).orDefault({});
4108 auto queuingStrategy = kj::mv(maybeQueuingStrategy).orDefault({});
4109 
4110 if (FeatureFlags::get(js).getPedanticWpt()) {
4111 // Per the spec, the type property for WritableStream's underlying sink must be undefined.
4112 // If it's anything else, throw a RangeError.
4113 JSG_REQUIRE(underlyingSink.type == kj::none, RangeError,
4114 "Invalid underlying sink type. Only undefined is valid.");
4115 }
4116 
4117 // We account for the memory usage of the WritableStreamDefaultController and AbortSignal together
4118 // because their lifetimes are identical and memory accounting itself has a memory overhead.
4119 auto controller = js.allocAccounted<WritableStreamDefaultController>(
4120 sizeof(WritableStreamDefaultController) + sizeof(AbortSignal), js, KJ_ASSERT_NONNULL(owner),
4121 js.alloc<AbortSignal>());
4122 auto& controllerRef = *controller;
4123 state.transitionTo<Controller>(kj::mv(controller));
4124 controllerRef.setup(js, kj::mv(underlyingSink), kj::mv(queuingStrategy));
4125}
4126 
4127kj::Maybe<jsg::Promise<void>> WritableStreamJsController::tryPipeFrom(
4128 jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) {
4129 JSG_REQUIRE_NONNULL(
4130 ioContext, Error, "Unable to pipe to a WritableStream created outside of a request");
4131 
4132 // The ReadableStream source here can be either a JavaScript-backed ReadableStream
4133 // or ReadableStreamSource-backed. In either case, however, this WritableStream is
4134 // JavaScript-based and must use a JavaScript promise-based data flow for piping data.
4135 // We'll treat all ReadableStreams as if they are JavaScript-backed.
4136 //
4137 // This method will return a JavaScript promise that is resolved when the pipe operation
4138 // completes, or is rejected if the pipe operation is aborted or errored.
4139 
4140 // Let's also acquire the destination pipe lock.
4141 lock.pipeLock(KJ_ASSERT_NONNULL(owner), kj::mv(source), options);
4142 
4143 return pipeLoop(js).then(js, JSG_VISITABLE_LAMBDA((ref = addRef()), (ref), (auto& js){}));
4144}
4145 
4146jsg::Promise<void> WritableStreamJsController::pipeLoop(jsg::Lock& js) {
4147 auto maybePipeLock = lock.tryGetPipe();
4148 if (maybePipeLock == kj::none) return js.resolvedPromise();
4149 auto& pipeLock = KJ_REQUIRE_NONNULL(maybePipeLock);
4150 
4151 auto preventAbort = pipeLock.flags.preventAbort;
4152 auto preventCancel = pipeLock.flags.preventCancel;
4153 auto preventClose = pipeLock.flags.preventClose;
4154 auto pipeThrough = pipeLock.flags.pipeThrough;
4155 auto& source = pipeLock.source;
4156 // At the start of each pipe step, we check to see if either the source or
4157 // the destination has closed or errored and propagate that on to the other.
4158 KJ_IF_SOME(promise, pipeLock.checkSignal(js, *this)) {
4159 lock.releasePipeLock();
4160 return kj::mv(promise);
4161 }
4162 
4163 KJ_IF_SOME(errored, pipeLock.source.tryGetErrored(js)) {
4164 source.release(js);
4165 lock.releasePipeLock();
4166 if (!preventAbort) {
4167 auto onSuccess = JSG_VISITABLE_LAMBDA(
4168 (pipeThrough, reason = js.v8Ref(errored)), (reason), (jsg::Lock& js) {
4169 return rejectedMaybeHandledPromise<void>(js, reason.getHandle(js), pipeThrough);
4170 });
4171 auto promise = abort(js, errored);
4172 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
4173 return promise.then(js, ioContext.addFunctor(kj::mv(onSuccess)));
4174 } else {
4175 return promise.then(js, kj::mv(onSuccess));
4176 }
4177 }
4178 return rejectedMaybeHandledPromise<void>(js, errored, pipeThrough);
4179 }
4180 
4181 KJ_IF_SOME(errored, state.tryGetUnsafe<StreamStates::Errored>()) {
4182 lock.releasePipeLock();
4183 auto reason = errored.getHandle(js);
4184 if (!preventCancel) {
4185 source.release(js, reason);
4186 } else {
4187 source.release(js);
4188 }
4189 return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough);
4190 }
4191 
4192 KJ_IF_SOME(erroring, isErroring(js)) {
4193 lock.releasePipeLock();
4194 if (!preventCancel) {
4195 source.release(js, erroring);
4196 } else {
4197 source.release(js);
4198 }
4199 return rejectedMaybeHandledPromise<void>(js, erroring, pipeThrough);
4200 }
4201 
4202 if (source.isClosed()) {
4203 source.release(js);
4204 lock.releasePipeLock();
4205 if (!preventClose) {
4206 auto promise = close(js);
4207 if (pipeThrough) {
4208 promise.markAsHandled(js);
4209 }
4210 return kj::mv(promise);
4211 }
4212 return js.resolvedPromise();
4213 }
4214 
4215 if (state.is<StreamStates::Closed>()) {
4216 lock.releasePipeLock();
4217 auto reason = js.v8TypeError("This destination writable stream is closed."_kj);
4218 if (!preventCancel) {
4219 source.release(js, reason);
4220 } else {
4221 source.release(js);
4222 }
4223 
4224 return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough);
4225 }
4226 
4227 // Assuming we get by that, we perform a read on the source. If the read errors,
4228 // we propagate the error to the destination, depending on options and reject
4229 // the pipe promise. If the read is successful then we'll get a ReadResult
4230 // back. If the ReadResult indicates done, then we close the destination
4231 // depending on options and resolve the pipe promise. If the ReadResult is
4232 // not done, we write the value on to the destination. If the write operation
4233 // fails, we reject the pipe promise and propagate the error back to the
4234 // source (again, depending on options). If the write operation is successful,
4235 // we call pipeLoop again to move on to the next iteration.
4236 
4237 auto onSuccess = JSG_VISITABLE_LAMBDA((this, ref = addRef(), preventCancel, pipeThrough), (ref),
4238 (jsg::Lock & js, ReadResult result)->jsg::Promise<void> {
4239 auto maybePipeLock = lock.tryGetPipe();
4240 if (maybePipeLock == kj::none) return js.resolvedPromise();
4241 auto& pipeLock = KJ_REQUIRE_NONNULL(maybePipeLock);
4242 
4243 KJ_IF_SOME(promise, pipeLock.checkSignal(js, *this)) {
4244 lock.releasePipeLock();
4245 return kj::mv(promise);
4246 } else {
4247 } // Trailing else() is squash compiler warning
4248 
4249 if (result.done) {
4250 // We'll handle the close at the start of the next iteration.
4251 return pipeLoop(js);
4252 }
4253 
4254 auto onSuccess = JSG_VISITABLE_LAMBDA(
4255 (this, ref=addRef()), (ref) , (jsg::Lock& js) {
4256 return pipeLoop(js);
4257 } );
4258 
4259 auto onFailure = JSG_VISITABLE_LAMBDA(
4260 (this, ref=addRef(), preventCancel, pipeThrough),
4261 (ref) , (jsg::Lock& js, jsg::Value value) {
4262 // The write failed. We need to release the source if the pipe lock still exists.
4263 auto reason = value.getHandle(js);
4264 KJ_IF_SOME(pipeLock, lock.tryGetPipe()) {
4265 if (!preventCancel) {
4266 pipeLock.source.release(js, reason);
4267 } else {
4268 pipeLock.source.release(js);
4269 }
4270 } else {} // Trailing else() to squash compiler warning
4271 return rejectedMaybeHandledPromise<void>(js, reason, pipeThrough);
4272 } );
4273 
4274 auto promise =
4275 write(js, result.value.map([&](jsg::Value& value) { return value.getHandle(js); }));
4276 
4277 return maybeAddFunctor(js, kj::mv(promise), kj::mv(onSuccess), kj::mv(onFailure));
4278 });
4279 
4280 auto onFailure =
4281 JSG_VISITABLE_LAMBDA((this, ref = addRef()), (ref), (jsg::Lock& js, jsg::Value value) {
4282 // The read failed. We will handle the error at the start of the next iteration.
4283 return pipeLoop(js);
4284 });
4285 
4286 return maybeAddFunctor(js, pipeLock.source.read(js), kj::mv(onSuccess), kj::mv(onFailure));
4287}
4288 
4289void WritableStreamJsController::updateBackpressure(jsg::Lock& js, bool backpressure) {
4290 KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe<WriterLocked>()) {
4291 if (backpressure) {
4292 // Per the spec, when backpressure is updated and is true, we replace the existing
4293 // ready promise on the writer with a new pending promise, regardless of whether
4294 // the existing one is resolved or not.
4295 auto prp = js.newPromiseAndResolver<void>();
4296 prp.promise.markAsHandled(js);
4297 return writerLock.setReadyFulfiller(js, prp);
4298 }
4299 
4300 // When backpressure is updated and is false, we resolve the ready promise on the writer
4301 maybeResolvePromise(js, writerLock.getReadyFulfiller());
4302 }
4303}
4304 
4305jsg::Promise<void> WritableStreamJsController::write(
4306 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) {
4307 KJ_SWITCH_ONEOF(state) {
4308 KJ_CASE_ONEOF(initial, Initial) {
4309 return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj));
4310 }
4311 KJ_CASE_ONEOF(closed, StreamStates::Closed) {
4312 return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj));
4313 }
4314 KJ_CASE_ONEOF(errored, StreamStates::Errored) {
4315 return js.rejectedPromise<void>(errored.addRef(js));
4316 }
4317 KJ_CASE_ONEOF(controller, Controller) {
4318 return controller->write(js, value.orDefault([&] { return js.undefined(); }));
4319 }
4320 }
4321 KJ_UNREACHABLE;
4322}
4323 
4324void WritableStreamJsController::visitForGc(jsg::GcVisitor& visitor) {
4325 state.visitForGc(visitor);
4326 visitor.visit(maybeAbortPromise, lock);
4327}
4328 
4329// =======================================================================================
4330 
4331TransformStreamDefaultController::TransformStreamDefaultController(jsg::Lock& js)
4332 : ioContext(tryGetIoContext()),
4333 startPromise(js.newPromiseAndResolver<void>()) {}
4334 
4335kj::Maybe<int> TransformStreamDefaultController::getDesiredSize() {
4336 KJ_IF_SOME(readableController, tryGetReadableController()) {
4337 return readableController.getDesiredSize();
4338 }
4339 return kj::none;
4340}
4341 
4342void TransformStreamDefaultController::enqueue(jsg::Lock& js, v8::Local<v8::Value> chunk) {
4343 auto& readableController = JSG_REQUIRE_NONNULL(tryGetReadableController(), TypeError,
4344 "The readable side of this TransformStream is no longer readable.");
4345 // Hold a strong reference to the readable controller for the duration of this
4346 // method. The readableController.enqueue() call below invokes the user-provided
4347 // size algorithm, which can re-enter JS and call error() on this transform
4348 // controller, dropping the jsg::Ref held by this->readable and the one held by
4349 // the ReadableStreamJsController's ValueReadable. Without this ref the
4350 // ReadableStreamDefaultController would be freed while its enqueue() method is
4351 // still on the stack.
4352 auto readableControllerRef = kj::addRef(readableController);
4353 
4354 JSG_REQUIRE(readableController.canCloseOrEnqueue(), TypeError,
4355 "The readable side of this TransformStream is no longer readable.");
4356 js.tryCatch([&] { readableController.enqueue(js, chunk); }, [&](jsg::Value exception) {
4357 errorWritableAndUnblockWrite(js, exception.getHandle(js));
4358 js.throwException(kj::mv(exception));
4359 });
4360 
4361 // If the controller was errored during the enqueue (e.g. by the size callback
4362 // calling error()), skip the backpressure update โ€” the stream is already torn down.
4363 if (!readableController.canCloseOrEnqueue()) {
4364 return;
4365 }
4366 
4367 bool newBackpressure = readableController.hasBackpressure();
4368 if (newBackpressure != backpressure) {
4369 KJ_ASSERT(newBackpressure);
4370 // Unfortunately the original implementation forgot to actually set the backpressure
4371 // here so the backpressure signaling failed to work correctly. This is unfortunate
4372 // because applying the backpressure here could break existing code, so we need to
4373 // put the fix behind a compat flag. Doh!
4374 if (FeatureFlags::get(js).getFixupTransformStreamBackpressure()) {
4375 setBackpressure(js, true);
4376 }
4377 }
4378}
4379 
4380void TransformStreamDefaultController::error(jsg::Lock& js, v8::Local<v8::Value> reason) {
4381 KJ_IF_SOME(readableController, tryGetReadableController()) {
4382 readableController.error(js, reason);
4383 readable = kj::none;
4384 }
4385 errorWritableAndUnblockWrite(js, reason);
4386}
4387 
4388void TransformStreamDefaultController::terminate(jsg::Lock& js) {
4389 KJ_IF_SOME(readableController, tryGetReadableController()) {
4390 readableController.close(js);
4391 readable = kj::none;
4392 }
4393 errorWritableAndUnblockWrite(js, js.v8TypeError("The transform stream has been terminated"_kj));
4394}
4395 
4396jsg::Promise<void> TransformStreamDefaultController::write(
4397 jsg::Lock& js, v8::Local<v8::Value> chunk) {
4398 KJ_IF_SOME(writableController, tryGetWritableController()) {
4399 KJ_IF_SOME(error, writableController.isErroredOrErroring(js)) {
4400 return js.rejectedPromise<void>(error);
4401 }
4402 
4403 KJ_ASSERT(writableController.isWritable());
4404 
4405 if (backpressure) {
4406 auto chunkRef = js.v8Ref(chunk);
4407 return KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.whenResolved(js).then(js,
4408 JSG_VISITABLE_LAMBDA((chunkRef = kj::mv(chunkRef), ref=JSG_THIS),
4409 (chunkRef, ref), (jsg::Lock& js) mutable -> jsg::Promise<void> {
4410 KJ_IF_SOME(writableController, ref->tryGetWritableController()) {
4411 KJ_IF_SOME(error, writableController.isErroring(js)) {
4412 return js.rejectedPromise<void>(error);
4413 } else {
4414 // Else block to avert dangling else compiler warning.
4415 }
4416 } else {
4417 // Else block to avert dangling else compiler warning.
4418 }
4419 return ref->performTransform(js, chunkRef.getHandle(js));
4420 }));
4421 }
4422 return performTransform(js, chunk);
4423 } else {
4424 return js.rejectedPromise<void>(
4425 KJ_EXCEPTION(FAILED, "jsg.TypeError: Writing to the TransformStream failed."));
4426 }
4427}
4428 
4429jsg::Promise<void> TransformStreamDefaultController::abort(
4430 jsg::Lock& js, v8::Local<v8::Value> reason) {
4431 if (FeatureFlags::get(js).getPedanticWpt()) {
4432 // If a finish operation is already in progress, return the existing promise
4433 // or handle the case where we're being called synchronously from within another
4434 // finish operation.
4435 if (algorithms.finishStarted) {
4436 KJ_IF_SOME(finish, algorithms.maybeFinish) {
4437 return finish.whenResolved(js);
4438 }
4439 // finishStarted is true but maybeFinish is not set yet - this means we're being
4440 // called synchronously from within another finish operation (like cancel).
4441 // We need to error the stream with the abort reason so that both the current
4442 // operation and this abort reject with the abort reason.
4443 error(js, reason);
4444 return js.rejectedPromise<void>(js.v8Ref(reason));
4445 }
4446 
4447 // Mark that we're starting a finish operation before running the algorithm.
4448 algorithms.finishStarted = true;
4449 } else {
4450 KJ_IF_SOME(finish, algorithms.maybeFinish) {
4451 return finish.whenResolved(js);
4452 }
4453 }
4454 
4455 return algorithms.maybeFinish
4456 .emplace(maybeRunAlgorithm(js, algorithms.cancel,
4457 JSG_VISITABLE_LAMBDA(
4458 (this, ref = JSG_THIS, reason = jsg::JsRef(js, jsg::JsValue(reason))), (ref, reason),
4459 (jsg::Lock & js)->jsg::Promise<void> {
4460 // If the readable side is errored, return a rejected promise with the stored error
4461 {
4462 KJ_IF_SOME(err, getReadableErrorState(js)) {
4463 return js.rejectedPromise<void>(kj::mv(err));
4464 } else {
4465 // Else block to avert dangling else compiler warning.
4466 }
4467 }
4468 // Otherwise... error with the given reason and resolve the abort promise
4469 error(js, reason.getHandle(js));
4470 return js.resolvedPromise();
4471 }),
4472 JSG_VISITABLE_LAMBDA((this, ref = JSG_THIS), (ref),
4473 (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> {
4474 error(js, reason.getHandle(js));
4475 return js.rejectedPromise<void>(kj::mv(reason));
4476 }),
4477 jsg::JsValue(reason)))
4478 .whenResolved(js);
4479}
4480 
4481jsg::Promise<void> TransformStreamDefaultController::close(jsg::Lock& js) {
4482 auto flags = FeatureFlags::get(js);
4483 if (flags.getPedanticWpt()) {
4484 // If a finish operation is already in progress (e.g., from cancel or abort),
4485 // we should not run flush. Per the WHATWG streams spec, close/flush should
4486 // coordinate with cancel to avoid calling both.
4487 if (algorithms.finishStarted) {
4488 KJ_IF_SOME(finish, algorithms.maybeFinish) {
4489 return finish.whenResolved(js);
4490 }
4491 // finishStarted is true but maybeFinish is not set yet - this means we're being
4492 // called synchronously from within another finish operation. If the stream was
4493 // errored during that operation, return a rejected promise with the error.
4494 KJ_IF_SOME(writableController, tryGetWritableController()) {
4495 KJ_IF_SOME(err, writableController.isErroredOrErroring(js)) {
4496 return js.rejectedPromise<void>(err);
4497 }
4498 }
4499 KJ_IF_SOME(err, getReadableErrorState(js)) {
4500 return js.rejectedPromise<void>(kj::mv(err));
4501 }
4502 return js.resolvedPromise();
4503 }
4504 
4505 // Mark that we're starting a finish operation before running the algorithm,
4506 // since the algorithm may synchronously call other finish operations.
4507 algorithms.finishStarted = true;
4508 }
4509 
4510 auto onSuccess =
4511 JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), (jsg::Lock & js)->jsg::Promise<void> {
4512 // If the stream was errored during the flush algorithm (e.g., by controller.error()
4513 // or by a parallel cancel() calling abort()), we should reject with that error.
4514 if (FeatureFlags::get(js).getPedanticWpt()) {
4515 KJ_IF_SOME(err, ref->getReadableErrorState(js)) {
4516 return js.rejectedPromise<void>(kj::mv(err));
4517 } else {
4518 // Else block to avert dangling else compiler warning.
4519 }
4520 }
4521 // Allows for a graceful close of the readable side. Close will
4522 // complete once all of the queued data is read or the stream
4523 // errors. Only close if the stream can still be closed (e.g.,
4524 // it wasn't closed by a cancel operation from within flush).
4525 {
4526 KJ_IF_SOME(readableController, ref->tryGetReadableController()) {
4527 if (readableController.canCloseOrEnqueue()) {
4528 readableController.close(js);
4529 }
4530 } else {
4531 // Else block to avert dangling else compiler warning.
4532 }
4533 }
4534 return js.resolvedPromise();
4535 });
4536 
4537 auto onFailure = JSG_VISITABLE_LAMBDA(
4538 (ref = JSG_THIS), (ref), (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> {
4539 ref->error(js, reason.getHandle(js));
4540 return js.rejectedPromise<void>(kj::mv(reason));
4541 });
4542 
4543 if (flags.getPedanticWpt()) {
4544 return algorithms.maybeFinish
4545 .emplace(
4546 maybeRunAlgorithm(js, algorithms.flush, kj::mv(onSuccess), kj::mv(onFailure), JSG_THIS))
4547 .whenResolved(js);
4548 }
4549 
4550 return maybeRunAlgorithm(js, algorithms.flush, kj::mv(onSuccess), kj::mv(onFailure), JSG_THIS);
4551}
4552 
4553jsg::Promise<void> TransformStreamDefaultController::pull(jsg::Lock& js) {
4554 KJ_ASSERT(backpressure);
4555 setBackpressure(js, false);
4556 return KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.whenResolved(js);
4557}
4558 
4559jsg::Promise<void> TransformStreamDefaultController::cancel(
4560 jsg::Lock& js, v8::Local<v8::Value> reason) {
4561 if (FeatureFlags::get(js).getPedanticWpt()) {
4562 // If a finish operation is already in progress, return the existing promise
4563 // or check for errors if we're being called synchronously from within another
4564 // finish operation.
4565 if (algorithms.finishStarted) {
4566 KJ_IF_SOME(finish, algorithms.maybeFinish) {
4567 return finish.whenResolved(js);
4568 }
4569 // finishStarted is true but maybeFinish is not set yet - check if the stream
4570 // was errored during that operation.
4571 KJ_IF_SOME(err, getReadableErrorState(js)) {
4572 return js.rejectedPromise<void>(kj::mv(err));
4573 }
4574 return js.resolvedPromise();
4575 }
4576 
4577 // Mark that we're starting a finish operation before running the algorithm.
4578 algorithms.finishStarted = true;
4579 }
4580 
4581 return algorithms.maybeFinish
4582 .emplace(maybeRunAlgorithm(js, algorithms.cancel,
4583 JSG_VISITABLE_LAMBDA(
4584 (this, ref = JSG_THIS, reason = jsg::JsRef(js, jsg::JsValue(reason))), (ref, reason),
4585 (jsg::Lock & js)->jsg::Promise<void> {
4586 // If the stream was errored during the cancel algorithm (e.g., by controller.error()
4587 // or by a parallel abort()), we should reject with that error.
4588 if (FeatureFlags::get(js).getPedanticWpt()) {
4589 KJ_IF_SOME(err, getReadableErrorState(js)) {
4590 readable = kj::none;
4591 errorWritableAndUnblockWrite(js, reason.getHandle(js));
4592 return js.rejectedPromise<void>(kj::mv(err));
4593 } else {
4594 // Else block to avert dangling else compiler warning.
4595 }
4596 }
4597 readable = kj::none;
4598 errorWritableAndUnblockWrite(js, reason.getHandle(js));
4599 return js.resolvedPromise();
4600 }),
4601 JSG_VISITABLE_LAMBDA((this, ref = JSG_THIS), (ref),
4602 (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> {
4603 readable = kj::none;
4604 errorWritableAndUnblockWrite(js, reason.getHandle(js));
4605 return js.rejectedPromise<void>(kj::mv(reason));
4606 }),
4607 jsg::JsValue(reason)))
4608 .whenResolved(js);
4609}
4610 
4611jsg::Promise<void> TransformStreamDefaultController::performTransform(
4612 jsg::Lock& js, v8::Local<v8::Value> chunk) {
4613 if (algorithms.transform != kj::none) {
4614 return maybeRunAlgorithm(js, algorithms.transform,
4615 [](jsg::Lock& js) -> jsg::Promise<void> { return js.resolvedPromise(); },
4616 JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref),
4617 (jsg::Lock & js, jsg::Value reason)->jsg::Promise<void> {
4618 ref->error(js, reason.getHandle(js));
4619 return js.rejectedPromise<void>(kj::mv(reason));
4620 }),
4621 chunk, JSG_THIS);
4622 }
4623 // If we got here, there is no transform algorithm. Per the spec, the default
4624 // behavior then is to just pass along the value untransformed.
4625 return js.tryCatch([&] {
4626 enqueue(js, chunk);
4627 return js.resolvedPromise();
4628 }, [&](jsg::Value exception) { return js.rejectedPromise<void>(kj::mv(exception)); });
4629}
4630 
4631void TransformStreamDefaultController::setBackpressure(jsg::Lock& js, bool newBackpressure) {
4632 KJ_ASSERT(newBackpressure != backpressure);
4633 KJ_IF_SOME(prp, maybeBackpressureChange) {
4634 prp.resolver.resolve(js);
4635 }
4636 maybeBackpressureChange = js.newPromiseAndResolver<void>();
4637 KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.markAsHandled(js);
4638 backpressure = newBackpressure;
4639}
4640 
4641void TransformStreamDefaultController::errorWritableAndUnblockWrite(
4642 jsg::Lock& js, v8::Local<v8::Value> reason) {
4643 algorithms.clear();
4644 KJ_IF_SOME(writableController, tryGetWritableController()) {
4645 if (FeatureFlags::get(js).getPedanticWpt()) {
4646 // Use errorIfNeeded which goes through the proper error transition (Erroring -> Errored).
4647 // This allows close() to be called while the stream is "erroring" and reject with the
4648 // stored error, which is the expected behavior per the WHATWG streams spec.
4649 writableController.errorIfNeeded(js, reason);
4650 } else if (writableController.isWritable()) {
4651 writableController.doError(js, reason);
4652 }
4653 writable = kj::none;
4654 }
4655 if (backpressure) {
4656 setBackpressure(js, false);
4657 }
4658}
4659 
4660void TransformStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) {
4661 KJ_IF_SOME(backpressureChange, maybeBackpressureChange) {
4662 visitor.visit(backpressureChange.promise, backpressureChange.resolver);
4663 }
4664 visitor.visit(writable, readable, startPromise.resolver, startPromise.promise, algorithms);
4665}
4666 
4667void TransformStreamDefaultController::init(jsg::Lock& js,
4668 jsg::Ref<ReadableStream>& readable,
4669 jsg::Ref<WritableStream>& writable,
4670 jsg::Optional<Transformer> maybeTransformer) {
4671 KJ_ASSERT(this->readable == kj::none);
4672 KJ_ASSERT(this->writable == kj::none);
4673 
4674 this->writable = writable.addRef();
4675 
4676 // The TransformStreamDefaultController needs to have a reference to the underlying controller
4677 // and not just the readable because if the readable is teed, or passed off to source, etc,
4678 // the TransformStream has to make sure that it can continue to interface with the controller
4679 // to push data into it.
4680 auto& readableController = static_cast<ReadableStreamJsController&>(readable->getController());
4681 auto readableRef = KJ_ASSERT_NONNULL(readableController.getController());
4682 this->readable = KJ_ASSERT_NONNULL(readableRef.tryGet<DefaultController>()).addRef();
4683 
4684 auto transformer = kj::mv(maybeTransformer).orDefault({});
4685 
4686 // TODO(someday): The stream standard includes placeholders for supporting byte-oriented
4687 // TransformStreams but does not yet define them. For now, we are limiting our implementation
4688 // here to only support value-based transforms.
4689 JSG_REQUIRE(transformer.readableType == kj::none, TypeError,
4690 "transformer.readableType must be undefined.");
4691 JSG_REQUIRE(transformer.writableType == kj::none, TypeError,
4692 "transformer.writableType must be undefined.");
4693 
4694 KJ_IF_SOME(transform, transformer.transform) {
4695 algorithms.transform = kj::mv(transform);
4696 }
4697 
4698 KJ_IF_SOME(flush, transformer.flush) {
4699 algorithms.flush = kj::mv(flush);
4700 }
4701 
4702 KJ_IF_SOME(cancel, transformer.cancel) {
4703 algorithms.cancel = kj::mv(cancel);
4704 }
4705 
4706 setBackpressure(js, true);
4707 
4708 maybeRunAlgorithm(js, transformer.start,
4709 JSG_VISITABLE_LAMBDA(
4710 (ref = JSG_THIS), (ref), (jsg::Lock& js) { ref->startPromise.resolver.resolve(js); }),
4711 JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref),
4712 (jsg::Lock& js, jsg::Value reason) {
4713 ref->startPromise.resolver.reject(js, reason.getHandle(js));
4714 }),
4715 JSG_THIS);
4716}
4717 
4718kj::Maybe<ReadableStreamDefaultController&> TransformStreamDefaultController::
4719 tryGetReadableController() {
4720 KJ_IF_SOME(controller, readable) {
4721 return *controller;
4722 }
4723 return kj::none;
4724}
4725 
4726kj::Maybe<WritableStreamJsController&> TransformStreamDefaultController::
4727 tryGetWritableController() {
4728 KJ_IF_SOME(w, writable) {
4729 return static_cast<WritableStreamJsController&>(w->getController());
4730 }
4731 return kj::none;
4732}
4733 
4734kj::Maybe<jsg::Value> TransformStreamDefaultController::getReadableErrorState(jsg::Lock& js) {
4735 KJ_IF_SOME(controller, tryGetReadableController()) {
4736 return controller.getMaybeErrorState(js);
4737 }
4738 return kj::none;
4739}
4740 
4741template <class Self>
4742kj::StringPtr WritableImpl<Self>::jsgGetMemoryName() const {
4743 return "WritableImpl"_kjc;
4744}
4745 
4746template <class Self>
4747size_t WritableImpl<Self>::jsgGetMemorySelfSize() const {
4748 return sizeof(WritableImpl<Self>);
4749}
4750 
4751template <class Self>
4752void WritableImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
4753 tracker.trackField("signal", signal);
4754 
4755 KJ_SWITCH_ONEOF(state) {
4756 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
4757 KJ_CASE_ONEOF(error, StreamStates::Errored) {
4758 tracker.trackField("error", error);
4759 }
4760 KJ_CASE_ONEOF(erroring, StreamStates::Erroring) {
4761 tracker.trackField("erroring", erroring.reason);
4762 }
4763 KJ_CASE_ONEOF(writable, Writable) {}
4764 }
4765 
4766 tracker.trackField("abortAlgorithm", algorithms.abort);
4767 tracker.trackField("closeAlgorithm", algorithms.close);
4768 tracker.trackField("writeAlgorithm", algorithms.write);
4769 tracker.trackField("sizeAlgorithm", algorithms.size);
4770 
4771 for (auto& request: writeRequests) {
4772 tracker.trackField("pendingWrite", request);
4773 }
4774 
4775 tracker.trackField("inFlightWrite", inFlightWrite);
4776 tracker.trackField("inFlightClose", inFlightClose);
4777 tracker.trackField("closeRequest", closeRequest);
4778 tracker.trackField("maybePendingAbort", maybePendingAbort);
4779}
4780 
4781kj::StringPtr WritableStreamJsController::jsgGetMemoryName() const {
4782 return "WritableStreamJsController"_kjc;
4783}
4784 
4785size_t WritableStreamJsController::jsgGetMemorySelfSize() const {
4786 return sizeof(WritableStreamJsController);
4787}
4788 
4789void WritableStreamJsController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
4790 KJ_SWITCH_ONEOF(state) {
4791 KJ_CASE_ONEOF(initial, Initial) {}
4792 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
4793 KJ_CASE_ONEOF(error, StreamStates::Errored) {
4794 tracker.trackField("error", error);
4795 }
4796 KJ_CASE_ONEOF(controller, Controller) {
4797 tracker.trackField("controller", controller);
4798 }
4799 }
4800 tracker.trackField("lock", lock);
4801 tracker.trackField("maybeAbortPromise", maybeAbortPromise);
4802}
4803 
4804void WritableStreamDefaultController::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
4805 tracker.trackField("impl", impl);
4806}
4807 
4808kj::StringPtr ReadableStreamJsController::jsgGetMemoryName() const {
4809 return "ReadableStreamJsController"_kjc;
4810}
4811 
4812size_t ReadableStreamJsController::jsgGetMemorySelfSize() const {
4813 return sizeof(ReadableStreamJsController);
4814}
4815 
4816void ReadableStreamJsController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
4817 KJ_SWITCH_ONEOF(state) {
4818 KJ_CASE_ONEOF(initial, Initial) {}
4819 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
4820 KJ_CASE_ONEOF(error, StreamStates::Errored) {
4821 tracker.trackField("error", error);
4822 }
4823 KJ_CASE_ONEOF(readable, kj::Own<ValueReadable>) {
4824 tracker.trackField("readable", readable);
4825 }
4826 KJ_CASE_ONEOF(readable, kj::Own<ByteReadable>) {
4827 tracker.trackField("readable", readable);
4828 }
4829 }
4830 
4831 tracker.trackField("lock", lock);
4832 
4833 // Track pending error state if present (Closed has no trackable content)
4834 KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe<StreamStates::Errored>()) {
4835 tracker.trackField("pendingError", pendingError);
4836 }
4837}
4838 
4839template <class Self>
4840kj::StringPtr ReadableImpl<Self>::jsgGetMemoryName() const {
4841 return "ReadableImpl"_kjc;
4842}
4843 
4844template <class Self>
4845size_t ReadableImpl<Self>::jsgGetMemorySelfSize() const {
4846 return sizeof(ReadableImpl);
4847}
4848 
4849template <class Self>
4850void ReadableImpl<Self>::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const {
4851 KJ_SWITCH_ONEOF(state) {
4852 KJ_CASE_ONEOF(closed, StreamStates::Closed) {}
4853 KJ_CASE_ONEOF(error, StreamStates::Errored) {
4854 tracker.trackField("error", error);
4855 }
4856 KJ_CASE_ONEOF(queue, Queue) {
4857 tracker.trackField("queue", queue);
4858 }
4859 }
4860 
4861 tracker.trackField("startAlgorithm", algorithms.start);
4862 tracker.trackField("pullAlgorithm", algorithms.pull);
4863 tracker.trackField("cancelAlgorithm", algorithms.cancel);
4864 tracker.trackField("sizeAlgorithm", algorithms.size);
4865 tracker.trackField("pendingCancel", maybePendingCancel);
4866}
4867 
4868void ReadableStreamBYOBRequest::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
4869 KJ_IF_SOME(impl, maybeImpl) {
4870 tracker.trackField("readRequest", impl.readRequest);
4871 tracker.trackField("view", impl.view);
4872 }
4873}
4874 
4875void TransformStreamDefaultController::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
4876 tracker.trackField("startPromise", startPromise);
4877 tracker.trackField("maybeBackpressureChange", maybeBackpressureChange);
4878 tracker.trackField("transformAlgorithm", algorithms.transform);
4879 tracker.trackField("flushAlgorithm", algorithms.flush);
4880 tracker.trackField("writable", writable);
4881 tracker.trackField("readable", readable);
4882}
4883 
4884// ======================================================================================
4885 
4886jsg::Ref<ReadableStream> ReadableStream::from(
4887 jsg::Lock& js, jsg::AsyncGenerator<jsg::Value> generator) {
4888 
4889 // AsyncGenerator is not a refcounted type, so we need to wrap it in a refcounted
4890 // struct so that we can keep it alive through the various promise branches below.
4891 auto rcGenerator =
4892 kj::rc<kj::RefcountedWrapper<jsg::AsyncGenerator<jsg::Value>>>(kj::mv(generator));
4893 
4894 // clang-format off
4895 return constructor(js, UnderlyingSource{
4896 .pull = [generator = rcGenerator.addRef()](jsg::Lock& js, auto controller) mutable {
4897 auto& c = controller.template get<DefaultController>();
4898 return generator->getWrapped().next(js).then(js,
4899 JSG_VISITABLE_LAMBDA((controller = c.addRef(), generator = generator.addRef()),
4900 (controller),
4901 (jsg::Lock& js, kj::Maybe<jsg::Value> value) {
4902 KJ_IF_SOME(v, value) {
4903 auto handle = v.getHandle(js);
4904 // Per the ReadableStream.from spec, if the value is a promise,
4905 // the stream should wait for it to resolve and enqueue the
4906 // resolved value...
4907 // ... yes, this means that ReadableStream.from where the inputs
4908 // are promises will be slow, but that's the spec.
4909 if (handle->IsPromise()) {
4910 return js.toPromise(handle.As<v8::Promise>()).then(js,
4911 JSG_VISITABLE_LAMBDA(
4912 (controller=controller.addRef()),
4913 (controller),
4914 (jsg::Lock& js, jsg::Value val) mutable {
4915 controller->enqueue(js, val.getHandle(js));
4916 return js.resolvedPromise();
4917 }));
4918 }
4919 controller->enqueue(js, v.getHandle(js));
4920 } else {
4921 controller->close(js);
4922 }
4923 return js.resolvedPromise();
4924 }),
4925 JSG_VISITABLE_LAMBDA((controller = c.addRef(), generator = generator.addRef()),
4926 (controller), (jsg::Lock& js, jsg::Value reason) {
4927 controller->error(js, reason.getHandle(js));
4928 return js.rejectedPromise<void>(kj::mv(reason));
4929 }));
4930 },
4931 .cancel = [generator = rcGenerator.addRef()](jsg::Lock& js, auto reason) mutable {
4932 return generator->getWrapped().return_(js, js.v8Ref(reason))
4933 .then(js, [generator = kj::mv(generator)](auto& lock, auto) {
4934 // The generator might produce a value on return and might even want to continue,
4935 // but the stream has been canceled at this point, so we stop here.
4936 });
4937 },
4938 }, StreamQueuingStrategy{ .highWaterMark = 0 });
4939 // clang-format on
4940}
4941 
4942} // namespace workerd::api