Skip to content
File

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

cpp809 lines
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#pragma once
6 
7#include "common.h"
8#include "queue.h"
9 
10#include <workerd/jsg/jsg.h>
11#include <workerd/util/ring-buffer.h>
12#include <workerd/util/state-machine.h>
13#include <workerd/util/weak-refs.h>
14 
15namespace workerd::api {
16 
17// =======================================================================================
18// ReadableStreamJsController, WritableStreamJsController, and the rest here define the
19// implementation of JavaScript-backed ReadableStream and WritableStreams.
20//
21// A JavaScript-backed ReadableStream is backed by a ReadableStreamJsController that is either
22// Closed, Errored, or in a Readable state. When readable, the controller owns either a
23// ReadableStreamDefaultController or ReadableByteStreamController object that corresponds
24// to the identically named interfaces in the streams spec. These objects are responsible
25// for the bulk of the implementation detail, with the ReadableStreamJsController serving
26// only as a bridge between it and the ReadableStream object itself.
27//
28// * ReadableStream -> ReadableStreamJsController -> jsg::Ref<ReadableStreamDefaultController>
29// * ReadableStream -> ReadableStreamJsController -> jsg::Ref<ReadableByteStreamController>
30//
31// Contrast this with the implementation of internal streams using the
32// ReadableStreamInternalController:
33//
34// * ReadableStream -> ReadableStreamInternalController -> IoOwn<ReadableStreamSource>
35//
36// When user-code creates a JavaScript-backed ReadableStream using the `ReadableStream`
37// object constructor, they pass along an object called an "underlying source" that provides
38// JavaScript functions the ReadableStream will call to either initialize, close, or source
39// the data for the stream:
40//
41// const readable = new ReadableStream({
42// async start(controller) {
43// // Initialize the stream
44// },
45// async pull(controller) {
46// // Provide the stream data
47// },
48// async cancel(reason) {
49// // Cancel and de-initialize the stream
50// }
51// });
52//
53// By default, a JavaScript-backed ReadableStream is value-oriented -- that is, any JavaScript
54// type can be passed through the stream. It is not limited to bytes only. The implementation
55// of the pull method on the underlying source can push strings, booleans, numbers, even undefined
56// as values that can be read from the stream. In such streams, the `controller` used internally
57// (and owned by the ReadableStreamJsController) is the `ReadableStreamDefaultController`.
58//
59// To create a byte-oriented stream -- one that is capable only of working with bytes in the
60// form of ArrayBufferViews (e.g. `Uint8Array`, `Uint16Array`, `DataView`, etc), the underlying
61// source object passed into the `ReadableStream` constructor must have a property
62// `'type' = 'bytes'`.
63//
64// const readable = new ReadableStream({
65// type: 'bytes',
66// async start(controller) {
67// // Initialize the stream
68// },
69// async pull(controller) {
70// // Provide the stream data
71// },
72// async cancel(reason) {
73// // Cancel and de-initialize the stream
74// }
75// });
76//
77// From here on, we'll refer to these as either value streams or byte streams. And we'll refer to
78// ReadableStreamDefaultController as simply "DefaultController", and ReadableByteStreamController
79// as simply "ByobController".
80//
81// The DefaultController and ByobController each maintain an internal queue. When a read request
82// is received, if there is enough data in the internal queue to fulfill the read request, then
83// we do so. Otherwise, the controller will call the underlying source's pull method to ask it
84// to provide data to fulfill the read request.
85//
86// A critical aspect of the implementation here is that for JavaScript-backed streams, the entire
87// implementation never leaves the isolate lock, and we use JavaScript promises (via jsg::Promise)
88// instead of kj::Promise's to keep the implementation from having to bounce back and forth between
89// the two spaces. This means that with a JavaScript-backed ReadableStream, it is possible to read
90// and fully consume the stream entirely from within JavaScript without ever engaging the kj event
91// loop.
92//
93// When you tee() a JavaScript-backed ReadableStream, the stream is put into a locked state and
94// the data is funneled out through two separate "branches" (two new `ReadableStream`s).
95//
96// When anything reads from a tee branch, the underlying controller is asked to read from the
97// underlying source. When the underlying source responds to that read request, the
98// data is forwarded to all of the known branches.
99//
100// The story for JavaScript-backed writable streams is similar. User code passes what the
101// spec calls an "underlying sink" to the `WritableStream` object constructor. This provides
102// functions that are used to receive stream data.
103//
104// const writable = new WritableStream({
105// async start(controller) {
106// // initialize
107// },
108// async write(chunk, controller) {
109// // process the written chunk
110// },
111// async abort(reason) {},
112// async close(reason) {},
113// });
114//
115// It is important to note that JavaScript-backed WritableStream's are *always* value
116// oriented. It is up to the implementation of the underlying sink to determine if it is
117// capable of doing anything with whatever type of chunk it is given.
118//
119// JavaScript-backed WritableStreams are backed by the WritableStreamJsController and
120// WritableStreamDefaultController objects:
121//
122// WritableStream -> WritableStreamJsController -> jsg::Ref<WritableStreamDefaultController>
123//
124// All write operations on a JavaScript-backed WritableStream are processed within the
125// isolate lock using JavaScript promises instead of kj::Promises.
126 
127class ReadableStreamJsController;
128class WritableStreamJsController;
129 
130// =======================================================================================
131// The ReadableImpl provides implementation that is common to both the
132// ReadableStreamDefaultController and the ReadableByteStreamController.
133template <class Self>
134class ReadableImpl {
135 public:
136 using Consumer = Self::QueueType::Consumer;
137 using Entry = Self::QueueType::Entry;
138 using StateListener = Self::QueueType::ConsumerImpl::StateListener;
139 
140 ReadableImpl(UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy);
141 
142 // Invokes the start algorithm to initialize the underlying source.
143 void start(jsg::Lock& js, jsg::Ref<Self> self);
144 
145 // If the readable is not already closed or errored, initiates a cancellation.
146 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> maybeReason);
147 
148 // True if the readable is not closed, not errored, and close has not already been requested.
149 bool canCloseOrEnqueue();
150 
151 // Invokes the cancel algorithm to let the underlying source know that the
152 // readable has been canceled.
153 void doCancel(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason);
154 
155 // Close the queue if we are in a state where we can be closed.
156 void close(jsg::Lock& js);
157 
158 // Push a chunk of data into the queue.
159 void enqueue(jsg::Lock& js, kj::Rc<Entry> entry, jsg::Ref<Self> self);
160 
161 void doClose(jsg::Lock& js);
162 
163 // If it isn't already errored or closed, errors the queue, causing all consumers to be errored
164 // and detached.
165 void doError(jsg::Lock& js, jsg::Value reason);
166 
167 // When a negative number is returned, indicates that we are above the highwatermark
168 // and backpressure should be signaled.
169 kj::Maybe<int> getDesiredSize();
170 
171 // Invokes the pull algorithm only if we're in a state where the queue the
172 // queue is below the watermark and we actually need data right now.
173 void pullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self);
174 
175 // Like pullIfNeeded but bypasses the shouldCallPull() check. Used for draining reads
176 // which need to pull all available data regardless of backpressure settings.
177 void forcePullIfNeeded(jsg::Lock& js, jsg::Ref<Self> self);
178 
179 // True if the queue is current below the highwatermark.
180 bool shouldCallPull();
181 
182 // True if a pull is currently in progress (the pull promise is pending).
183 // Used by draining reads to determine if pumping completed synchronously.
184 bool isPulling() const {
185 return flags.pulling;
186 }
187 
188 // The consumer can be used to read from this readables queue so long as the queue
189 // is open. The consumer instance may outlive the readable but will be put into
190 // a closed state or errored state when the readable is destroyed.
191 kj::Own<Consumer> getConsumer(kj::Maybe<StateListener&> listener);
192 
193 // The number of consumers that exist for this readable.
194 size_t consumerCount();
195 
196 void visitForGc(jsg::GcVisitor& visitor);
197 
198 kj::StringPtr jsgGetMemoryName() const;
199 size_t jsgGetMemorySelfSize() const;
200 void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
201 
202 private:
203 struct Algorithms {
204 kj::Maybe<jsg::Function<UnderlyingSource::StartAlgorithm>> start;
205 kj::Maybe<jsg::Function<UnderlyingSource::PullAlgorithm>> pull;
206 kj::Maybe<jsg::Function<UnderlyingSource::CancelAlgorithm>> cancel;
207 kj::Maybe<jsg::Function<StreamQueuingStrategy::SizeAlgorithm>> size;
208 
209 Algorithms(UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy)
210 : start(kj::mv(underlyingSource.start)),
211 pull(kj::mv(underlyingSource.pull)),
212 cancel(kj::mv(underlyingSource.cancel)),
213 size(kj::mv(queuingStrategy.size)) {}
214 
215 Algorithms(Algorithms&& other) = default;
216 Algorithms& operator=(Algorithms&& other) = default;
217 
218 void clear() {
219 start = kj::none;
220 pull = kj::none;
221 cancel = kj::none;
222 size = kj::none;
223 }
224 
225 void visitForGc(jsg::GcVisitor& visitor) {
226 visitor.visit(start, pull, cancel, size);
227 }
228 };
229 
230 using Queue = Self::QueueType;
231 
232 // State machine for ReadableImpl:
233 // Queue is the active state where the stream can accept data
234 // Closed and Errored are terminal states (cannot transition back to Queue)
235 // Queue -> Closed (close() or doCancel() called)
236 // Queue -> Errored (doError() called)
237 using State = StateMachine<TerminalStates<StreamStates::Closed>,
238 ErrorState<StreamStates::Errored>,
239 ActiveState<Queue>,
240 StreamStates::Closed,
241 StreamStates::Errored,
242 Queue>;
243 State state;
244 Algorithms algorithms;
245 
246 size_t highWaterMark = 1;
247 
248 struct PendingCancel {
249 kj::Maybe<jsg::Promise<void>::Resolver> fulfiller;
250 jsg::Promise<void> promise;
251 JSG_MEMORY_INFO(PendingCancel) {
252 tracker.trackField("fulfiller", fulfiller);
253 tracker.trackField("promise", promise);
254 }
255 };
256 kj::Maybe<PendingCancel> maybePendingCancel;
257 
258 struct Flags {
259 uint8_t pullAgain : 1 = 0;
260 uint8_t pulling : 1 = 0;
261 uint8_t started : 1 = 0;
262 uint8_t starting : 1 = 0;
263 };
264 Flags flags{};
265 
266 friend Self;
267};
268 
269// Utility that provides the core implementation of WritableStreamJsController,
270// separated out for consistency with ReadableStreamJsController/ReadableImpl and
271// to enable it to be more easily reused should new kinds of WritableStream
272// controllers be introduced.
273template <class Self>
274class WritableImpl {
275 public:
276 using PendingAbort = WritableStreamController::PendingAbort;
277 
278 struct WriteRequest {
279 jsg::Promise<void>::Resolver resolver;
280 jsg::Value value;
281 size_t size;
282 
283 void visitForGc(jsg::GcVisitor& visitor) {
284 visitor.visit(resolver, value);
285 }
286 
287 JSG_MEMORY_INFO(WriteRequest) {
288 tracker.trackField("resolver", resolver);
289 tracker.trackField("value", value);
290 }
291 };
292 
293 WritableImpl(jsg::Lock& js, WritableStream& owner, jsg::Ref<AbortSignal> abortSignal);
294 
295 jsg::Promise<void> abort(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason);
296 
297 void advanceQueueIfNeeded(jsg::Lock& js, jsg::Ref<Self> self);
298 
299 jsg::Promise<void> close(jsg::Lock& js, jsg::Ref<Self> self);
300 
301 void dealWithRejection(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason);
302 
303 WriteRequest dequeueWriteRequest();
304 
305 void doClose(jsg::Lock& js);
306 
307 void doError(jsg::Lock& js, v8::Local<v8::Value> reason);
308 
309 void error(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason);
310 
311 void finishErroring(jsg::Lock& js, jsg::Ref<Self> self);
312 
313 void finishInFlightClose(
314 jsg::Lock& js, jsg::Ref<Self> self, kj::Maybe<v8::Local<v8::Value>> reason = kj::none);
315 
316 void finishInFlightWrite(
317 jsg::Lock& js, jsg::Ref<Self> self, kj::Maybe<v8::Local<v8::Value>> reason = kj::none);
318 
319 ssize_t getDesiredSize();
320 
321 bool isCloseQueuedOrInFlight();
322 
323 void rejectCloseAndClosedPromiseIfNeeded(jsg::Lock& js);
324 
325 kj::Maybe<WritableStreamJsController&> tryGetOwner();
326 
327 void setup(jsg::Lock& js,
328 jsg::Ref<Self> self,
329 UnderlyingSink underlyingSink,
330 StreamQueuingStrategy queuingStrategy);
331 
332 // Puts the writable into an erroring state. This allows any in flight write or
333 // close to complete before actually transitioning the writable.
334 void startErroring(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> reason);
335 
336 // Notifies the Writer of the current backpressure state. If the amount of data queued
337 // is equal to or above the highwatermark, then backpressure is applied.
338 void updateBackpressure(jsg::Lock& js);
339 
340 // Writes a chunk to the Writable, possibly queuing the chunk in the internal buffer
341 // if there are already other writes pending.
342 jsg::Promise<void> write(jsg::Lock& js, jsg::Ref<Self> self, v8::Local<v8::Value> value);
343 
344 // True if the writable is in a state where new chunks can be written
345 bool isWritable() const;
346 
347 void cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason);
348 
349 void visitForGc(jsg::GcVisitor& visitor);
350 
351 kj::StringPtr jsgGetMemoryName() const;
352 size_t jsgGetMemorySelfSize() const;
353 void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const;
354 
355 private:
356 struct Algorithms {
357 kj::Maybe<jsg::Function<UnderlyingSink::AbortAlgorithm>> abort;
358 kj::Maybe<jsg::Function<UnderlyingSink::CloseAlgorithm>> close;
359 kj::Maybe<jsg::Function<UnderlyingSink::WriteAlgorithm>> write;
360 kj::Maybe<jsg::Function<StreamQueuingStrategy::SizeAlgorithm>> size;
361 
362 Algorithms() {};
363 ~Algorithms() {
364 // Clear all algorithm references to break circular references
365 clear();
366 }
367 Algorithms(Algorithms&& other) = default;
368 Algorithms& operator=(Algorithms&& other) = default;
369 
370 void clear() {
371 abort = kj::none;
372 close = kj::none;
373 size = kj::none;
374 write = kj::none;
375 }
376 
377 void visitForGc(jsg::GcVisitor& visitor) {
378 visitor.visit(write, close, abort, size);
379 }
380 };
381 
382 struct Writable {
383 static constexpr kj::StringPtr NAME KJ_UNUSED = "writable"_kj;
384 };
385 
386 // State machine for WritableImpl:
387 // Writable is the active state where the stream can accept writes
388 // Erroring is a transitional state - waiting for in-flight ops before erroring
389 // Closed and Errored are terminal states
390 // Writable -> Erroring (startErroring() called)
391 // Writable -> Closed (finishInFlightClose() succeeds)
392 // Erroring -> Errored (finishErroring() called)
393 // Erroring -> Closed (finishInFlightClose() succeeds - close wins)
394 using State = StateMachine<TerminalStates<StreamStates::Closed>,
395 ErrorState<StreamStates::Errored>,
396 ActiveState<Writable>,
397 StreamStates::Closed,
398 StreamStates::Errored,
399 StreamStates::Erroring,
400 Writable>;
401 
402 // Sadly, we have to use a weak ref here rather than jsg::Ref. This is because
403 // the jsg::Ref<WritableStream> (via its internal WritableStreamJsController)
404 // holds a strong reference to the jsg::Ref<WritableStreamDefaultController> that
405 // uses this WritableImpl. This creates a strong circular reference between jsg::Refs
406 // that isn't allowed. GcTracing ends up with a stack overflow as the two jsg::Refs
407 // try tracing each other.
408 kj::Maybe<kj::Own<WeakRef<WritableStream>>> owner;
409 jsg::Ref<AbortSignal> signal;
410 State state = State::template create<Writable>();
411 Algorithms algorithms;
412 
413 size_t highWaterMark = 1;
414 size_t amountBuffered = 0;
415 
416 RingBuffer<WriteRequest, 8> writeRequests;
417 
418 kj::Maybe<WriteRequest> inFlightWrite;
419 kj::Maybe<jsg::Promise<void>::Resolver> inFlightClose;
420 kj::Maybe<jsg::Promise<void>::Resolver> closeRequest;
421 kj::Maybe<kj::Own<PendingAbort>> maybePendingAbort;
422 
423 struct Flags {
424 uint8_t started : 1 = 0;
425 uint8_t starting : 1 = 0;
426 uint8_t backpressure : 1 = 0;
427 uint8_t pedanticWpt : 1 = 0;
428 };
429 Flags flags{};
430 
431 friend Self;
432};
433 
434// =======================================================================================
435 
436// ReadableStreamDefaultController is a JavaScript object defined by the streams specification.
437// It is capable of streaming any JavaScript value through it, including typed arrays and
438// array buffers, but treats all values as opaque. BYOB reads are not supported.
439class ReadableStreamDefaultController: public jsg::Object {
440 public:
441 using QueueType = ValueQueue;
442 using ReadableImpl = ReadableImpl<ReadableStreamDefaultController>;
443 
444 ReadableStreamDefaultController(
445 UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy);
446 
447 void start(jsg::Lock& js);
448 
449 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason);
450 
451 void close(jsg::Lock& js);
452 
453 bool canCloseOrEnqueue();
454 bool hasBackpressure();
455 kj::Maybe<int> getDesiredSize();
456 
457 void enqueue(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> chunk);
458 
459 void error(jsg::Lock& js, v8::Local<v8::Value> reason);
460 
461 void pull(jsg::Lock& js);
462 
463 // Like pull(), but bypasses backpressure checks. Used for draining reads
464 // which need to pull all available data regardless of highWaterMark.
465 void forcePull(jsg::Lock& js);
466 
467 // True if a pull is currently in progress (the pull promise is pending).
468 bool isPulling() const {
469 return impl.isPulling();
470 }
471 
472 kj::Own<ValueQueue::Consumer> getConsumer(
473 kj::Maybe<ValueQueue::ConsumerImpl::StateListener&> stateListener);
474 
475 JSG_RESOURCE_TYPE(ReadableStreamDefaultController) {
476 JSG_READONLY_PROTOTYPE_PROPERTY(desiredSize, getDesiredSize);
477 JSG_METHOD(close);
478 JSG_METHOD(enqueue);
479 JSG_METHOD(error);
480 
481 JSG_TS_OVERRIDE(<R = any> {
482 enqueue(chunk?: R): void;
483 });
484 }
485 
486 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
487 tracker.trackField("impl", impl);
488 }
489 
490 kj::Maybe<StreamStates::Errored> getMaybeErrorState(jsg::Lock& js);
491 
492 private:
493 kj::Maybe<IoContext&> ioContext;
494 ReadableImpl impl;
495 
496 void visitForGc(jsg::GcVisitor& visitor);
497};
498 
499// The ReadableStreamBYOBRequest is provided by the ReadableByteStreamController
500// and is used by user code to fill a view provided by a BYOB read request.
501// Because we always support autoAllocateChunkSize in the ReadableByteStreamController,
502// there will always be a ReadableStreamBYOBRequest available when there is a pending
503// read.
504//
505// The ReadableStreamBYOBRequest is either in an attached or detached state.
506// The request is detached when invalidate() is called. Attempts to use the request
507// after it has been detached will fail.
508//
509// Note that the casing of the name (e.g. "BYOB" instead of the kj style "Byob") is
510// dictated by the streams specification since the class name is used as the exported
511// object name.
512class ReadableStreamBYOBRequest: public jsg::Object {
513 public:
514 ReadableStreamBYOBRequest(jsg::Lock& js,
515 kj::Own<ByteQueue::ByobRequest> readRequest,
516 kj::Rc<WeakRef<ReadableByteStreamController>> controller);
517 
518 KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamBYOBRequest);
519 
520 // getAtLeast is a non-standard Workers-specific extension that specifies
521 // the minimum number of bytes the stream should fill into the view. It is
522 // added to support the readAtLeast extension on the ReadableStreamBYOBReader.
523 kj::Maybe<int> getAtLeast();
524 
525 kj::Maybe<jsg::V8Ref<v8::Uint8Array>> getView(jsg::Lock& js);
526 
527 void invalidate(jsg::Lock& js);
528 
529 void respond(jsg::Lock& js, int bytesWritten);
530 
531 void respondWithNewView(jsg::Lock& js, jsg::BufferSource view);
532 
533 JSG_RESOURCE_TYPE(ReadableStreamBYOBRequest) {
534 JSG_READONLY_PROTOTYPE_PROPERTY(view, getView);
535 JSG_METHOD(respond);
536 JSG_METHOD(respondWithNewView);
537 
538 // atLeast is an Workers-specific extension used to support the
539 // readAtLeast API.
540 JSG_READONLY_PROTOTYPE_PROPERTY(atLeast, getAtLeast);
541 }
542 
543 bool isPartiallyFulfilled();
544 
545 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
546 
547 private:
548 struct Impl {
549 kj::Own<ByteQueue::ByobRequest> readRequest;
550 kj::Rc<WeakRef<ReadableByteStreamController>> controller;
551 jsg::V8Ref<v8::Uint8Array> view;
552 
553 size_t originalBufferByteLength;
554 size_t originalByteOffsetPlusBytesFilled;
555 
556 Impl(jsg::Lock& js,
557 kj::Own<ByteQueue::ByobRequest> readRequest,
558 kj::Rc<WeakRef<ReadableByteStreamController>> controller);
559 
560 void updateView(jsg::Lock& js);
561 };
562 
563 kj::Maybe<IoContext&> ioContext;
564 kj::Maybe<Impl> maybeImpl;
565 
566 void visitForGc(jsg::GcVisitor& visitor);
567};
568 
569// ReadableByteStreamController is a JavaScript object defined by the streams specification.
570// It is capable of only streaming byte data through it in the form of typed arrays.
571// BYOB reads are supported.
572class ReadableByteStreamController: public jsg::Object {
573 public:
574 using QueueType = ByteQueue;
575 using ReadableImpl = ReadableImpl<ReadableByteStreamController>;
576 
577 ReadableByteStreamController(
578 UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy);
579 ~ReadableByteStreamController() noexcept(false);
580 
581 jsg::Ref<ReadableByteStreamController> getSelf() {
582 return JSG_THIS;
583 }
584 
585 void start(jsg::Lock& js);
586 
587 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason);
588 
589 void close(jsg::Lock& js);
590 
591 void enqueue(jsg::Lock& js, jsg::BufferSource chunk);
592 
593 void error(jsg::Lock& js, v8::Local<v8::Value> reason);
594 
595 bool canCloseOrEnqueue();
596 bool hasBackpressure();
597 kj::Maybe<int> getDesiredSize();
598 
599 kj::Maybe<jsg::Ref<ReadableStreamBYOBRequest>> getByobRequest(jsg::Lock& js);
600 
601 void pull(jsg::Lock& js);
602 
603 // Like pull(), but bypasses backpressure checks. Used for draining reads
604 // which need to pull all available data regardless of highWaterMark.
605 void forcePull(jsg::Lock& js);
606 
607 // True if a pull is currently in progress (the pull promise is pending).
608 bool isPulling() const {
609 return impl.isPulling();
610 }
611 
612 kj::Own<ByteQueue::Consumer> getConsumer(
613 kj::Maybe<ByteQueue::ConsumerImpl::StateListener&> stateListener);
614 
615 JSG_RESOURCE_TYPE(ReadableByteStreamController) {
616 JSG_READONLY_PROTOTYPE_PROPERTY(byobRequest, getByobRequest);
617 JSG_READONLY_PROTOTYPE_PROPERTY(desiredSize, getDesiredSize);
618 JSG_METHOD(close);
619 JSG_METHOD(enqueue);
620 JSG_METHOD(error);
621 }
622 
623 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
624 tracker.trackField("impl", impl);
625 tracker.trackField("maybeByobRequest", maybeByobRequest);
626 }
627 
628 private:
629 kj::Rc<WeakRef<ReadableByteStreamController>> weakSelf;
630 kj::Maybe<IoContext&> ioContext;
631 ReadableImpl impl;
632 kj::Maybe<jsg::Ref<ReadableStreamBYOBRequest>> maybeByobRequest;
633 
634 void visitForGc(jsg::GcVisitor& visitor);
635 
636 friend class ReadableStreamBYOBRequest;
637 friend class ReadableStreamJsController;
638};
639 
640// =======================================================================================
641 
642// The WritableStreamDefaultController is an object defined by the stream specification.
643// Writable streams are always value oriented. It is up the underlying sink implementation
644// to determine whether it is capable of handling whatever type of JavaScript object it
645// is given.
646class WritableStreamDefaultController: public jsg::Object {
647 public:
648 using WritableImpl = WritableImpl<WritableStreamDefaultController>;
649 
650 explicit WritableStreamDefaultController(
651 jsg::Lock& js, WritableStream& owner, jsg::Ref<AbortSignal> abortSignal);
652 
653 ~WritableStreamDefaultController() noexcept(false);
654 
655 jsg::Promise<void> abort(jsg::Lock& js, v8::Local<v8::Value> reason);
656 
657 jsg::Promise<void> close(jsg::Lock& js);
658 
659 void error(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason);
660 
661 kj::Maybe<ssize_t> getDesiredSize();
662 
663 jsg::Ref<AbortSignal> getSignal();
664 
665 kj::Maybe<v8::Local<v8::Value>> isErroring(jsg::Lock& js);
666 
667 // Returns true if the stream is in the erroring state. Unlike the overload
668 // that takes a lock, this method does not require a lock since it doesn't
669 // return the error reason.
670 bool isErroring() const;
671 
672 bool isStarted() {
673 return impl.flags.started;
674 }
675 
676 bool hasBackpressure() {
677 return impl.flags.backpressure;
678 }
679 
680 void setup(jsg::Lock& js, UnderlyingSink underlyingSink, StreamQueuingStrategy queuingStrategy);
681 
682 jsg::Promise<void> write(jsg::Lock& js, v8::Local<v8::Value> value);
683 
684 JSG_RESOURCE_TYPE(WritableStreamDefaultController) {
685 JSG_READONLY_PROTOTYPE_PROPERTY(signal, getSignal);
686 JSG_METHOD(error);
687 }
688 
689 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
690 
691 void cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason);
692 
693 // Clear algorithms to break circular references during destruction
694 void clearAlgorithms();
695 
696 private:
697 kj::Maybe<IoContext&> ioContext;
698 WritableImpl impl;
699 
700 void visitForGc(jsg::GcVisitor& visitor);
701};
702 
703// =======================================================================================
704 
705// The relationship between the TransformStreamDefaultController and the
706// readable/writable streams associated with it can be complicated.
707// Strong references to the TransformStreamDefaultController are held by
708// the *algorithms* passed into the readable and writable streams using
709// JSG_VISITABLE_LAMBDAs. When those algorithms are cleared, the strong
710// references holding the TransformStreamDefaultController are freed.
711// However, user code can do silly things like hold the Transform controller
712// long after both the readable and writable sides have been GC'ed.
713class TransformStreamDefaultController: public jsg::Object {
714 public:
715 TransformStreamDefaultController(jsg::Lock& js);
716 
717 void init(jsg::Lock& js,
718 jsg::Ref<ReadableStream>& readable,
719 jsg::Ref<WritableStream>& writable,
720 jsg::Optional<Transformer> maybeTransformer);
721 
722 // The startPromise is used by both the readable and writable sides in their respective
723 // start algorithms. The promise itself is resolved within the init function when the
724 // transformers own start algorithm completes.
725 inline jsg::Promise<void> getStartPromise(jsg::Lock& js) {
726 return startPromise.promise.whenResolved(js);
727 }
728 
729 kj::Maybe<int> getDesiredSize();
730 
731 void enqueue(jsg::Lock& js, v8::Local<v8::Value> chunk);
732 
733 void error(jsg::Lock& js, v8::Local<v8::Value> reason);
734 
735 void terminate(jsg::Lock& js);
736 
737 JSG_RESOURCE_TYPE(TransformStreamDefaultController) {
738 JSG_READONLY_PROTOTYPE_PROPERTY(desiredSize, getDesiredSize);
739 JSG_METHOD(enqueue);
740 JSG_METHOD(error);
741 JSG_METHOD(terminate);
742 
743 JSG_TS_OVERRIDE(<O = any> {
744 enqueue(chunk?: O): void;
745 });
746 }
747 
748 jsg::Promise<void> write(jsg::Lock& js, v8::Local<v8::Value> chunk);
749 jsg::Promise<void> abort(jsg::Lock& js, v8::Local<v8::Value> reason);
750 jsg::Promise<void> close(jsg::Lock& js);
751 jsg::Promise<void> pull(jsg::Lock& js);
752 jsg::Promise<void> cancel(jsg::Lock& js, v8::Local<v8::Value> reason);
753 
754 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
755 
756 private:
757 struct Algorithms {
758 kj::Maybe<jsg::Function<Transformer::TransformAlgorithm>> transform;
759 kj::Maybe<jsg::Function<Transformer::FlushAlgorithm>> flush;
760 kj::Maybe<jsg::Function<Transformer::CancelAlgorithm>> cancel;
761 
762 kj::Maybe<jsg::Promise<void>> maybeFinish = kj::none;
763 // This flag is set to true at the start of a finish operation (close/cancel/abort)
764 // before the algorithm runs. This is needed because emplace() evaluates its argument
765 // before setting maybeFinish, so if the algorithm calls another finish operation
766 // synchronously, maybeFinish wouldn't be set yet.
767 bool finishStarted = false;
768 
769 Algorithms() {};
770 Algorithms(Algorithms&& other) = default;
771 Algorithms& operator=(Algorithms&& other) = default;
772 
773 inline void clear() {
774 transform = kj::none;
775 flush = kj::none;
776 cancel = kj::none;
777 }
778 
779 inline void visitForGc(jsg::GcVisitor& visitor) {
780 visitor.visit(transform, flush, cancel, maybeFinish);
781 }
782 };
783 
784 void errorWritableAndUnblockWrite(jsg::Lock& js, v8::Local<v8::Value> reason);
785 jsg::Promise<void> performTransform(jsg::Lock& js, v8::Local<v8::Value> chunk);
786 void setBackpressure(jsg::Lock& js, bool newBackpressure);
787 
788 kj::Maybe<IoContext&> ioContext;
789 jsg::PromiseResolverPair<void> startPromise;
790 
791 kj::Maybe<ReadableStreamDefaultController&> tryGetReadableController();
792 kj::Maybe<WritableStreamJsController&> tryGetWritableController();
793 
794 kj::Maybe<jsg::Value> getReadableErrorState(jsg::Lock& js);
795 
796 // Currently, JS-backed transform streams only support value-oriented streams.
797 // In the future, that may change and this will need to become a kj::OneOf
798 // that includes a ReadableByteStreamController.
799 kj::Maybe<jsg::Ref<ReadableStreamDefaultController>> readable;
800 kj::Maybe<jsg::Ref<WritableStream>> writable;
801 Algorithms algorithms;
802 bool backpressure = false;
803 kj::Maybe<jsg::PromiseResolverPair<void>> maybeBackpressureChange;
804 
805 void visitForGc(jsg::GcVisitor& visitor);
806};
807 
808} // namespace workerd::api