// 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 "common.h" #include #include namespace workerd::api { class ReadableStreamDefaultReader; class ReadableStreamBYOBReader; class ReaderImpl final { public: ReaderImpl(ReadableStreamController::Reader& reader); ~ReaderImpl() noexcept(false); void attach(ReadableStreamController& controller, jsg::Promise closedPromise); jsg::Promise cancel(jsg::Lock& js, jsg::Optional> maybeReason); void detach(); jsg::MemoizedIdentity>& getClosed(); void lockToStream(jsg::Lock& js, ReadableStream& stream); jsg::Promise read(jsg::Lock& js, kj::Maybe byobOptions); void releaseLock(jsg::Lock& js); void visitForGc(jsg::GcVisitor& visitor); kj::StringPtr jsgGetMemoryName() const; size_t jsgGetMemorySelfSize() const; void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const; private: struct Initial { static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; }; // While a Reader is attached to a ReadableStream, it holds a strong reference to the // ReadableStream to prevent it from being GC'ed so long as the Reader is available. // Once the reader is closed, released, or GC'ed the reference to the ReadableStream // is cleared and the ReadableStream can be GC'ed if there are no other references to // it being held anywhere. If the reader is still attached to the ReadableStream when // it is destroyed, the ReadableStream's reference to the reader is cleared but the // ReadableStream remains in the "reader locked" state, per the spec. struct Attached { static constexpr kj::StringPtr NAME KJ_UNUSED = "attached"_kj; jsg::Ref stream; }; // Released: The user explicitly called releaseLock() to detach the reader from the stream. // The stream remains usable and can be locked by a new reader. struct Released { static constexpr kj::StringPtr NAME KJ_UNUSED = "released"_kj; }; // Closed: The underlying stream ended (closed or errored) while the reader was attached. // The stream is no longer usable. struct Closed { static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; }; // State machine for ReaderImpl: // Initial -> Attached (attach() called) // Attached -> Closed (detach() called when stream closes) // Attached -> Released (releaseLock() called) // Closed and Released are terminal states. // Initial is not terminal but most methods assert if called in this state. using ReaderState = StateMachine, ActiveState, Initial, Attached, Closed, Released>; kj::Maybe ioContext; ReadableStreamController::Reader& reader; ReaderState state; inline void assertAttachedOrTerminal() const { KJ_ASSERT(!state.is(), "this reader was never attached"); } kj::Maybe>> closedPromise; friend class ReadableStreamDefaultReader; friend class ReadableStreamBYOBReader; }; class ReadableStreamDefaultReader : public jsg::Object, public ReadableStreamController::Reader { public: explicit ReadableStreamDefaultReader(); // JavaScript API static jsg::Ref constructor( jsg::Lock& js, jsg::Ref stream); jsg::MemoizedIdentity>& getClosed(); jsg::Promise cancel(jsg::Lock& js, jsg::Optional> reason); jsg::Promise read(jsg::Lock& js); void releaseLock(jsg::Lock& js); JSG_RESOURCE_TYPE(ReadableStreamDefaultReader, CompatibilityFlags::Reader flags) { if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed); } else { JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed); } JSG_METHOD(cancel); JSG_METHOD(read); JSG_METHOD(releaseLock); JSG_TS_OVERRIDE( { read(): Promise>; }); } // Internal API void attach(ReadableStreamController& controller, jsg::Promise closedPromise) override; void detach() override; void lockToStream(jsg::Lock& js, ReadableStream& stream); inline bool isByteOriented() const override { return false; } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("impl", impl); } private: ReaderImpl impl; void visitForGc(jsg::GcVisitor& visitor); }; class ReadableStreamBYOBReader: public jsg::Object, public ReadableStreamController::Reader { public: explicit ReadableStreamBYOBReader(); // JavaScript API static jsg::Ref constructor( jsg::Lock& js, jsg::Ref stream); jsg::MemoizedIdentity>& getClosed(); jsg::Promise cancel(jsg::Lock& js, jsg::Optional> reason); struct ReadableStreamBYOBReaderReadOptions { jsg::Optional min; JSG_STRUCT(min); }; jsg::Promise read(jsg::Lock& js, v8::Local byobBuffer, jsg::Optional options = kj::none); // Non-standard extension so that reads can specify a minimum number of elements to read. It's a // struct so that we could eventually add things like timeouts if we need to. Since there's no // existing spec that's a leading contender, this is behind a different method name to avoid // conflicts with any changes to `read`. Fewer than `minElements` may be returned if EOF is hit // or the underlying stream is closed/errors out. In all cases the read result is either // {value: theChunk, done: false} or {value: undefined, done: true} as with read. // TODO(soon): Like fetch() and Cache.match(), readAtLeast() returns a promise for a V8 object. jsg::Promise readAtLeast(jsg::Lock& js, int minElements, v8::Local byobBuffer); void releaseLock(jsg::Lock& js); JSG_RESOURCE_TYPE(ReadableStreamBYOBReader, CompatibilityFlags::Reader flags) { if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed); } else { JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed); } JSG_METHOD(cancel); JSG_METHOD(read); JSG_METHOD(releaseLock); // Non-standard extension that should only apply to BYOB byte streams. JSG_METHOD(readAtLeast); JSG_TS_OVERRIDE(ReadableStreamBYOBReader { read(view: T): Promise>; readAtLeast(minElements: number, view: T): Promise>; }); } // Internal API void attach( ReadableStreamController& controller, jsg::Promise closedPromise) override; void detach() override; void lockToStream(jsg::Lock& js, ReadableStream& stream); inline bool isByteOriented() const override { return true; } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("impl", impl); } private: ReaderImpl impl; void visitForGc(jsg::GcVisitor& visitor); }; // DrainingReader is a C++ only reader (not exposed to JavaScript) that performs // draining reads. It locks the stream like standard readers but uses drainingRead() // instead of regular read() to drain all synchronously available data at once. // This is intended for optimized pipe operations. class DrainingReader: public ReadableStreamController::Reader { public: explicit DrainingReader(); // Factory method to create and lock to a stream. Returns nullptr if stream is locked. static kj::Maybe> create(jsg::Lock& js, ReadableStream& stream); virtual ~DrainingReader() noexcept(false); // Performs a draining read, returning all synchronously available data as bytes. // The maxRead parameter is a soft limit - see ReadableStreamController::drainingRead. jsg::Promise read(jsg::Lock& js, size_t maxRead = kj::maxValue); // Cancels the stream. jsg::Promise cancel(jsg::Lock& js, jsg::Optional> maybeReason); // Releases the lock on the stream. void releaseLock(jsg::Lock& js); // Returns whether this reader is still attached to a stream. bool isAttached() const; // ReadableStreamController::Reader interface void attach(ReadableStreamController& controller, jsg::Promise closedPromise) override; void detach() override; bool isByteOriented() const override { return false; } void visitForGc(jsg::GcVisitor& visitor); private: struct Initial {}; using Attached = jsg::Ref; struct Released {}; kj::Maybe ioContext; kj::OneOf state = Initial(); kj::Maybe>> closedPromise; }; class ReadableStream: public jsg::Object { private: struct AsyncIteratorState { kj::Maybe ioContext; jsg::Ref reader; bool preventCancel; }; static jsg::Promise> nextFunction( jsg::Lock& js, AsyncIteratorState& state); static jsg::Promise returnFunction( jsg::Lock& js, AsyncIteratorState& state, jsg::Optional& value); public: explicit ReadableStream(IoContext& ioContext, kj::Own source); explicit ReadableStream(kj::Own controller); ReadableStreamController& getController(); jsg::Ref addRef(); bool isDisturbed(); // --------------------------------------------------------------------------- // JS interface // Creates a new JS-backed ReadableStream using the provided source and strategy. // We use v8::Local's here instead of jsg structs because we need // to preserve the object references within the implementation. static jsg::Ref constructor( jsg::Lock& js, jsg::Optional underlyingSource, jsg::Optional queuingStrategy); static jsg::Ref from(jsg::Lock& js, jsg::AsyncGenerator generator); bool isLocked(); // Closes the stream. All present and future read requests are fulfilled with successful empty // results. `reason` will be passed to the underlying source's cancel algorithm -- if this // readable stream is one side of a transform stream, then its cancel algorithm causes the // transform's writable side to become errored with `reason`. jsg::Promise cancel(jsg::Lock& js, jsg::Optional> reason); using Reader = kj::OneOf, jsg::Ref>; struct GetReaderOptions { jsg::Optional mode; // can be "byob" or undefined JSG_STRUCT(mode); JSG_STRUCT_TS_OVERRIDE({ mode: "byob" }); // Intentionally required, so we can use `GetReaderOptions` directly in the // `ReadableStream#getReader()` overload. }; Reader getReader(jsg::Lock& js, jsg::Optional options); // Options specifically for the values() function. struct ValuesOptions { jsg::Optional preventCancel = false; JSG_STRUCT(preventCancel); }; JSG_ASYNC_ITERATOR_WITH_OPTIONS(ReadableStreamAsyncIterator, values, jsg::Value, AsyncIteratorState, nextFunction, returnFunction, ValuesOptions); struct Transform { jsg::Ref readable; jsg::Ref writable; JSG_STRUCT(readable, writable); JSG_STRUCT_TS_OVERRIDE(ReadableWritablePair { readable: ReadableStream; writable: WritableStream; }); }; jsg::Ref pipeThrough( jsg::Lock& js, Transform transform, jsg::Optional options); jsg::Promise pipeTo( jsg::Lock& js, jsg::Ref destination, jsg::Optional options); // Locks the stream and returns a pair of two new ReadableStreams, each of which read the same // data as this ReadableStream would. kj::Array> tee(jsg::Lock& js); jsg::JsString inspectState(jsg::Lock& js); bool inspectSupportsBYOB(); jsg::Optional inspectLength(); JSG_RESOURCE_TYPE(ReadableStream, CompatibilityFlags::Reader flags) { if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_READONLY_PROTOTYPE_PROPERTY(locked, isLocked); } else { JSG_READONLY_INSTANCE_PROPERTY(locked, isLocked); } JSG_METHOD(cancel); JSG_METHOD(getReader); JSG_METHOD(pipeThrough); JSG_METHOD(pipeTo); JSG_METHOD(tee); JSG_METHOD(values); JSG_STATIC_METHOD(from); JSG_INSPECT_PROPERTY(state, inspectState); JSG_INSPECT_PROPERTY(supportsBYOB, inspectSupportsBYOB); JSG_INSPECT_PROPERTY(length, inspectLength); JSG_ASYNC_ITERABLE(values); if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_TS_DEFINE(interface ReadableStream { get locked(): boolean; cancel(reason?: any): Promise; getReader(): ReadableStreamDefaultReader; getReader(options: ReadableStreamGetReaderOptions): ReadableStreamBYOBReader; pipeThrough(transform: ReadableWritablePair, options?: StreamPipeOptions): ReadableStream; pipeTo(destination: WritableStream, options?: StreamPipeOptions): Promise; tee(): [ReadableStream, ReadableStream]; values(options?: ReadableStreamValuesOptions): AsyncIterableIterator; [Symbol.asyncIterator](options?: ReadableStreamValuesOptions): AsyncIterableIterator; }); } else { JSG_TS_DEFINE(interface ReadableStream { readonly locked: boolean; cancel(reason?: any): Promise; getReader(): ReadableStreamDefaultReader; getReader(options: ReadableStreamGetReaderOptions): ReadableStreamBYOBReader; pipeThrough(transform: ReadableWritablePair, options?: StreamPipeOptions): ReadableStream; pipeTo(destination: WritableStream, options?: StreamPipeOptions): Promise; tee(): [ReadableStream, ReadableStream]; values(options?: ReadableStreamValuesOptions): AsyncIterableIterator; [Symbol.asyncIterator](options?: ReadableStreamValuesOptions): AsyncIterableIterator; }); } // Replace ReadableStream class with an interface and const, so we can have // two constructors with differing type parameters for byte-oriented and // value-oriented streams. JSG_TS_OVERRIDE(const ReadableStream: { prototype: ReadableStream; new (underlyingSource: UnderlyingByteSource, strategy?: QueuingStrategy): ReadableStream; new (underlyingSource?: UnderlyingSource, strategy?: QueuingStrategy): ReadableStream; }); } // Detaches this ReadableStream from its underlying controller state, returning a // new ReadableStream instance that takes over the underlying state. This is used to // support the "create a proxy" of a ReadableStream algorithm in the streams spec // (see https://streams.spec.whatwg.org/#readablestream-create-a-proxy). In that // algorithm, it says to create a proxy of a stream by creating a new TransformStream // and piping the original through it. The readable side of the created transform // becomes the proxy. That is quite inefficient so instead, we create a new // ReadableStream that will take over ownership of the internal state of this one, // leaving this ReadableStream locked and disturbed so that it is no longer usable. // The name "detach" here is used in the sense of "detaching the internal state". jsg::Ref detach(jsg::Lock& js, bool ignoreDisturbed=false); kj::Maybe tryGetLength(StreamEncoding encoding); // A potentially optimized version of pipe that sends this stream's data to the given // sink. The entire stream is consumed. The ReadableStream will be left locked and // disturbed and the DeferredProxy returned will take over ownership of the internal // state of the readable. kj::Promise> pumpTo(jsg::Lock& js, kj::Own sink, bool end); // Initializes signalling mechanism for EOF detection. Returns a promise that will resolve when // EOF is reached. // // This method should only be called once. jsg::Promise onEof(jsg::Lock& js); // Used by ReadableStreamInternalController to signal EOF being reached. Can be called even if // `onEof` wasn't called. void signalEof(jsg::Lock& js); void serialize(jsg::Lock& js, jsg::Serializer& serializer); static jsg::Ref deserialize( jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer); JSG_SERIALIZABLE(rpc::SerializationTag::READABLE_STREAM); void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; private: kj::Maybe ioContext; kj::Own controller; // Used to signal when this ReadableStream reads EOF. This signal is required for TCP sockets. kj::Maybe> eofResolverPair; void visitForGc(jsg::GcVisitor& visitor); }; struct QueuingStrategyInit { double highWaterMark; JSG_STRUCT(highWaterMark); }; using QueuingStrategySizeFunction = jsg::Optional(jsg::Optional>); // Utility class defined by the streams spec that uses byteLength to calculate // backpressure changes. class ByteLengthQueuingStrategy: public jsg::Object { public: ByteLengthQueuingStrategy(QueuingStrategyInit init) : init(init) {} static jsg::Ref constructor(jsg::Lock& js, QueuingStrategyInit init) { return js.alloc(init); } double getHighWaterMark() const { return init.highWaterMark; } jsg::Function getSize() const { return &size; } JSG_RESOURCE_TYPE(ByteLengthQueuingStrategy) { JSG_READONLY_PROTOTYPE_PROPERTY(highWaterMark, getHighWaterMark); JSG_READONLY_PROTOTYPE_PROPERTY(size, getSize); // QueuingStrategy requires the result of the size function to be defined JSG_TS_OVERRIDE(implements QueuingStrategy { get size(): (chunk?: any) => number; }); } private: static jsg::Optional size(jsg::Lock& js, jsg::Optional>); QueuingStrategyInit init; }; // Utility class defined by the streams spec that uses a fixed value of 1 to calculate // backpressure change class CountQueuingStrategy: public jsg::Object { public: CountQueuingStrategy(QueuingStrategyInit init) : init(init) {} static jsg::Ref constructor(jsg::Lock& js, QueuingStrategyInit init) { return js.alloc(init); } double getHighWaterMark() const { return init.highWaterMark; } jsg::Function getSize() const { return &size; } JSG_RESOURCE_TYPE(CountQueuingStrategy) { JSG_READONLY_PROTOTYPE_PROPERTY(highWaterMark, getHighWaterMark); JSG_READONLY_PROTOTYPE_PROPERTY(size, getSize); // QueuingStrategy requires the result of the size function to be defined JSG_TS_OVERRIDE(implements QueuingStrategy { get size(): (chunk?: any) => number; }); } private: static jsg::Optional size(jsg::Lock& js, jsg::Optional>) { return 1; } QueuingStrategyInit init; }; } // namespace workerd::api