Skip to content
File

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

cpp962 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 "../basics.h"
8 
9#include <workerd/io/io-context.h>
10#include <workerd/io/worker-interface.capnp.h>
11#include <workerd/jsg/jsg.h>
12 
13#if _MSC_VER
14using ssize_t = long long;
15#endif
16 
17namespace workerd::api {
18 
19class ReadableStream;
20class ReadableStreamController;
21class ReadableStreamSource;
22class ReadableStreamDefaultController;
23class ReadableByteStreamController;
24 
25class WritableStream;
26class WritableStreamController;
27class WritableStreamSink;
28class WritableStreamDefaultController;
29 
30class TransformStreamDefaultController;
31 
32using rpc::StreamEncoding;
33 
34enum class ReadAllTextOption : uint8_t {
35 NONE = 0,
36 NULL_TERMINATE = 1 << 0,
37 STRIP_BOM = 1 << 1,
38};
39 
40inline ReadAllTextOption operator|(ReadAllTextOption a, ReadAllTextOption b) {
41 return static_cast<ReadAllTextOption>(static_cast<uint8_t>(a) | static_cast<uint8_t>(b));
42}
43 
44inline ReadAllTextOption& operator|=(ReadAllTextOption& a, ReadAllTextOption b) {
45 return a = a | b;
46}
47 
48inline bool operator&(ReadAllTextOption a, ReadAllTextOption b) {
49 return (static_cast<uint8_t>(a) & static_cast<uint8_t>(b)) != 0;
50}
51 
52static constexpr kj::byte UTF8_BOM[] = {0xEF, 0xBB, 0xBF};
53static constexpr size_t UTF8_BOM_SIZE = sizeof(UTF8_BOM);
54 
55inline bool hasUtf8Bom(kj::ArrayPtr<const kj::byte> data) {
56 return data.size() >= UTF8_BOM_SIZE && memcmp(data.begin(), UTF8_BOM, UTF8_BOM_SIZE) == 0;
57}
58 
59struct ReadResult {
60 jsg::Optional<jsg::Value> value;
61 bool done;
62 
63 JSG_STRUCT(value, done);
64 JSG_STRUCT_TS_OVERRIDE(type ReadableStreamReadResult<R = any> =
65 | { done: false, value: R; }
66 | { done: true; value?: undefined; }
67 );
68 
69 void visitForGc(jsg::GcVisitor& visitor) {
70 visitor.visit(value);
71 }
72};
73 
74// Result type for draining read operations. Always returns bytes, even for value streams.
75// Used by DrainingReader for optimized pipe-to operations with vectored writes.
76// This is a C++ only type - not exposed to JavaScript.
77struct DrainingReadResult {
78 kj::Array<kj::Array<kj::byte>> chunks; // Multiple byte arrays for vectored writes
79 bool done = false; // True if stream is closed/closing
80};
81 
82struct StreamQueuingStrategy {
83 using SizeAlgorithm = uint64_t(v8::Local<v8::Value>);
84 
85 jsg::Optional<uint64_t> highWaterMark;
86 jsg::Optional<jsg::Function<SizeAlgorithm>> size;
87 
88 JSG_STRUCT(highWaterMark, size);
89 JSG_STRUCT_TS_OVERRIDE(QueuingStrategy<T = any> {
90 size?: (chunk: T) => number | bigint;
91 });
92};
93 
94struct UnderlyingSource {
95 using Controller =
96 kj::OneOf<jsg::Ref<ReadableStreamDefaultController>, jsg::Ref<ReadableByteStreamController>>;
97 using StartAlgorithm = jsg::Promise<void>(Controller);
98 using PullAlgorithm = jsg::Promise<void>(Controller);
99 using CancelAlgorithm = jsg::Promise<void>(v8::Local<v8::Value> reason);
100 
101 // The autoAllocateChunkSize mechanism allows byte streams to operate as if a BYOB
102 // reader is being used even if it is just a default reader. Support is optional
103 // per the streams spec but our implementation will always enable it. Specifically,
104 // if user code does not provide an explicit autoAllocateChunkSize, we'll assume
105 // this default.
106 static constexpr int DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE = 4096;
107 
108 // We want to increase the default auto allocate chunk size but we need to do
109 // so carefully to avoid introducing memory regressions and causing workers to
110 // hit OOM errors. We'll use an autogate to roll out the new default.
111 static constexpr int DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2 = 16 * 1024;
112 
113 // Per the spec, the type property for the UnderlyingSource should be either
114 // undefined, the empty string, or "bytes". When undefined, the empty string is
115 // used as the default. When type is the empty string, the stream is considered
116 // to be value-oriented rather than byte-oriented.
117 jsg::Optional<kj::String> type;
118 
119 // Used only when type is equal to "bytes", the autoAllocateChunkSize defines
120 // the size of automatically allocated buffer that is created when a default
121 // mode read is performed on a byte-oriented ReadableStream that supports
122 // BYOB reads. The stream standard makes this optional to support and defines
123 // no default value. We've chosen to use a default value of 4096. If given,
124 // the value must be greater than zero.
125 jsg::Optional<int> autoAllocateChunkSize;
126 
127 jsg::Optional<jsg::Function<StartAlgorithm>> start;
128 jsg::Optional<jsg::Function<PullAlgorithm>> pull;
129 jsg::Optional<jsg::Function<CancelAlgorithm>> cancel;
130 
131 // The expectedLength is a non-standard extension used to support specifying the
132 // content-length when using a ReadableStream as the body of a request or response.
133 jsg::Optional<uint64_t> expectedLength;
134 
135 JSG_STRUCT(type, autoAllocateChunkSize, start, pull, cancel, expectedLength);
136 JSG_STRUCT_TS_DEFINE(interface UnderlyingByteSource {
137 type: "bytes";
138 autoAllocateChunkSize?: number;
139 start?: (controller: ReadableByteStreamController) => void | Promise<void>;
140 pull?: (controller: ReadableByteStreamController) => void | Promise<void>;
141 cancel?: (reason: any) => void | Promise<void>;
142 });
143 JSG_STRUCT_TS_OVERRIDE(<R = any> {
144 type?: "" | undefined;
145 autoAllocateChunkSize: never;
146 start?: (controller: ReadableStreamDefaultController<R>) => void | Promise<void>;
147 pull?: (controller: ReadableStreamDefaultController<R>) => void | Promise<void>;
148 cancel?: (reason: any) => void | Promise<void>;
149 });
150};
151 
152struct UnderlyingSink {
153 using Controller = jsg::Ref<WritableStreamDefaultController>;
154 using StartAlgorithm = jsg::Promise<void>(Controller);
155 using WriteAlgorithm = jsg::Promise<void>(v8::Local<v8::Value>, Controller);
156 using AbortAlgorithm = jsg::Promise<void>(v8::Local<v8::Value> reason);
157 using CloseAlgorithm = jsg::Promise<void>();
158 
159 // Per the spec, the type property for the UnderlyingSink should always be either
160 // undefined or the empty string. Any other value will trigger a TypeError.
161 jsg::Optional<kj::String> type;
162 
163 jsg::Optional<jsg::Function<StartAlgorithm>> start;
164 jsg::Optional<jsg::Function<WriteAlgorithm>> write;
165 jsg::Optional<jsg::Function<AbortAlgorithm>> abort;
166 jsg::Optional<jsg::Function<CloseAlgorithm>> close;
167 
168 JSG_STRUCT(type, start, write, abort, close);
169 
170 // TODO(cleanup): Get rid of this override and parse the type directly in param-extractor.rs
171 JSG_STRUCT_TS_OVERRIDE(<W = any> {
172 write?: (chunk: W, controller: WritableStreamDefaultController) => void | Promise<void>;
173 start?: (controller: WritableStreamDefaultController) => void | Promise<void>;
174 abort?: (reason: any) => void | Promise<void>;
175 close?: () => void | Promise<void>;
176 });
177};
178 
179struct Transformer {
180 using Controller = jsg::Ref<TransformStreamDefaultController>;
181 using StartAlgorithm = jsg::Promise<void>(Controller);
182 using TransformAlgorithm = jsg::Promise<void>(v8::Local<v8::Value>, Controller);
183 using FlushAlgorithm = jsg::Promise<void>(Controller);
184 using CancelAlgorithm = jsg::Promise<void>(jsg::JsValue reason);
185 
186 jsg::Optional<kj::String> readableType;
187 jsg::Optional<kj::String> writableType;
188 
189 jsg::Optional<jsg::Function<StartAlgorithm>> start;
190 jsg::Optional<jsg::Function<TransformAlgorithm>> transform;
191 jsg::Optional<jsg::Function<FlushAlgorithm>> flush;
192 jsg::Optional<jsg::Function<CancelAlgorithm>> cancel;
193 
194 // The expectedLength is a non-standard extension used to support specifying the
195 // content-length when using a TransformStream readable side as the body of a
196 // request or response.
197 jsg::Optional<uint64_t> expectedLength;
198 
199 JSG_STRUCT(readableType, writableType, start, transform, flush, cancel, expectedLength);
200 JSG_STRUCT_TS_OVERRIDE(<I = any, O = any> {
201 start?: (controller: TransformStreamDefaultController<O>) => void | Promise<void>;
202 transform?: (chunk: I, controller: TransformStreamDefaultController<O>) => void | Promise<void>;
203 flush?: (controller: TransformStreamDefaultController<O>) => void | Promise<void>;
204 cancel?: (reason: any) => void | Promise<void>;
205 expectedLength?: number;
206 });
207};
208 
209// ReadableStreamSource and WritableStreamSink
210//
211// These are implementation interfaces for ReadableStream and WritableStream. If you just need to
212// use a ReadableStream or WritableStream, you can safely skip reading this. If you need to
213// implement a new kind of stream, read on.
214 
215// In the original Workers streams implementation, a ReadableStream would have a
216// ReadableStreamSource backing it. Likewise, a WritableStream would have a WritableStreamSink.
217// The ReadableStreamSource and WritableStreamSink are kj heap objects that provide a thin
218// wrapper on internal native stream sources originating from within the Workers runtime.
219//
220// With implementation of full streams standard support, we introduce the new abstraction APIs
221// ReadableStreamController and WritableStreamController, which will provide the underlying
222// implementation for both ReadableStream and WritableStream, respectively.
223//
224// When creating a new kind of *internal* ReadableStream, where the data is originating internally
225// from a kj stream, you will still implement the ReadableStreamSource API, just as before.
226// Likewise, when creating a new kind of *internal* WritableStream, where the data destination is
227// a kj stream, you will implement the WritableStreamSink API.
228 
229class WritableStreamSink {
230 public:
231 virtual kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) KJ_WARN_UNUSED_RESULT = 0;
232 virtual kj::Promise<void> write(
233 kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) KJ_WARN_UNUSED_RESULT = 0;
234 
235 virtual kj::Promise<void> end() KJ_WARN_UNUSED_RESULT = 0;
236 // Must call to flush and finish the stream.
237 
238 virtual kj::Maybe<kj::Promise<DeferredProxy<void>>> tryPumpFrom(
239 ReadableStreamSource& input, bool end);
240 
241 virtual void abort(kj::Exception reason) = 0;
242 // TODO(conform): abort() should return a promise after which closed fulfillers should be
243 // rejected. This may necessitate an "erroring" state.
244 
245 // Tells the sink that it is no longer to be responsible for encoding in the correct format.
246 // Instead, the caller takes responsibility. The expected encoding is returned; the caller
247 // promises that all future writes will use this encoding. The default implementation returns
248 // IDENTITY, which is always correct since that's the encoding write()s should have used if
249 // this weren't called at all.
250 virtual StreamEncoding disownEncodingResponsibility() {
251 return StreamEncoding::IDENTITY;
252 }
253};
254 
255class ReadableStreamSource {
256 public:
257 virtual kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) = 0;
258 
259 // The ReadableStreamSource version of pumpTo() has no `amount` parameter, since the Streams spec
260 // only defines pumping everything.
261 //
262 // If `end` is true, then `output.end()` will be called after pumping. Note that it's especially
263 // important to take advantage of this when using deferred proxying since calling `end()`
264 // directly might attempt to use the `IoContext` to call `registerPendingEvent()`.
265 virtual kj::Promise<DeferredProxy<void>> pumpTo(WritableStreamSink& output, bool end);
266 
267 // If pumpTo() pumps to a system stream, what is the best encoding for that system stream to
268 // use? This is just a hint.
269 virtual StreamEncoding getPreferredEncoding() {
270 return StreamEncoding::IDENTITY;
271 };
272 
273 virtual kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding);
274 
275 kj::Promise<kj::Array<byte>> readAllBytes(uint64_t limit);
276 kj::Promise<kj::String> readAllText(
277 uint64_t limit, ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE);
278 
279 // Hook to inform this ReadableStreamSource that the ReadableStream has been canceled. This only
280 // really means anything to TransformStreams, which are supposed to propagate the error to the
281 // writable side, and custom ReadableStreams, which we don't implement yet.
282 //
283 // NOTE: By "propagate the error back to the writable stream", I mean: if the WritableStream is in
284 // the Writable state, set it to the Errored state and reject its closed fulfiller with
285 // `reason`. I'm not sure how I'm going to do this yet.
286 virtual void cancel(kj::Exception reason);
287 // TODO(conform): Should return promise.
288 //
289 // TODO(conform): `reason` should be allowed to be any JS value, and not just an exception.
290 // That is, something silly like `stream.cancel(42)` should be allowed and trigger a
291 // rejection with the integer `42`.
292 
293 struct Tee {
294 kj::Own<ReadableStreamSource> branches[2];
295 };
296 
297 // Implement this if your ReadableStreamSource has a better way to tee a stream than the naive
298 // method, which relies upon `tryRead()`. The default implementation returns nullptr.
299 virtual kj::Maybe<Tee> tryTee(uint64_t limit);
300};
301 
302struct PipeToOptions {
303 jsg::Optional<bool> preventAbort;
304 jsg::Optional<bool> preventCancel;
305 jsg::Optional<bool> preventClose;
306 jsg::Optional<jsg::Ref<AbortSignal>> signal;
307 
308 JSG_STRUCT(preventAbort, preventCancel, preventClose, signal);
309 JSG_STRUCT_TS_OVERRIDE(StreamPipeOptions);
310 
311 // An additional, internal only property that is used to indicate
312 // when the pipe operation is used for a pipeThrough rather than
313 // a pipeTo. We use this information, for instance, to identify
314 // when we should mark returned rejected promises as handled.
315 bool pipeThrough = false;
316};
317 
318namespace StreamStates {
319struct Closed {
320 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
321};
322using Errored = jsg::Value;
323struct Erroring {
324 static constexpr kj::StringPtr NAME KJ_UNUSED = "erroring"_kj;
325 jsg::Value reason;
326 
327 Erroring(jsg::Value reason): reason(kj::mv(reason)) {}
328 
329 void visitForGc(jsg::GcVisitor& visitor) {
330 visitor.visit(reason);
331 }
332};
333} // namespace StreamStates
334 
335// A ReadableStreamController provides the underlying implementation for a ReadableStream.
336// We will generally have three implementations:
337// * ReadableStreamDefaultController
338// * ReadableByteStreamController
339// * ReadableStreamInternalController
340//
341// The ReadableStreamDefaultController and ReadableByteStreamController are defined by the
342// streams standard and source all of the stream data from JavaScript functions provided by
343// user code.
344//
345// The ReadableStreamInternalController is Workers runtime specific and provides a bridge
346// to the existing ReadableStreamSource API. At the API contract layer, the
347// ReadableByteStreamController and ReadableStreamInternalController will appear to be
348// identical. Internally, however, they will be very different from one another.
349//
350// The ReadableStreamController instance is meant to be a private member of the ReadableStream,
351// e.g.
352// class ReadableStream {
353// public:
354// // ...
355// private:
356// ReadableStreamController controller;
357// // ...
358// }
359//
360// As such, it exists within the V8 heap (it's allocated directly as a member of the
361// ReadableStream) and will always execute within the V8 isolate lock.
362//
363// The methods here return jsg::Promise rather than kj::Promise because the controller
364// operations here do not always require passing through the kj mechanisms or kj event loop.
365// Likewise, we do not make use of kj::Exception in these interfaces because the stream
366// standard dictates that streams can be canceled/aborted/errored using any arbitrary JavaScript
367// value, not just Errors.
368class ReadableStreamController {
369 public:
370 // The ReadableStreamController::Reader interface is a base for all ReadableStream reader
371 // implementations and is used solely as a means of attaching a Reader implementation to
372 // the internal state of the controller. See the ReadableStream::*Reader classes for the
373 // full Reader API.
374 class Reader {
375 public:
376 // True if the reader is a BYOB reader.
377 virtual bool isByteOriented() const = 0;
378 
379 // When a Reader is locked to a controller, the controller will attach itself to the reader,
380 // passing along the closed promise that will be used to communicate state to the
381 // user code.
382 //
383 // The Reader will hold a reference to the controller that will be cleared when the reader
384 // is released or destroyed. The controller is guaranteed to either outlive or detach the
385 // reader so the ReadableStreamController& reference should remain valid.
386 virtual void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) = 0;
387 
388 // When a Reader lock is released, the controller will signal to the reader that it has been
389 // detached.
390 virtual void detach() = 0;
391 };
392 
393 struct ByobOptions {
394 static constexpr size_t DEFAULT_AT_LEAST = 1;
395 
396 jsg::V8Ref<v8::ArrayBufferView> bufferView;
397 size_t byteOffset = 0;
398 size_t byteLength;
399 
400 // The minimum number of elements that should be read. When not specified, the default
401 // is DEFAULT_AT_LEAST. This is a non-standard, Workers-specific extension to
402 // support the readAtLeast method on the ReadableStreamBYOBReader object.
403 // ReaderImpl::read() converts this to bytes by multiplying by element size.
404 kj::Maybe<size_t> atLeast = DEFAULT_AT_LEAST;
405 
406 // True if the given buffer should be detached. Per the spec, we should always be
407 // detaching a BYOB buffer but the original Workers implementation did not.
408 // To avoid breaking backwards compatibility, a compatibility flag is provided to turn
409 // detach on/off as appropriate.
410 bool detachBuffer = true;
411 };
412 
413 struct Tee {
414 jsg::Ref<ReadableStream> branch1;
415 jsg::Ref<ReadableStream> branch2;
416 };
417 
418 // Abstract API for ReadableStreamController implementations that provide their own
419 // tee implementations that are not backed by kj's tee. Each branch of the tee uses
420 // the TeeController to interface with the shared underlying source, and the
421 // TeeController ensures that each Branch receives the data that is read.
422 class TeeController {
423 public:
424 // Represents an individual ReadableStreamController tee branch registered with
425 // a TeeController. One or more branches is registered with the TeeController.
426 class Branch {
427 public:
428 virtual ~Branch() noexcept(false) {}
429 
430 virtual void doClose(jsg::Lock& js) = 0;
431 virtual void doError(jsg::Lock& js, v8::Local<v8::Value> reason) = 0;
432 virtual void handleData(jsg::Lock& js, ReadResult result) = 0;
433 };
434 
435 class BranchPtr {
436 public:
437 inline BranchPtr(Branch* branch): inner(branch) {
438 KJ_ASSERT(inner != nullptr);
439 }
440 BranchPtr(BranchPtr&& other) = default;
441 BranchPtr& operator=(BranchPtr&&) = default;
442 BranchPtr(BranchPtr& other) = default;
443 
444 inline void doClose(jsg::Lock& js) {
445 inner->doClose(js);
446 }
447 
448 inline void doError(jsg::Lock& js, v8::Local<v8::Value> reason) {
449 inner->doError(js, reason);
450 }
451 
452 inline void handleData(jsg::Lock& js, ReadResult result) {
453 inner->handleData(js, kj::mv(result));
454 }
455 
456 inline uint hashCode() {
457 return kj::hashCode(inner);
458 }
459 inline bool operator==(BranchPtr& other) const {
460 return inner == other.inner;
461 }
462 
463 private:
464 Branch* inner;
465 };
466 
467 virtual ~TeeController() noexcept(false) {}
468 
469 virtual void addBranch(Branch* branch) = 0;
470 
471 virtual void close(jsg::Lock& js) = 0;
472 
473 virtual void error(jsg::Lock& js, v8::Local<v8::Value> reason) = 0;
474 
475 virtual void ensurePulling(jsg::Lock& js) = 0;
476 
477 // maybeJs will be nullptr when the isolate lock is not available.
478 // If maybeJs is set, any operations pending for the branch will be canceled.
479 virtual void removeBranch(Branch* branch, kj::Maybe<jsg::Lock&> maybeJs) = 0;
480 };
481 
482 // The PipeController simplifies the abstraction between ReadableStreamController
483 // and WritableStreamController so that the pipeTo/pipeThrough/tryPipeTo can work
484 // without caring about what kind of controller it is working with.
485 class PipeController {
486 public:
487 virtual ~PipeController() noexcept(false) {}
488 virtual bool isClosed() = 0;
489 virtual kj::Maybe<v8::Local<v8::Value>> tryGetErrored(jsg::Lock& js) = 0;
490 virtual void cancel(jsg::Lock& js, v8::Local<v8::Value> reason) = 0;
491 virtual void close(jsg::Lock& js) = 0;
492 virtual void error(jsg::Lock& js, v8::Local<v8::Value> reason) = 0;
493 virtual void release(jsg::Lock& js, kj::Maybe<v8::Local<v8::Value>> maybeError = kj::none) = 0;
494 virtual kj::Maybe<kj::Promise<void>> tryPumpTo(WritableStreamSink& sink, bool end) = 0;
495 virtual jsg::Promise<ReadResult> read(jsg::Lock& js) = 0;
496 };
497 
498 virtual ~ReadableStreamController() noexcept(false) {}
499 
500 virtual void setOwnerRef(ReadableStream& stream) = 0;
501 
502 virtual jsg::Ref<ReadableStream> addRef() = 0;
503 
504 // Returns true if the underlying source for this controller is byte-oriented and
505 // therefore supports the pull into API. When false, the stream can be used to pass
506 // any arbitrary JavaScript value through.
507 virtual bool isByteOriented() const = 0;
508 
509 // Reads data from the stream. If the stream is byte-oriented, then the ByobOptions can be
510 // specified to provide a v8::ArrayBuffer to be filled by the read operation. If the ByobOptions
511 // are provided and the stream is not byte-oriented, the operation will return a rejected promise.
512 virtual kj::Maybe<jsg::Promise<ReadResult>> read(
513 jsg::Lock& js, kj::Maybe<ByobOptions> byobOptions) = 0;
514 
515 // Performs a draining read operation that:
516 // 1. Drains all currently buffered data from the queue
517 // 2. Pumps the controller for synchronously available data (respecting pull promise state)
518 // 3. Returns bytes even for value streams (converting ArrayBuffer/ArrayBufferView/string)
519 // 4. Has mutual exclusion with regular reads - returns rejected promise if pending regular reads
520 // 5. Returns done: true with final data when stream is closing
521 //
522 // This is a C++ only API (not exposed to JavaScript) intended for optimized pipe operations.
523 // Returns kj::none if the stream is locked in a way that prevents the read.
524 //
525 // The maxRead parameter provides a soft limit on how much data to read. Both the initial
526 // buffer drain and subsequent synchronous pump attempts stop when the total bytes read
527 // reaches maxRead (after finishing the current item). This prevents unbounded memory
528 // accumulation when a fast producer outpaces a slow consumer.
529 virtual kj::Maybe<jsg::Promise<DrainingReadResult>> drainingRead(
530 jsg::Lock& js, size_t maxRead = kj::maxValue) = 0;
531 
532 // The pipeTo implementation fully consumes the stream by directing all of its data at the
533 // destination. Controllers should try to be as efficient as possible here. For instance, if
534 // a ReadableStreamInternalController is piping to a WritableStreamInternalController, then
535 // a more efficient kj pipe should be possible.
536 virtual jsg::Promise<void> pipeTo(
537 jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) = 0;
538 
539 // Indicates that the consumer no longer has any interest in the streams data.
540 virtual jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) = 0;
541 
542 // Branches the ReadableStreamController into two ReadableStream instances that will receive
543 // this streams data. The specific details of how the branching occurs is entirely up to the
544 // controller implementation.
545 virtual Tee tee(jsg::Lock& js) = 0;
546 
547 virtual bool isClosedOrErrored() const = 0;
548 
549 virtual bool isClosed() const = 0;
550 
551 virtual bool isDisturbed() = 0;
552 
553 // True if a Reader has been locked to this controller.
554 virtual bool isLockedToReader() const = 0;
555 
556 // Locks this controller to the given reader, returning true if the lock was successful, or false
557 // if the controller was already locked.
558 virtual bool lockReader(jsg::Lock& js, Reader& reader) = 0;
559 
560 // Removes the lock and releases the reader from this controller.
561 // maybeJs will be nullptr when the isolate lock is not available.
562 // If maybeJs is set, the reader's closed promise will be resolved.
563 virtual void releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) = 0;
564 
565 virtual kj::Maybe<PipeController&> tryPipeLock() = 0;
566 
567 virtual void visitForGc(jsg::GcVisitor& visitor) {};
568 
569 // Fully consumes the ReadableStream. If the stream is already locked to a reader or
570 // errored, the returned JS promise will reject. If the stream is already closed, the
571 // returned JS promise will resolve with a zero-length result. Importantly, this will
572 // lock the stream and will fully consume it.
573 //
574 // limit specifies an upper maximum bound on the number of bytes permitted to be read.
575 // The promise will reject if the read will produce more bytes than the limit.
576 virtual jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit) = 0;
577 
578 // Fully consumes the ReadableStream. If the stream is already locked to a reader or
579 // errored, the returned JS promise will reject. If the stream is already closed, the
580 // returned JS promise will resolve with a zero-length result. Importantly, this will
581 // lock the stream and will fully consume it.
582 //
583 // limit specifies an upper maximum bound on the number of bytes permitted to be read.
584 // The promise will reject if the read will produce more bytes than the limit.
585 virtual jsg::Promise<kj::String> readAllText(jsg::Lock& js, uint64_t limit) = 0;
586 
587 virtual kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) = 0;
588 
589 virtual void setup(jsg::Lock& js,
590 jsg::Optional<UnderlyingSource> maybeUnderlyingSource,
591 jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) {}
592 
593 virtual kj::Promise<DeferredProxy<void>> pumpTo(
594 jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) = 0;
595 
596 // If pumpTo() pumps to a system stream, what is the best encoding for that system stream to
597 // use? This is just a hint.
598 virtual StreamEncoding getPreferredEncoding() {
599 return StreamEncoding::IDENTITY;
600 }
601 
602 virtual kj::Own<ReadableStreamController> detach(jsg::Lock& js, bool ignoreDisturbed) = 0;
603 
604 // Used by sockets to signal that the ReadableStream shouldn't allow reads due to pending
605 // closure.
606 virtual void setPendingClosure() = 0;
607 
608 virtual kj::StringPtr jsgGetMemoryName() const = 0;
609 virtual size_t jsgGetMemorySelfSize() const = 0;
610 virtual void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const = 0;
611};
612 
613kj::Own<ReadableStreamController> newReadableStreamJsController();
614kj::Own<ReadableStreamController> newReadableStreamInternalController(
615 IoContext& ioContext, kj::Own<ReadableStreamSource> source);
616 
617// A WritableStreamController provides the underlying implementation for a WritableStream.
618// We will generally have two implementations:
619// * WritableStreamDefaultController
620// * WritableStreamInternalController
621//
622// The WritableStreamDefaultController is defined by the streams standard and directs all
623// of the stream data to JavaScript functions provided by user code.
624//
625// The WritableStreamInternalController is Workers runtime specific and provides a bridge
626// to the existing WritableStreamSink API.
627//
628// The WritableStreamController instance is meant to be a private member of the WritableStream,
629// e.g.
630// class WritableStream {
631// public:
632// // ...
633// private:
634// WritableStreamController controller;
635// };
636//
637// As such, it exists within the V8 heap (it's allocated directly as a member of the
638// WritableStream) and will always execute within the V8 isolate lock.
639//
640// The methods here return jsg::Promise rather than kj::Promise because the controller
641// operations here do not always require passing through the kj mechanisms or kj event loop.
642// Likewise, we do not make use of kj::Exception in these interfaces because the stream
643// standard dictates that streams can be canceled/aborted/errored using any arbitrary JavaScript
644// value, not just Errors.
645class WritableStreamController {
646 public:
647 // The WritableStreamController::Writer interface is a base for all WritableStream writer
648 // implementations and is used solely as a means of attaching a Writer implementation to
649 // the internal state of the controller. See the WritableStream::*Writer classes for the
650 // full Writer API.
651 class Writer {
652 public:
653 // When a Writer is locked to a controller, the controller will attach itself to the writer,
654 // passing along the closed and ready promises that will be used to communicate state to the
655 // user code.
656 //
657 // The controller is guaranteed to either outlive the Writer or will detach the Writer so the
658 // WritableStreamController& reference should always remain valid.
659 virtual void attach(jsg::Lock& js,
660 WritableStreamController& controller,
661 jsg::Promise<void> closedPromise,
662 jsg::Promise<void> readyPromise) = 0;
663 
664 // When a Writer lock is released, the controller will signal to the writer that is has been
665 // detached.
666 virtual void detach() = 0;
667 
668 // The ready promise can be replaced whenever backpressure is signaled by the underlying
669 // controller.
670 virtual void replaceReadyPromise(jsg::Lock& js, jsg::Promise<void> readyPromise) = 0;
671 };
672 
673 struct PendingAbort {
674 kj::Maybe<jsg::Promise<void>::Resolver> resolver;
675 jsg::Promise<void> promise;
676 jsg::Value reason;
677 bool reject = false;
678 
679 PendingAbort(jsg::Lock& js,
680 jsg::PromiseResolverPair<void> prp,
681 v8::Local<v8::Value> reason,
682 bool reject);
683 
684 PendingAbort(jsg::Lock& js, v8::Local<v8::Value> reason, bool reject);
685 
686 void complete(jsg::Lock& js);
687 
688 void fail(jsg::Lock& js, v8::Local<v8::Value> reason);
689 
690 inline jsg::Promise<void> whenResolved(jsg::Lock& js) {
691 return promise.whenResolved(js);
692 }
693 
694 inline jsg::Promise<void> whenResolved(auto&& func) {
695 return promise.whenResolved(kj::fwd(func));
696 }
697 
698 inline jsg::Promise<void> whenResolved(auto&& func, auto&& errFunc) {
699 return promise.whenResolved(kj::fwd(func), kj::fwd(errFunc));
700 }
701 
702 void visitForGc(jsg::GcVisitor& visitor) {
703 visitor.visit(resolver, promise, reason);
704 }
705 
706 static kj::Maybe<kj::Own<PendingAbort>> dequeue(
707 kj::Maybe<kj::Own<PendingAbort>>& maybePendingAbort);
708 
709 JSG_MEMORY_INFO(PendingAbort) {
710 tracker.trackField("resolver", resolver);
711 tracker.trackField("promise", promise);
712 tracker.trackField("reason", reason);
713 }
714 };
715 
716 virtual ~WritableStreamController() noexcept(false) {}
717 
718 virtual void setOwnerRef(WritableStream& stream) = 0;
719 
720 virtual jsg::Ref<WritableStream> addRef() = 0;
721 
722 // The controller implementation will determine what kind of JavaScript data
723 // it is capable of writing, returning a rejected promise if the written
724 // data type is not supported.
725 virtual jsg::Promise<void> write(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) = 0;
726 
727 // Indicates that no additional data will be written to the controller. All
728 // existing pending writes should be allowed to complete.
729 virtual jsg::Promise<void> close(jsg::Lock& js, bool markAsHandled = false) = 0;
730 
731 // Waits for pending data to be written. The returned promise is resolved when all pending writes
732 // have completed.
733 virtual jsg::Promise<void> flush(jsg::Lock& js, bool markAsHandled = false) = 0;
734 
735 // Immediately interrupts existing pending writes and errors the stream.
736 virtual jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) = 0;
737 
738 // The tryPipeFrom attempts to establish a data pipe where source's data
739 // is delivered to this WritableStreamController as efficiently as possible.
740 virtual kj::Maybe<jsg::Promise<void>> tryPipeFrom(
741 jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) = 0;
742 
743 // Only byte-oriented WritableStreamController implementations will have a WritableStreamSink
744 // that can be detached using removeSink. A nullptr should be returned by any controller that
745 // does not support removing the sink. After the WritableStreamSink has been released, all other
746 // methods on the controller should fail with an exception as the WritableStreamSink should be
747 // the only way to interact with the underlying sink.
748 virtual kj::Maybe<kj::Own<WritableStreamSink>> removeSink(jsg::Lock& js) = 0;
749 
750 // Detaches the WritableStreamController from its underlying implementation, leaving the
751 // writable stream locked and in a state where no further writes can be made.
752 virtual void detach(jsg::Lock& js) = 0;
753 
754 virtual kj::Maybe<int> getDesiredSize() = 0;
755 
756 // True if a Writer has been locked to this controller.
757 virtual bool isLockedToWriter() const = 0;
758 
759 // Locks this controller to the given writer, returning true if the lock was successful, or false
760 // if the controller was already locked.
761 virtual bool lockWriter(jsg::Lock& js, Writer& writer) = 0;
762 
763 // Removes the lock and releases the writer from this controller.
764 // maybeJs will be nullptr when the isolate lock is not available.
765 // If maybeJs is set, the writer's closed and ready promises will be resolved.
766 virtual void releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) = 0;
767 
768 virtual kj::Maybe<v8::Local<v8::Value>> isErroring(jsg::Lock& js) = 0;
769 
770 virtual void visitForGc(jsg::GcVisitor& visitor) {};
771 
772 virtual void setup(jsg::Lock& js,
773 jsg::Optional<UnderlyingSink> underlyingSink,
774 jsg::Optional<StreamQueuingStrategy> queuingStrategy) {}
775 
776 virtual bool isClosedOrClosing() = 0;
777 virtual bool isErrored() = 0;
778 
779 // True is this controller requires ArrayBuffer(Views) to be written to it.
780 virtual bool isByteOriented() const = 0;
781 
782 // Used by sockets to signal that the WritableStream shouldn't allow writes due to pending
783 // closure.
784 virtual void setPendingClosure() = 0;
785 
786 // For memory tracking
787 virtual kj::StringPtr jsgGetMemoryName() const = 0;
788 virtual size_t jsgGetMemorySelfSize() const = 0;
789 virtual void jsgGetMemoryInfo(jsg::MemoryTracker& info) const = 0;
790};
791 
792kj::Own<WritableStreamController> newWritableStreamJsController();
793kj::Own<WritableStreamController> newWritableStreamInternalController(IoContext& ioContext,
794 kj::Own<WritableStreamSink> sink,
795 kj::Maybe<kj::Own<ByteStreamObserver>> observer,
796 kj::Maybe<uint64_t> maybeHighWaterMark = kj::none,
797 kj::Maybe<jsg::Promise<void>> maybeClosureWaitable = kj::none);
798 
799struct Unlocked {
800 static constexpr kj::StringPtr NAME KJ_UNUSED = "unlocked"_kj;
801};
802struct Locked {
803 static constexpr kj::StringPtr NAME KJ_UNUSED = "locked"_kj;
804};
805 
806// When a reader is locked to a ReadableStream, a ReaderLock instance
807// is used internally to represent the locked state in the ReadableStreamController.
808class ReaderLocked {
809 public:
810 static constexpr kj::StringPtr NAME KJ_UNUSED = "reader-locked"_kj;
811 ReaderLocked(ReadableStreamController::Reader& reader,
812 jsg::Promise<void>::Resolver closedFulfiller,
813 kj::Maybe<IoOwn<kj::Canceler>> canceler = kj::none)
814 : reader(reader),
815 closedFulfiller(kj::mv(closedFulfiller)),
816 canceler(kj::mv(canceler)) {}
817 
818 ReaderLocked(ReaderLocked&&) = default;
819 ~ReaderLocked() noexcept(false) {
820 KJ_IF_SOME(r, reader) {
821 r.detach();
822 }
823 }
824 KJ_DISALLOW_COPY(ReaderLocked);
825 
826 void visitForGc(jsg::GcVisitor& visitor) {
827 visitor.visit(closedFulfiller);
828 }
829 
830 ReadableStreamController::Reader& getReader() {
831 return KJ_ASSERT_NONNULL(reader);
832 }
833 
834 kj::Maybe<jsg::Promise<void>::Resolver>& getClosedFulfiller() {
835 return closedFulfiller;
836 }
837 
838 kj::Maybe<IoOwn<kj::Canceler>>& getCanceler() {
839 return canceler;
840 }
841 
842 void clear() {
843 reader = kj::none;
844 closedFulfiller = kj::none;
845 canceler = kj::none;
846 }
847 
848 JSG_MEMORY_INFO(ReaderLocked) {
849 tracker.trackField("closedFulfiller", closedFulfiller);
850 tracker.trackFieldWithSize("IoOwn<kj::Canceler>", sizeof(IoOwn<kj::Canceler>));
851 }
852 
853 private:
854 kj::Maybe<ReadableStreamController::Reader&> reader;
855 kj::Maybe<jsg::Promise<void>::Resolver> closedFulfiller;
856 kj::Maybe<IoOwn<kj::Canceler>> canceler;
857};
858 
859// When a writer is locked to a WritableStream, a WriterLock instance
860// is used internally to represent the locked state in the WritableStreamController.
861class WriterLocked {
862 public:
863 static constexpr kj::StringPtr NAME KJ_UNUSED = "writer-locked"_kj;
864 WriterLocked(WritableStreamController::Writer& writer,
865 jsg::Promise<void>::Resolver closedFulfiller,
866 kj::Maybe<jsg::Promise<void>::Resolver> readyFulfiller = kj::none)
867 : writer(writer),
868 closedFulfiller(kj::mv(closedFulfiller)),
869 readyFulfiller(kj::mv(readyFulfiller)) {}
870 
871 WriterLocked(WriterLocked&&) = default;
872 ~WriterLocked() noexcept(false) {
873 KJ_IF_SOME(w, writer) {
874 w.detach();
875 }
876 }
877 
878 void visitForGc(jsg::GcVisitor& visitor) {
879 visitor.visit(closedFulfiller, readyFulfiller);
880 }
881 
882 WritableStreamController::Writer& getWriter() {
883 return KJ_ASSERT_NONNULL(writer);
884 }
885 
886 kj::Maybe<jsg::Promise<void>::Resolver>& getClosedFulfiller() {
887 return closedFulfiller;
888 }
889 
890 kj::Maybe<jsg::Promise<void>::Resolver>& getReadyFulfiller() {
891 return readyFulfiller;
892 }
893 
894 void setReadyFulfiller(jsg::Lock& js, jsg::PromiseResolverPair<void>& pair) {
895 KJ_IF_SOME(w, writer) {
896 readyFulfiller = kj::mv(pair.resolver);
897 w.replaceReadyPromise(js, kj::mv(pair.promise));
898 }
899 }
900 
901 void clear() {
902 writer = kj::none;
903 closedFulfiller = kj::none;
904 readyFulfiller = kj::none;
905 }
906 
907 JSG_MEMORY_INFO(WriterLocked) {
908 tracker.trackField("closedFulfiller", closedFulfiller);
909 tracker.trackField("readyFulfiller", readyFulfiller);
910 }
911 
912 private:
913 kj::Maybe<WritableStreamController::Writer&> writer;
914 kj::Maybe<jsg::Promise<void>::Resolver> closedFulfiller;
915 kj::Maybe<jsg::Promise<void>::Resolver> readyFulfiller;
916};
917 
918template <typename T>
919void maybeResolvePromise(
920 jsg::Lock& js, kj::Maybe<typename jsg::Promise<T>::Resolver>& maybeResolver, T&& t) {
921 KJ_IF_SOME(resolver, maybeResolver) {
922 resolver.resolve(js, kj::fwd<T>(t));
923 maybeResolver = nullptr;
924 }
925}
926 
927inline void maybeResolvePromise(
928 jsg::Lock& js, kj::Maybe<jsg::Promise<void>::Resolver>& maybeResolver) {
929 KJ_IF_SOME(resolver, maybeResolver) {
930 resolver.resolve(js);
931 maybeResolver = kj::none;
932 }
933}
934 
935template <typename T>
936void maybeRejectPromise(jsg::Lock& js,
937 kj::Maybe<typename jsg::Promise<T>::Resolver>& maybeResolver,
938 v8::Local<v8::Value> reason) {
939 KJ_IF_SOME(resolver, maybeResolver) {
940 resolver.reject(js, reason);
941 maybeResolver = kj::none;
942 }
943}
944 
945template <typename T>
946jsg::Promise<T> rejectedMaybeHandledPromise(
947 jsg::Lock& js, v8::Local<v8::Value> reason, bool handled) {
948 auto prp = js.newPromiseAndResolver<T>();
949 if (handled) {
950 prp.promise.markAsHandled(js);
951 }
952 prp.resolver.reject(js, reason);
953 return kj::mv(prp.promise);
954}
955 
956inline kj::Maybe<IoContext&> tryGetIoContext() {
957 // TODO(cleanup): This function is obsolete; callers should just call IoContext::tryCurrent()
958 return IoContext::tryCurrent();
959}
960 
961} // namespace workerd::api