// 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 "writable.h" #include #include #include #include #include namespace workerd::api { // ======================================================================================= // The ReadableStreamInternalController and WritableStreamInternalController provide the // internal (original) implementation of the ReadableStream/WritableStream objects and are // each backed by the ReadableStreamSource and WritableStreamSink respectively. Every stream // implementation that originates from *within* the Workers runtime will use these. // // It is important to understand that the behavior of these are not entirely compliant with // the streams specification. // The ReadableStreamInternalController is always in one of three states: Readable, Closed, // or Errored. When the state is Readable, the controller has an associated ReadableStreamSource. // When the state is Errored, the ReadableStreamSource has been released and the controller // stores a jsg::Value with whatever value was used to error. When Closed, the // ReadableStreamSource has been released. // Likewise, the WritableStreamInternalController is always either Writable, Closed, or Errored. // When the state is Writable, the controller has an associated WritableStreamSink. In either of // the other two states, the sink has been released. class WritableStreamInternalController; class ReadableStreamInternalController: public ReadableStreamController { public: using Readable = IoOwn; explicit ReadableStreamInternalController(StreamStates::Closed closed) : state(State::create()) {} explicit ReadableStreamInternalController(StreamStates::Errored errored) : state(State::create(kj::mv(errored))) {} explicit ReadableStreamInternalController(Readable readable) : state(State::create(kj::mv(readable))) {} KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamInternalController); ~ReadableStreamInternalController() noexcept(false) override; void setOwnerRef(ReadableStream& stream) override { owner = stream; } jsg::Ref addRef() override; bool isByteOriented() const override { return true; } kj::Maybe> read( jsg::Lock& js, kj::Maybe byobOptions) override; kj::Maybe> drainingRead( jsg::Lock& js, size_t maxRead = kj::maxValue) override; jsg::Promise pipeTo( jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) override; jsg::Promise cancel(jsg::Lock& js, jsg::Optional> reason) override; Tee tee(jsg::Lock& js) override; kj::Maybe> removeSource( jsg::Lock& js, bool ignoreDisturbed = false); bool isClosedOrErrored() const override { return state.is() || state.is(); } bool isClosed() const override { return state.is(); } bool isDisturbed() override { return disturbed; } bool isLockedToReader() const override { return !readState.is(); } bool lockReader(jsg::Lock& js, Reader& reader) override; void releaseReader(Reader& reader, kj::Maybe maybeJs) override; // See the comment for releaseReader in common.h for details on the use of maybeJs kj::Maybe tryPipeLock() override; void visitForGc(jsg::GcVisitor& visitor) override; jsg::Promise readAllBytes(jsg::Lock& js, uint64_t limit) override; jsg::Promise readAllText(jsg::Lock& js, uint64_t limit) override; kj::Maybe tryGetLength(StreamEncoding encoding) override; kj::Promise> pumpTo( jsg::Lock& js, kj::Own sink, bool end) override; StreamEncoding getPreferredEncoding() override; kj::Own detach(jsg::Lock& js, bool ignoreDisturbed) override; void setPendingClosure() override { isPendingClosure = true; } kj::StringPtr jsgGetMemoryName() const override; size_t jsgGetMemorySelfSize() const override; void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override; private: void doCancel(jsg::Lock& js, jsg::Optional> reason); void doClose(jsg::Lock& js); void doError(jsg::Lock& js, v8::Local reason); class PipeLocked: public PipeController { public: static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj; PipeLocked(ReadableStreamInternalController& inner): inner(inner) {} bool isClosed() override; kj::Maybe> tryGetErrored(jsg::Lock& js) override; void cancel(jsg::Lock& js, v8::Local reason) override; void close(jsg::Lock& js) override; void error(jsg::Lock& js, v8::Local reason) override; void release(jsg::Lock& js, kj::Maybe> maybeError = kj::none) override; kj::Maybe> tryPumpTo(WritableStreamSink& sink, bool end) override; jsg::Promise read(jsg::Lock& js) override; private: ReadableStreamInternalController& inner; }; kj::Maybe owner; // State machine for ReadableStreamInternalController: // Closed is terminal, Errored is implicitly terminal via ErrorState. // Readable is the active state (stream has data). using State = StateMachine, ErrorState, ActiveState, StreamStates::Closed, StreamStates::Errored, Readable>; State state; // Lock state machine for ReadableStreamInternalController: // All states can transition to any other state (no terminal states). // Unlocked -> Locked (removeSink() or pumpTo() called) // Unlocked -> ReaderLocked (lockReader() called) // Unlocked -> PipeLocked (tryPipeLock() called) // ReaderLocked -> Unlocked (releaseReader() called) // PipeLocked -> Unlocked (release() or doClose/doError called) // Locked -> (remains until stream is done) using ReadLockState = StateMachine; ReadLockState readState = ReadLockState::create(); bool disturbed = false; bool readPending = false; // Used by Sockets code to signal to the ReadableStream that it should error when read from // because the socket is currently being closed. bool isPendingClosure = false; friend class ReadableStream; friend class WritableStreamInternalController; friend class PipeLocked; }; class WritableStreamInternalController: public WritableStreamController { public: struct Writable { kj::Own sink; kj::Canceler canceler; Writable(kj::Own sink): sink(kj::mv(sink)) {} void abort(kj::Exception&& ex); }; explicit WritableStreamInternalController(StreamStates::Closed closed) : state(State::create()) {} explicit WritableStreamInternalController(StreamStates::Errored errored) : state(State::create(kj::mv(errored))) {} explicit WritableStreamInternalController(kj::Own writable, kj::Maybe> observer, kj::Maybe maybeHighWaterMark = kj::none, kj::Maybe> maybeClosureWaitable = kj::none) : state(State::create>( IoContext::current().addObject(kj::heap(kj::mv(writable))))), observer(kj::mv(observer)), maybeHighWaterMark(maybeHighWaterMark), maybeClosureWaitable(kj::mv(maybeClosureWaitable)) {} WritableStreamInternalController(WritableStreamInternalController&& other) = default; WritableStreamInternalController& operator=(WritableStreamInternalController&& other) = default; ~WritableStreamInternalController() noexcept(false) override; void setOwnerRef(WritableStream& stream) override { owner = stream; } jsg::Ref addRef() override; jsg::Promise write(jsg::Lock& js, jsg::Optional> value) override; jsg::Promise close(jsg::Lock& js, bool markAsHandled = false) override; jsg::Promise flush(jsg::Lock& js, bool markAsHandled = false) override; jsg::Promise abort(jsg::Lock& js, jsg::Optional> reason) override; kj::Maybe> tryPipeFrom( jsg::Lock& js, jsg::Ref source, PipeToOptions options) override; kj::Maybe> removeSink(jsg::Lock& js) override; void detach(jsg::Lock& js) override; kj::Maybe getDesiredSize() override; bool isLockedToWriter() const override { return !writeState.is(); } bool lockWriter(jsg::Lock& js, Writer& writer) override; void releaseWriter(Writer& writer, kj::Maybe maybeJs) override; // See the comment for releaseWriter in common.h for details on the use of maybeJs kj::Maybe> isErroring(jsg::Lock& js) override { // TODO(later): The internal controller has no concept of an "erroring" // state, so for now we just return kj::none here. return kj::none; } void visitForGc(jsg::GcVisitor& visitor) override; void setHighWaterMark(uint64_t highWaterMark); bool isClosedOrClosing() override; bool isPiping(); bool isErrored() override; inline bool isByteOriented() const override { return true; } void setPendingClosure() override { isPendingClosure = true; } kj::StringPtr jsgGetMemoryName() const override; size_t jsgGetMemorySelfSize() const override; void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override; private: struct AbortOptions { bool reject = false; bool handled = false; }; jsg::Promise doAbort(jsg::Lock& js, v8::Local reason, AbortOptions options = {.reject = false, .handled = false}); void doClose(jsg::Lock& js); void doError(jsg::Lock& js, v8::Local reason); void ensureWriting(jsg::Lock& js); jsg::Promise writeLoop(jsg::Lock& js, IoContext& ioContext); jsg::Promise writeLoopAfterFrontOutputLock(jsg::Lock& js); void drain(jsg::Lock& js, v8::Local reason); void finishClose(jsg::Lock& js); void finishError(jsg::Lock& js, v8::Local reason); jsg::Promise closeImpl(jsg::Lock& js, bool markAsHandled); struct PipeLocked { static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj; ReadableStream& ref; }; kj::Maybe owner; // State machine for WritableStreamInternalController: // Closed is terminal, Errored is implicitly terminal via ErrorState. // IoOwn is the active state (stream is writable). using State = StateMachine, ErrorState, ActiveState>, StreamStates::Closed, StreamStates::Errored, IoOwn>; State state; // Lock state machine for WritableStreamInternalController: // All states can transition to any other state (no terminal states). // Unlocked -> Locked (removeSink() or detach() called) // Unlocked -> WriterLocked (lockWriter() called) // Unlocked -> PipeLocked (tryPipeFrom() called) // WriterLocked -> Unlocked (releaseWriter() called) // WriterLocked -> Locked (doClose/doError called - stream closed but writer still attached) // PipeLocked -> Unlocked (pipe completes) using WriteLockState = StateMachine; WriteLockState writeState = WriteLockState::create(); kj::Maybe> observer; kj::Maybe> maybePendingAbort; uint64_t currentWriteBufferSize = 0; // The highWaterMark is the total amount of data currently buffered in // the controller waiting to be flushed out to the underlying WritableStreamSink. // It is used to implement backpressure signaling using desiredSize and the ready // promise on the writer. kj::Maybe maybeHighWaterMark; // Used by Sockets code to ensure the connection is established before the associated // WritableStream is closed. kj::Maybe> maybeClosureWaitable; bool waitingOnClosureWritableAlready = false; // Used by Sockets code to signal to the WritableStream that it should error when written to // because the socket is currently being closed. bool isPendingClosure = false; void adjustWriteBufferSize(jsg::Lock& js, int64_t amount); void updateBackpressure(jsg::Lock& js, bool backpressure); struct Write { kj::Maybe::Resolver> promise; size_t totalBytes; kj::Array ownBytes; kj::ArrayPtr bytes; JSG_MEMORY_INFO(Write) { tracker.trackField("resolver", promise); if (ownBytes != nullptr) { tracker.trackFieldWithSize("backing", totalBytes); } } }; struct Close { kj::Maybe::Resolver> promise; JSG_MEMORY_INFO(Close) { tracker.trackField("promise", promise); } }; struct Flush { kj::Maybe::Resolver> promise; JSG_MEMORY_INFO(Flush) { tracker.trackField("promise", promise); } }; struct Pipe { // PipeState is ref-counted so that it can be safely captured by lambdas in pipeLoop(). // When drain() destroys the Pipe, the state survives as long as pending callbacks need it. // The `aborted` flag is set when the Pipe is destroyed. struct State: public kj::Refcounted { WritableStreamInternalController& parent; ReadableStreamController::PipeController& source; kj::Maybe::Resolver> promise; kj::Maybe> maybeSignal; bool preventAbort; bool preventClose; bool preventCancel; // True when the Pipe is being destroyed bool aborted = false; State(WritableStreamInternalController& parent, ReadableStreamController::PipeController& source, kj::Maybe::Resolver> promise, bool preventAbort, bool preventClose, bool preventCancel, kj::Maybe> maybeSignal) : parent(parent), source(source), promise(kj::mv(promise)), maybeSignal(kj::mv(maybeSignal)), preventAbort(preventAbort), preventClose(preventClose), preventCancel(preventCancel) {} bool checkSignal(jsg::Lock& js); jsg::Promise pipeLoop(jsg::Lock& js); jsg::Promise write(v8::Local value); JSG_MEMORY_INFO(State) { tracker.trackField("resolver", promise); tracker.trackField("signal", maybeSignal); } }; kj::Own state; Pipe(WritableStreamInternalController& parent, ReadableStreamController::PipeController& source, kj::Maybe::Resolver> promise, bool preventAbort, bool preventClose, bool preventCancel, kj::Maybe> maybeSignal) : state(kj::refcounted(parent, source, kj::mv(promise), preventAbort, preventClose, preventCancel, kj::mv(maybeSignal))) {} ~Pipe() noexcept(false) { state->aborted = true; } WritableStreamInternalController& parent() { return state->parent; } ReadableStreamController::PipeController& source() { return state->source; } kj::Maybe::Resolver>& promise() { return state->promise; } bool preventAbort() const { return state->preventAbort; } bool preventClose() const { return state->preventClose; } bool preventCancel() const { return state->preventCancel; } kj::Maybe>& maybeSignal() { return state->maybeSignal; } bool checkSignal(jsg::Lock& js) { return state->checkSignal(js); } jsg::Promise pipeLoop(jsg::Lock& js) { return state->pipeLoop(js); } jsg::Promise write(v8::Local value) { return state->write(value); } JSG_MEMORY_INFO(Pipe) { tracker.trackField("state", state); } }; struct WriteEvent { kj::Maybe>> outputLock; // must wait for this before actually writing kj::OneOf, kj::Own, kj::Own, kj::Own> event; JSG_MEMORY_INFO(WriteEvent) { if (outputLock != kj::none) { tracker.trackFieldWithSize("outputLock", sizeof(IoOwn>)); } KJ_SWITCH_ONEOF(event) { KJ_CASE_ONEOF(w, kj::Own) { tracker.trackField("inner", w); } KJ_CASE_ONEOF(p, kj::Own) { tracker.trackField("inner", p); } KJ_CASE_ONEOF(c, kj::Own) { tracker.trackField("inner", c); } KJ_CASE_ONEOF(f, kj::Own) { tracker.trackField("inner", f); } } } }; RingBuffer queue; }; } // namespace workerd::api