// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include "../basics.h" #include #include #include #if _MSC_VER using ssize_t = long long; #endif namespace workerd::api { class ReadableStream; class ReadableStreamController; class ReadableStreamSource; class ReadableStreamDefaultController; class ReadableByteStreamController; class WritableStream; class WritableStreamController; class WritableStreamSink; class WritableStreamDefaultController; class TransformStreamDefaultController; using rpc::StreamEncoding; enum class ReadAllTextOption : uint8_t { NONE = 0, NULL_TERMINATE = 1 << 0, STRIP_BOM = 1 << 1, }; inline ReadAllTextOption operator|(ReadAllTextOption a, ReadAllTextOption b) { return static_cast(static_cast(a) | static_cast(b)); } inline ReadAllTextOption& operator|=(ReadAllTextOption& a, ReadAllTextOption b) { return a = a | b; } inline bool operator&(ReadAllTextOption a, ReadAllTextOption b) { return (static_cast(a) & static_cast(b)) != 0; } static constexpr kj::byte UTF8_BOM[] = {0xEF, 0xBB, 0xBF}; static constexpr size_t UTF8_BOM_SIZE = sizeof(UTF8_BOM); inline bool hasUtf8Bom(kj::ArrayPtr data) { return data.size() >= UTF8_BOM_SIZE && memcmp(data.begin(), UTF8_BOM, UTF8_BOM_SIZE) == 0; } struct ReadResult { jsg::Optional value; bool done; JSG_STRUCT(value, done); JSG_STRUCT_TS_OVERRIDE(type ReadableStreamReadResult = | { done: false, value: R; } | { done: true; value?: undefined; } ); void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(value); } }; // Result type for draining read operations. Always returns bytes, even for value streams. // Used by DrainingReader for optimized pipe-to operations with vectored writes. // This is a C++ only type - not exposed to JavaScript. struct DrainingReadResult { kj::Array> chunks; // Multiple byte arrays for vectored writes bool done = false; // True if stream is closed/closing }; struct StreamQueuingStrategy { using SizeAlgorithm = uint64_t(v8::Local); jsg::Optional highWaterMark; jsg::Optional> size; JSG_STRUCT(highWaterMark, size); JSG_STRUCT_TS_OVERRIDE(QueuingStrategy { size?: (chunk: T) => number | bigint; }); }; struct UnderlyingSource { using Controller = kj::OneOf, jsg::Ref>; using StartAlgorithm = jsg::Promise(Controller); using PullAlgorithm = jsg::Promise(Controller); using CancelAlgorithm = jsg::Promise(v8::Local reason); // The autoAllocateChunkSize mechanism allows byte streams to operate as if a BYOB // reader is being used even if it is just a default reader. Support is optional // per the streams spec but our implementation will always enable it. Specifically, // if user code does not provide an explicit autoAllocateChunkSize, we'll assume // this default. static constexpr int DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE = 4096; // We want to increase the default auto allocate chunk size but we need to do // so carefully to avoid introducing memory regressions and causing workers to // hit OOM errors. We'll use an autogate to roll out the new default. static constexpr int DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2 = 16 * 1024; // Per the spec, the type property for the UnderlyingSource should be either // undefined, the empty string, or "bytes". When undefined, the empty string is // used as the default. When type is the empty string, the stream is considered // to be value-oriented rather than byte-oriented. jsg::Optional type; // Used only when type is equal to "bytes", the autoAllocateChunkSize defines // the size of automatically allocated buffer that is created when a default // mode read is performed on a byte-oriented ReadableStream that supports // BYOB reads. The stream standard makes this optional to support and defines // no default value. We've chosen to use a default value of 4096. If given, // the value must be greater than zero. jsg::Optional autoAllocateChunkSize; jsg::Optional> start; jsg::Optional> pull; jsg::Optional> cancel; // The expectedLength is a non-standard extension used to support specifying the // content-length when using a ReadableStream as the body of a request or response. jsg::Optional expectedLength; JSG_STRUCT(type, autoAllocateChunkSize, start, pull, cancel, expectedLength); JSG_STRUCT_TS_DEFINE(interface UnderlyingByteSource { type: "bytes"; autoAllocateChunkSize?: number; start?: (controller: ReadableByteStreamController) => void | Promise; pull?: (controller: ReadableByteStreamController) => void | Promise; cancel?: (reason: any) => void | Promise; }); JSG_STRUCT_TS_OVERRIDE( { type?: "" | undefined; autoAllocateChunkSize: never; start?: (controller: ReadableStreamDefaultController) => void | Promise; pull?: (controller: ReadableStreamDefaultController) => void | Promise; cancel?: (reason: any) => void | Promise; }); }; struct UnderlyingSink { using Controller = jsg::Ref; using StartAlgorithm = jsg::Promise(Controller); using WriteAlgorithm = jsg::Promise(v8::Local, Controller); using AbortAlgorithm = jsg::Promise(v8::Local reason); using CloseAlgorithm = jsg::Promise(); // Per the spec, the type property for the UnderlyingSink should always be either // undefined or the empty string. Any other value will trigger a TypeError. jsg::Optional type; jsg::Optional> start; jsg::Optional> write; jsg::Optional> abort; jsg::Optional> close; JSG_STRUCT(type, start, write, abort, close); // TODO(cleanup): Get rid of this override and parse the type directly in param-extractor.rs JSG_STRUCT_TS_OVERRIDE( { write?: (chunk: W, controller: WritableStreamDefaultController) => void | Promise; start?: (controller: WritableStreamDefaultController) => void | Promise; abort?: (reason: any) => void | Promise; close?: () => void | Promise; }); }; struct Transformer { using Controller = jsg::Ref; using StartAlgorithm = jsg::Promise(Controller); using TransformAlgorithm = jsg::Promise(v8::Local, Controller); using FlushAlgorithm = jsg::Promise(Controller); using CancelAlgorithm = jsg::Promise(jsg::JsValue reason); jsg::Optional readableType; jsg::Optional writableType; jsg::Optional> start; jsg::Optional> transform; jsg::Optional> flush; jsg::Optional> cancel; // The expectedLength is a non-standard extension used to support specifying the // content-length when using a TransformStream readable side as the body of a // request or response. jsg::Optional expectedLength; JSG_STRUCT(readableType, writableType, start, transform, flush, cancel, expectedLength); JSG_STRUCT_TS_OVERRIDE( { start?: (controller: TransformStreamDefaultController) => void | Promise; transform?: (chunk: I, controller: TransformStreamDefaultController) => void | Promise; flush?: (controller: TransformStreamDefaultController) => void | Promise; cancel?: (reason: any) => void | Promise; expectedLength?: number; }); }; // ReadableStreamSource and WritableStreamSink // // These are implementation interfaces for ReadableStream and WritableStream. If you just need to // use a ReadableStream or WritableStream, you can safely skip reading this. If you need to // implement a new kind of stream, read on. // In the original Workers streams implementation, a ReadableStream would have a // ReadableStreamSource backing it. Likewise, a WritableStream would have a WritableStreamSink. // The ReadableStreamSource and WritableStreamSink are kj heap objects that provide a thin // wrapper on internal native stream sources originating from within the Workers runtime. // // With implementation of full streams standard support, we introduce the new abstraction APIs // ReadableStreamController and WritableStreamController, which will provide the underlying // implementation for both ReadableStream and WritableStream, respectively. // // When creating a new kind of *internal* ReadableStream, where the data is originating internally // from a kj stream, you will still implement the ReadableStreamSource API, just as before. // Likewise, when creating a new kind of *internal* WritableStream, where the data destination is // a kj stream, you will implement the WritableStreamSink API. class WritableStreamSink { public: virtual kj::Promise write(kj::ArrayPtr buffer) KJ_WARN_UNUSED_RESULT = 0; virtual kj::Promise write( kj::ArrayPtr> pieces) KJ_WARN_UNUSED_RESULT = 0; virtual kj::Promise end() KJ_WARN_UNUSED_RESULT = 0; // Must call to flush and finish the stream. virtual kj::Maybe>> tryPumpFrom( ReadableStreamSource& input, bool end); virtual void abort(kj::Exception reason) = 0; // TODO(conform): abort() should return a promise after which closed fulfillers should be // rejected. This may necessitate an "erroring" state. // Tells the sink that it is no longer to be responsible for encoding in the correct format. // Instead, the caller takes responsibility. The expected encoding is returned; the caller // promises that all future writes will use this encoding. The default implementation returns // IDENTITY, which is always correct since that's the encoding write()s should have used if // this weren't called at all. virtual StreamEncoding disownEncodingResponsibility() { return StreamEncoding::IDENTITY; } }; class ReadableStreamSource { public: virtual kj::Promise tryRead(void* buffer, size_t minBytes, size_t maxBytes) = 0; // The ReadableStreamSource version of pumpTo() has no `amount` parameter, since the Streams spec // only defines pumping everything. // // If `end` is true, then `output.end()` will be called after pumping. Note that it's especially // important to take advantage of this when using deferred proxying since calling `end()` // directly might attempt to use the `IoContext` to call `registerPendingEvent()`. virtual kj::Promise> pumpTo(WritableStreamSink& output, bool end); // If pumpTo() pumps to a system stream, what is the best encoding for that system stream to // use? This is just a hint. virtual StreamEncoding getPreferredEncoding() { return StreamEncoding::IDENTITY; }; virtual kj::Maybe tryGetLength(StreamEncoding encoding); kj::Promise> readAllBytes(uint64_t limit); kj::Promise readAllText( uint64_t limit, ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE); // Hook to inform this ReadableStreamSource that the ReadableStream has been canceled. This only // really means anything to TransformStreams, which are supposed to propagate the error to the // writable side, and custom ReadableStreams, which we don't implement yet. // // NOTE: By "propagate the error back to the writable stream", I mean: if the WritableStream is in // the Writable state, set it to the Errored state and reject its closed fulfiller with // `reason`. I'm not sure how I'm going to do this yet. virtual void cancel(kj::Exception reason); // TODO(conform): Should return promise. // // TODO(conform): `reason` should be allowed to be any JS value, and not just an exception. // That is, something silly like `stream.cancel(42)` should be allowed and trigger a // rejection with the integer `42`. struct Tee { kj::Own branches[2]; }; // Implement this if your ReadableStreamSource has a better way to tee a stream than the naive // method, which relies upon `tryRead()`. The default implementation returns nullptr. virtual kj::Maybe tryTee(uint64_t limit); }; struct PipeToOptions { jsg::Optional preventAbort; jsg::Optional preventCancel; jsg::Optional preventClose; jsg::Optional> signal; JSG_STRUCT(preventAbort, preventCancel, preventClose, signal); JSG_STRUCT_TS_OVERRIDE(StreamPipeOptions); // An additional, internal only property that is used to indicate // when the pipe operation is used for a pipeThrough rather than // a pipeTo. We use this information, for instance, to identify // when we should mark returned rejected promises as handled. bool pipeThrough = false; }; namespace StreamStates { struct Closed { static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; }; using Errored = jsg::Value; struct Erroring { static constexpr kj::StringPtr NAME KJ_UNUSED = "erroring"_kj; jsg::Value reason; Erroring(jsg::Value reason): reason(kj::mv(reason)) {} void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(reason); } }; } // namespace StreamStates // A ReadableStreamController provides the underlying implementation for a ReadableStream. // We will generally have three implementations: // * ReadableStreamDefaultController // * ReadableByteStreamController // * ReadableStreamInternalController // // The ReadableStreamDefaultController and ReadableByteStreamController are defined by the // streams standard and source all of the stream data from JavaScript functions provided by // user code. // // The ReadableStreamInternalController is Workers runtime specific and provides a bridge // to the existing ReadableStreamSource API. At the API contract layer, the // ReadableByteStreamController and ReadableStreamInternalController will appear to be // identical. Internally, however, they will be very different from one another. // // The ReadableStreamController instance is meant to be a private member of the ReadableStream, // e.g. // class ReadableStream { // public: // // ... // private: // ReadableStreamController controller; // // ... // } // // As such, it exists within the V8 heap (it's allocated directly as a member of the // ReadableStream) and will always execute within the V8 isolate lock. // // The methods here return jsg::Promise rather than kj::Promise because the controller // operations here do not always require passing through the kj mechanisms or kj event loop. // Likewise, we do not make use of kj::Exception in these interfaces because the stream // standard dictates that streams can be canceled/aborted/errored using any arbitrary JavaScript // value, not just Errors. class ReadableStreamController { public: // The ReadableStreamController::Reader interface is a base for all ReadableStream reader // implementations and is used solely as a means of attaching a Reader implementation to // the internal state of the controller. See the ReadableStream::*Reader classes for the // full Reader API. class Reader { public: // True if the reader is a BYOB reader. virtual bool isByteOriented() const = 0; // When a Reader is locked to a controller, the controller will attach itself to the reader, // passing along the closed promise that will be used to communicate state to the // user code. // // The Reader will hold a reference to the controller that will be cleared when the reader // is released or destroyed. The controller is guaranteed to either outlive or detach the // reader so the ReadableStreamController& reference should remain valid. virtual void attach(ReadableStreamController& controller, jsg::Promise closedPromise) = 0; // When a Reader lock is released, the controller will signal to the reader that it has been // detached. virtual void detach() = 0; }; struct ByobOptions { static constexpr size_t DEFAULT_AT_LEAST = 1; jsg::V8Ref bufferView; size_t byteOffset = 0; size_t byteLength; // The minimum number of elements that should be read. When not specified, the default // is DEFAULT_AT_LEAST. This is a non-standard, Workers-specific extension to // support the readAtLeast method on the ReadableStreamBYOBReader object. // ReaderImpl::read() converts this to bytes by multiplying by element size. kj::Maybe atLeast = DEFAULT_AT_LEAST; // True if the given buffer should be detached. Per the spec, we should always be // detaching a BYOB buffer but the original Workers implementation did not. // To avoid breaking backwards compatibility, a compatibility flag is provided to turn // detach on/off as appropriate. bool detachBuffer = true; }; struct Tee { jsg::Ref branch1; jsg::Ref branch2; }; // Abstract API for ReadableStreamController implementations that provide their own // tee implementations that are not backed by kj's tee. Each branch of the tee uses // the TeeController to interface with the shared underlying source, and the // TeeController ensures that each Branch receives the data that is read. class TeeController { public: // Represents an individual ReadableStreamController tee branch registered with // a TeeController. One or more branches is registered with the TeeController. class Branch { public: virtual ~Branch() noexcept(false) {} virtual void doClose(jsg::Lock& js) = 0; virtual void doError(jsg::Lock& js, v8::Local reason) = 0; virtual void handleData(jsg::Lock& js, ReadResult result) = 0; }; class BranchPtr { public: inline BranchPtr(Branch* branch): inner(branch) { KJ_ASSERT(inner != nullptr); } BranchPtr(BranchPtr&& other) = default; BranchPtr& operator=(BranchPtr&&) = default; BranchPtr(BranchPtr& other) = default; inline void doClose(jsg::Lock& js) { inner->doClose(js); } inline void doError(jsg::Lock& js, v8::Local reason) { inner->doError(js, reason); } inline void handleData(jsg::Lock& js, ReadResult result) { inner->handleData(js, kj::mv(result)); } inline uint hashCode() { return kj::hashCode(inner); } inline bool operator==(BranchPtr& other) const { return inner == other.inner; } private: Branch* inner; }; virtual ~TeeController() noexcept(false) {} virtual void addBranch(Branch* branch) = 0; virtual void close(jsg::Lock& js) = 0; virtual void error(jsg::Lock& js, v8::Local reason) = 0; virtual void ensurePulling(jsg::Lock& js) = 0; // maybeJs will be nullptr when the isolate lock is not available. // If maybeJs is set, any operations pending for the branch will be canceled. virtual void removeBranch(Branch* branch, kj::Maybe maybeJs) = 0; }; // The PipeController simplifies the abstraction between ReadableStreamController // and WritableStreamController so that the pipeTo/pipeThrough/tryPipeTo can work // without caring about what kind of controller it is working with. class PipeController { public: virtual ~PipeController() noexcept(false) {} virtual bool isClosed() = 0; virtual kj::Maybe> tryGetErrored(jsg::Lock& js) = 0; virtual void cancel(jsg::Lock& js, v8::Local reason) = 0; virtual void close(jsg::Lock& js) = 0; virtual void error(jsg::Lock& js, v8::Local reason) = 0; virtual void release(jsg::Lock& js, kj::Maybe> maybeError = kj::none) = 0; virtual kj::Maybe> tryPumpTo(WritableStreamSink& sink, bool end) = 0; virtual jsg::Promise read(jsg::Lock& js) = 0; }; virtual ~ReadableStreamController() noexcept(false) {} virtual void setOwnerRef(ReadableStream& stream) = 0; virtual jsg::Ref addRef() = 0; // Returns true if the underlying source for this controller is byte-oriented and // therefore supports the pull into API. When false, the stream can be used to pass // any arbitrary JavaScript value through. virtual bool isByteOriented() const = 0; // Reads data from the stream. If the stream is byte-oriented, then the ByobOptions can be // specified to provide a v8::ArrayBuffer to be filled by the read operation. If the ByobOptions // are provided and the stream is not byte-oriented, the operation will return a rejected promise. virtual kj::Maybe> read( jsg::Lock& js, kj::Maybe byobOptions) = 0; // Performs a draining read operation that: // 1. Drains all currently buffered data from the queue // 2. Pumps the controller for synchronously available data (respecting pull promise state) // 3. Returns bytes even for value streams (converting ArrayBuffer/ArrayBufferView/string) // 4. Has mutual exclusion with regular reads - returns rejected promise if pending regular reads // 5. Returns done: true with final data when stream is closing // // This is a C++ only API (not exposed to JavaScript) intended for optimized pipe operations. // Returns kj::none if the stream is locked in a way that prevents the read. // // The maxRead parameter provides a soft limit on how much data to read. Both the initial // buffer drain and subsequent synchronous pump attempts stop when the total bytes read // reaches maxRead (after finishing the current item). This prevents unbounded memory // accumulation when a fast producer outpaces a slow consumer. virtual kj::Maybe> drainingRead( jsg::Lock& js, size_t maxRead = kj::maxValue) = 0; // The pipeTo implementation fully consumes the stream by directing all of its data at the // destination. Controllers should try to be as efficient as possible here. For instance, if // a ReadableStreamInternalController is piping to a WritableStreamInternalController, then // a more efficient kj pipe should be possible. virtual jsg::Promise pipeTo( jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) = 0; // Indicates that the consumer no longer has any interest in the streams data. virtual jsg::Promise cancel(jsg::Lock& js, jsg::Optional> reason) = 0; // Branches the ReadableStreamController into two ReadableStream instances that will receive // this streams data. The specific details of how the branching occurs is entirely up to the // controller implementation. virtual Tee tee(jsg::Lock& js) = 0; virtual bool isClosedOrErrored() const = 0; virtual bool isClosed() const = 0; virtual bool isDisturbed() = 0; // True if a Reader has been locked to this controller. virtual bool isLockedToReader() const = 0; // Locks this controller to the given reader, returning true if the lock was successful, or false // if the controller was already locked. virtual bool lockReader(jsg::Lock& js, Reader& reader) = 0; // Removes the lock and releases the reader from this controller. // maybeJs will be nullptr when the isolate lock is not available. // If maybeJs is set, the reader's closed promise will be resolved. virtual void releaseReader(Reader& reader, kj::Maybe maybeJs) = 0; virtual kj::Maybe tryPipeLock() = 0; virtual void visitForGc(jsg::GcVisitor& visitor) {}; // Fully consumes the ReadableStream. If the stream is already locked to a reader or // errored, the returned JS promise will reject. If the stream is already closed, the // returned JS promise will resolve with a zero-length result. Importantly, this will // lock the stream and will fully consume it. // // limit specifies an upper maximum bound on the number of bytes permitted to be read. // The promise will reject if the read will produce more bytes than the limit. virtual jsg::Promise readAllBytes(jsg::Lock& js, uint64_t limit) = 0; // Fully consumes the ReadableStream. If the stream is already locked to a reader or // errored, the returned JS promise will reject. If the stream is already closed, the // returned JS promise will resolve with a zero-length result. Importantly, this will // lock the stream and will fully consume it. // // limit specifies an upper maximum bound on the number of bytes permitted to be read. // The promise will reject if the read will produce more bytes than the limit. virtual jsg::Promise readAllText(jsg::Lock& js, uint64_t limit) = 0; virtual kj::Maybe tryGetLength(StreamEncoding encoding) = 0; virtual void setup(jsg::Lock& js, jsg::Optional maybeUnderlyingSource, jsg::Optional maybeQueuingStrategy) {} virtual kj::Promise> pumpTo( jsg::Lock& js, kj::Own sink, bool end) = 0; // If pumpTo() pumps to a system stream, what is the best encoding for that system stream to // use? This is just a hint. virtual StreamEncoding getPreferredEncoding() { return StreamEncoding::IDENTITY; } virtual kj::Own detach(jsg::Lock& js, bool ignoreDisturbed) = 0; // Used by sockets to signal that the ReadableStream shouldn't allow reads due to pending // closure. virtual void setPendingClosure() = 0; virtual kj::StringPtr jsgGetMemoryName() const = 0; virtual size_t jsgGetMemorySelfSize() const = 0; virtual void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const = 0; }; kj::Own newReadableStreamJsController(); kj::Own newReadableStreamInternalController( IoContext& ioContext, kj::Own source); // A WritableStreamController provides the underlying implementation for a WritableStream. // We will generally have two implementations: // * WritableStreamDefaultController // * WritableStreamInternalController // // The WritableStreamDefaultController is defined by the streams standard and directs all // of the stream data to JavaScript functions provided by user code. // // The WritableStreamInternalController is Workers runtime specific and provides a bridge // to the existing WritableStreamSink API. // // The WritableStreamController instance is meant to be a private member of the WritableStream, // e.g. // class WritableStream { // public: // // ... // private: // WritableStreamController controller; // }; // // As such, it exists within the V8 heap (it's allocated directly as a member of the // WritableStream) and will always execute within the V8 isolate lock. // // The methods here return jsg::Promise rather than kj::Promise because the controller // operations here do not always require passing through the kj mechanisms or kj event loop. // Likewise, we do not make use of kj::Exception in these interfaces because the stream // standard dictates that streams can be canceled/aborted/errored using any arbitrary JavaScript // value, not just Errors. class WritableStreamController { public: // The WritableStreamController::Writer interface is a base for all WritableStream writer // implementations and is used solely as a means of attaching a Writer implementation to // the internal state of the controller. See the WritableStream::*Writer classes for the // full Writer API. class Writer { public: // When a Writer is locked to a controller, the controller will attach itself to the writer, // passing along the closed and ready promises that will be used to communicate state to the // user code. // // The controller is guaranteed to either outlive the Writer or will detach the Writer so the // WritableStreamController& reference should always remain valid. virtual void attach(jsg::Lock& js, WritableStreamController& controller, jsg::Promise closedPromise, jsg::Promise readyPromise) = 0; // When a Writer lock is released, the controller will signal to the writer that is has been // detached. virtual void detach() = 0; // The ready promise can be replaced whenever backpressure is signaled by the underlying // controller. virtual void replaceReadyPromise(jsg::Lock& js, jsg::Promise readyPromise) = 0; }; struct PendingAbort { kj::Maybe::Resolver> resolver; jsg::Promise promise; jsg::Value reason; bool reject = false; PendingAbort(jsg::Lock& js, jsg::PromiseResolverPair prp, v8::Local reason, bool reject); PendingAbort(jsg::Lock& js, v8::Local reason, bool reject); void complete(jsg::Lock& js); void fail(jsg::Lock& js, v8::Local reason); inline jsg::Promise whenResolved(jsg::Lock& js) { return promise.whenResolved(js); } inline jsg::Promise whenResolved(auto&& func) { return promise.whenResolved(kj::fwd(func)); } inline jsg::Promise whenResolved(auto&& func, auto&& errFunc) { return promise.whenResolved(kj::fwd(func), kj::fwd(errFunc)); } void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(resolver, promise, reason); } static kj::Maybe> dequeue( kj::Maybe>& maybePendingAbort); JSG_MEMORY_INFO(PendingAbort) { tracker.trackField("resolver", resolver); tracker.trackField("promise", promise); tracker.trackField("reason", reason); } }; virtual ~WritableStreamController() noexcept(false) {} virtual void setOwnerRef(WritableStream& stream) = 0; virtual jsg::Ref addRef() = 0; // The controller implementation will determine what kind of JavaScript data // it is capable of writing, returning a rejected promise if the written // data type is not supported. virtual jsg::Promise write(jsg::Lock& js, jsg::Optional> value) = 0; // Indicates that no additional data will be written to the controller. All // existing pending writes should be allowed to complete. virtual jsg::Promise close(jsg::Lock& js, bool markAsHandled = false) = 0; // Waits for pending data to be written. The returned promise is resolved when all pending writes // have completed. virtual jsg::Promise flush(jsg::Lock& js, bool markAsHandled = false) = 0; // Immediately interrupts existing pending writes and errors the stream. virtual jsg::Promise abort(jsg::Lock& js, jsg::Optional> reason) = 0; // The tryPipeFrom attempts to establish a data pipe where source's data // is delivered to this WritableStreamController as efficiently as possible. virtual kj::Maybe> tryPipeFrom( jsg::Lock& js, jsg::Ref source, PipeToOptions options) = 0; // Only byte-oriented WritableStreamController implementations will have a WritableStreamSink // that can be detached using removeSink. A nullptr should be returned by any controller that // does not support removing the sink. After the WritableStreamSink has been released, all other // methods on the controller should fail with an exception as the WritableStreamSink should be // the only way to interact with the underlying sink. virtual kj::Maybe> removeSink(jsg::Lock& js) = 0; // Detaches the WritableStreamController from its underlying implementation, leaving the // writable stream locked and in a state where no further writes can be made. virtual void detach(jsg::Lock& js) = 0; virtual kj::Maybe getDesiredSize() = 0; // True if a Writer has been locked to this controller. virtual bool isLockedToWriter() const = 0; // Locks this controller to the given writer, returning true if the lock was successful, or false // if the controller was already locked. virtual bool lockWriter(jsg::Lock& js, Writer& writer) = 0; // Removes the lock and releases the writer from this controller. // maybeJs will be nullptr when the isolate lock is not available. // If maybeJs is set, the writer's closed and ready promises will be resolved. virtual void releaseWriter(Writer& writer, kj::Maybe maybeJs) = 0; virtual kj::Maybe> isErroring(jsg::Lock& js) = 0; virtual void visitForGc(jsg::GcVisitor& visitor) {}; virtual void setup(jsg::Lock& js, jsg::Optional underlyingSink, jsg::Optional queuingStrategy) {} virtual bool isClosedOrClosing() = 0; virtual bool isErrored() = 0; // True is this controller requires ArrayBuffer(Views) to be written to it. virtual bool isByteOriented() const = 0; // Used by sockets to signal that the WritableStream shouldn't allow writes due to pending // closure. virtual void setPendingClosure() = 0; // For memory tracking virtual kj::StringPtr jsgGetMemoryName() const = 0; virtual size_t jsgGetMemorySelfSize() const = 0; virtual void jsgGetMemoryInfo(jsg::MemoryTracker& info) const = 0; }; kj::Own newWritableStreamJsController(); kj::Own newWritableStreamInternalController(IoContext& ioContext, kj::Own sink, kj::Maybe> observer, kj::Maybe maybeHighWaterMark = kj::none, kj::Maybe> maybeClosureWaitable = kj::none); struct Unlocked { static constexpr kj::StringPtr NAME KJ_UNUSED = "unlocked"_kj; }; struct Locked { static constexpr kj::StringPtr NAME KJ_UNUSED = "locked"_kj; }; // When a reader is locked to a ReadableStream, a ReaderLock instance // is used internally to represent the locked state in the ReadableStreamController. class ReaderLocked { public: static constexpr kj::StringPtr NAME KJ_UNUSED = "reader-locked"_kj; ReaderLocked(ReadableStreamController::Reader& reader, jsg::Promise::Resolver closedFulfiller, kj::Maybe> canceler = kj::none) : reader(reader), closedFulfiller(kj::mv(closedFulfiller)), canceler(kj::mv(canceler)) {} ReaderLocked(ReaderLocked&&) = default; ~ReaderLocked() noexcept(false) { KJ_IF_SOME(r, reader) { r.detach(); } } KJ_DISALLOW_COPY(ReaderLocked); void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(closedFulfiller); } ReadableStreamController::Reader& getReader() { return KJ_ASSERT_NONNULL(reader); } kj::Maybe::Resolver>& getClosedFulfiller() { return closedFulfiller; } kj::Maybe>& getCanceler() { return canceler; } void clear() { reader = kj::none; closedFulfiller = kj::none; canceler = kj::none; } JSG_MEMORY_INFO(ReaderLocked) { tracker.trackField("closedFulfiller", closedFulfiller); tracker.trackFieldWithSize("IoOwn", sizeof(IoOwn)); } private: kj::Maybe reader; kj::Maybe::Resolver> closedFulfiller; kj::Maybe> canceler; }; // When a writer is locked to a WritableStream, a WriterLock instance // is used internally to represent the locked state in the WritableStreamController. class WriterLocked { public: static constexpr kj::StringPtr NAME KJ_UNUSED = "writer-locked"_kj; WriterLocked(WritableStreamController::Writer& writer, jsg::Promise::Resolver closedFulfiller, kj::Maybe::Resolver> readyFulfiller = kj::none) : writer(writer), closedFulfiller(kj::mv(closedFulfiller)), readyFulfiller(kj::mv(readyFulfiller)) {} WriterLocked(WriterLocked&&) = default; ~WriterLocked() noexcept(false) { KJ_IF_SOME(w, writer) { w.detach(); } } void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(closedFulfiller, readyFulfiller); } WritableStreamController::Writer& getWriter() { return KJ_ASSERT_NONNULL(writer); } kj::Maybe::Resolver>& getClosedFulfiller() { return closedFulfiller; } kj::Maybe::Resolver>& getReadyFulfiller() { return readyFulfiller; } void setReadyFulfiller(jsg::Lock& js, jsg::PromiseResolverPair& pair) { KJ_IF_SOME(w, writer) { readyFulfiller = kj::mv(pair.resolver); w.replaceReadyPromise(js, kj::mv(pair.promise)); } } void clear() { writer = kj::none; closedFulfiller = kj::none; readyFulfiller = kj::none; } JSG_MEMORY_INFO(WriterLocked) { tracker.trackField("closedFulfiller", closedFulfiller); tracker.trackField("readyFulfiller", readyFulfiller); } private: kj::Maybe writer; kj::Maybe::Resolver> closedFulfiller; kj::Maybe::Resolver> readyFulfiller; }; template void maybeResolvePromise( jsg::Lock& js, kj::Maybe::Resolver>& maybeResolver, T&& t) { KJ_IF_SOME(resolver, maybeResolver) { resolver.resolve(js, kj::fwd(t)); maybeResolver = nullptr; } } inline void maybeResolvePromise( jsg::Lock& js, kj::Maybe::Resolver>& maybeResolver) { KJ_IF_SOME(resolver, maybeResolver) { resolver.resolve(js); maybeResolver = kj::none; } } template void maybeRejectPromise(jsg::Lock& js, kj::Maybe::Resolver>& maybeResolver, v8::Local reason) { KJ_IF_SOME(resolver, maybeResolver) { resolver.reject(js, reason); maybeResolver = kj::none; } } template jsg::Promise rejectedMaybeHandledPromise( jsg::Lock& js, v8::Local reason, bool handled) { auto prp = js.newPromiseAndResolver(); if (handled) { prp.promise.markAsHandled(js); } prp.resolver.reject(js, reason); return kj::mv(prp.promise); } inline kj::Maybe tryGetIoContext() { // TODO(cleanup): This function is obsolete; callers should just call IoContext::tryCurrent() return IoContext::tryCurrent(); } } // namespace workerd::api