// 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 #include "standard.h" #include "readable.h" #include "writable.h" #include #include #include #include #include #include #include namespace workerd::api { using DefaultController = jsg::Ref; using ByobController = jsg::Ref; namespace { struct ValueReadable; struct ByteReadable; } // namespace // ======================================================================================= // The Unlocked, Locked, ReaderLocked, and WriterLocked structs // are used to track the current lock status of JavaScript-backed streams. // All readable and writable streams begin in the Unlocked state. When a // reader or writer are attached, the streams will transition into the // ReaderLocked or WriterLocked state. When the reader is released, those // will transition back to Unlocked. // // When a readable is piped to a writable, both will enter the PipeLocked state. // (PipeLocked is defined within the ReadableLockImpl and WritableLockImpl classes // below) When the pipe completes, both will transition back to Unlocked. // // When a ReadableStreamJsController is tee()'d, it will enter the locked state. namespace { // A utility class used by ReadableStreamJsController // for implementing the reader lock in a consistent way (without duplicating any code). template class ReadableLockImpl { public: using PipeController = ReadableStreamController::PipeController; using Reader = ReadableStreamController::Reader; bool isLockedToReader() const { return !state.template is(); } bool lockReader(jsg::Lock& js, Controller& self, Reader& reader); // See the comment for releaseReader in common.h for details on the use of maybeJs void releaseReader(Controller& self, Reader& reader, kj::Maybe maybeJs); bool lock(); void onClose(jsg::Lock& js); void onError(jsg::Lock& js, v8::Local reason); kj::Maybe tryPipeLock(Controller& self); void visitForGc(jsg::GcVisitor& visitor); kj::StringPtr jsgGetMemoryName() const { return "ReadableLockImpl"_kjc; } size_t jsgGetMemorySelfSize() const { return sizeof(ReadableLockImpl); } void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(locked, Locked) {} KJ_CASE_ONEOF(unlocked, Unlocked) {} KJ_CASE_ONEOF(pipeLocked, PipeLocked) {} KJ_CASE_ONEOF(readerLocked, ReaderLocked) { tracker.trackField("readerLocked", readerLocked); } } } private: class PipeLocked final: public PipeController { public: static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj; explicit PipeLocked(Controller& inner): inner(inner) {} bool isClosed() override { return inner.state.template is(); } kj::Maybe> tryGetErrored(jsg::Lock& js) override { KJ_IF_SOME(errored, inner.state.template tryGetUnsafe()) { return errored.getHandle(js); } return kj::none; } void cancel(jsg::Lock& js, v8::Local reason) override { // Cancel here returns a Promise but we do not need to propagate it. // We can safely drop it on the floor here. auto promise KJ_UNUSED = inner.cancel(js, reason); } void close(jsg::Lock& js) override { inner.doClose(js); } void error(jsg::Lock& js, v8::Local reason) override { inner.doError(js, reason); } void release(jsg::Lock& js, kj::Maybe> maybeError = kj::none) override { KJ_IF_SOME(error, maybeError) { cancel(js, error); } inner.lock.state.template transitionTo(); } kj::Maybe> tryPumpTo(WritableStreamSink& sink, bool end) override; jsg::Promise read(jsg::Lock& js) override; private: Controller& inner; friend Controller; }; // State machine for ReadableLockImpl: // All states can transition to any other state (no terminal states). // Unlocked -> Locked (lock() called for tee) // Unlocked -> ReaderLocked (lockReader() called) // Unlocked -> PipeLocked (tryPipeLock() called) // ReaderLocked -> Unlocked (releaseReader() called) // PipeLocked -> Unlocked (release() or onClose/onError called) // Locked -> (remains until stream is done) using LockState = StateMachine; LockState state = LockState::template create(); friend Controller; }; // A utility class used by WritableStreamJsController to implement the writer lock // mechanism. Extracted for consistency with ReadableStreamJsController and to // eventually allow it to be shared also with WritableStreamInternalController. template class WritableLockImpl { public: using Writer = WritableStreamController::Writer; bool isLockedToWriter() const; bool lockWriter(jsg::Lock& js, Controller& self, Writer& writer); // See the comment for releaseWriter in common.h for details on the use of maybeJs void releaseWriter(Controller& self, Writer& writer, kj::Maybe maybeJs); void visitForGc(jsg::GcVisitor& visitor); bool pipeLock(WritableStream& owner, jsg::Ref source, PipeToOptions& options); void releasePipeLock(); JSG_MEMORY_INFO(WritableLockImpl) { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(unlocked, Unlocked) {} KJ_CASE_ONEOF(locked, Locked) {} KJ_CASE_ONEOF(writerLocked, WriterLocked) { tracker.trackField("writerLocked", writerLocked); } KJ_CASE_ONEOF(pipeLocked, PipeLocked) { tracker.trackField("pipeLocked", pipeLocked); } } } private: struct PipeLocked { static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj; ReadableStreamController::PipeController& source; jsg::Ref readableStreamRef; kj::Maybe> maybeSignal; kj::Maybe> checkSignal(jsg::Lock& js, Controller& self); struct Flags { uint8_t preventAbort : 1 = 0; uint8_t preventCancel : 1 = 0; uint8_t preventClose : 1 = 0; uint8_t pipeThrough : 1 = 0; }; Flags flags{}; JSG_MEMORY_INFO(PipeLocked) { tracker.trackField("readableStreamRef", readableStreamRef); tracker.trackField("signal", maybeSignal); } }; // State machine for WritableLockImpl: // All states can transition to any other state (no terminal states). // Unlocked -> Locked (not currently used) // Unlocked -> WriterLocked (lockWriter() called) // Unlocked -> PipeLocked (pipeLock() called) // WriterLocked -> Unlocked (releaseWriter() called) // PipeLocked -> Unlocked (releasePipeLock() called) using LockState = StateMachine; LockState state = LockState::template create(); inline kj::Maybe tryGetPipe() { KJ_IF_SOME(locked, state.template tryGetUnsafe()) { return locked; } return kj::none; } friend Controller; }; // ====================================================================================== template bool ReadableLockImpl::lock() { if (isLockedToReader()) { return false; } state.template transitionTo(); return true; } template bool ReadableLockImpl::lockReader(jsg::Lock& js, Controller& self, Reader& reader) { if (isLockedToReader()) { return false; } auto prp = js.newPromiseAndResolver(); prp.promise.markAsHandled(js); auto lock = ReaderLocked(reader, kj::mv(prp.resolver)); if (self.state.template is()) { maybeResolvePromise(js, lock.getClosedFulfiller()); } else KJ_IF_SOME(errored, self.state.template tryGetUnsafe()) { maybeRejectPromise(js, lock.getClosedFulfiller(), errored.getHandle(js)); } state.template transitionTo(kj::mv(lock)); reader.attach(self, kj::mv(prp.promise)); return true; } template void ReadableLockImpl::releaseReader( Controller& self, Reader& reader, kj::Maybe maybeJs) { KJ_IF_SOME(locked, state.template tryGetUnsafe()) { KJ_ASSERT(&locked.getReader() == &reader); KJ_IF_SOME(js, maybeJs) { auto reason = js.typeError("This ReadableStream reader has been released."_kj); KJ_SWITCH_ONEOF(self.state) { KJ_CASE_ONEOF(initial, typename Controller::Initial) {} KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(errored, StreamStates::Errored) {} KJ_CASE_ONEOF(consumer, kj::Own) { consumer->cancelPendingReads(js, reason); } KJ_CASE_ONEOF(consumer, kj::Own) { consumer->cancelPendingReads(js, reason); } } maybeRejectPromise(js, locked.getClosedFulfiller(), reason); } // Keep the locked.clear() after the isolate and hasPendingReadRequests check above. // Clearing will release the references and we don't want to do that if the // hasPendingReadRequests check fails. locked.clear(); // When maybeJs is nullptr, that means releaseReader was called when the reader is // being deconstructed and not as the result of explicitly calling releaseLock and // we do not have an isolate lock. In that case, we don't want to change the lock // state itself. Moving the lock above will free the lock state while keeping the // ReadableStream marked as locked. if (maybeJs != kj::none) { state.template transitionTo(); } } } template kj::Maybe ReadableLockImpl::tryPipeLock( Controller& self) { if (isLockedToReader()) { return kj::none; } return state.template transitionTo(self); } template void ReadableLockImpl::visitForGc(jsg::GcVisitor& visitor) { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(locked, Locked) {} KJ_CASE_ONEOF(locked, Unlocked) {} KJ_CASE_ONEOF(locked, PipeLocked) {} KJ_CASE_ONEOF(locked, ReaderLocked) { visitor.visit(locked); } } } template void ReadableLockImpl::onClose(jsg::Lock& js) { KJ_IF_SOME(locked, state.template tryGetUnsafe()) { try { maybeResolvePromise(js, locked.getClosedFulfiller()); } catch (jsg::JsExceptionThrown&) { // Resolving the promise could end up throwing an exception in some cases, // causing a jsg::JsExceptionThrown to be thrown. At this point, however, // we are already in the process of closing the stream and an error at this // point is not recoverable. Log and move on. LOG_NOSENTRY(ERROR, "Error resolving ReadableStream reader closed promise"); }; } else { (void)state.template transitionFromTo(); } } template void ReadableLockImpl::onError(jsg::Lock& js, v8::Local reason) { KJ_IF_SOME(locked, state.template tryGetUnsafe()) { try { maybeRejectPromise(js, locked.getClosedFulfiller(), reason); } catch (jsg::JsExceptionThrown&) { // Rejecting the promise could end up throwing an exception in some cases, // causing a jsg::JsExceptionThrown to be thrown. At this point, however, // we are already in the process of closing the stream and an error at this // point is not recoverable. Log and move on. LOG_NOSENTRY(ERROR, "Error rejecting ReadableStream reader closed promise"); } } else { (void)state.template transitionFromTo(); } } template kj::Maybe> ReadableLockImpl::PipeLocked::tryPumpTo( WritableStreamSink& sink, bool end) { // We return nullptr here because this controller does not support kj's pumpTo. return kj::none; } template jsg::Promise ReadableLockImpl::PipeLocked::read(jsg::Lock& js) { return KJ_ASSERT_NONNULL(inner.read(js, kj::none)); } // ====================================================================================== template bool WritableLockImpl::isLockedToWriter() const { return !state.template is(); } template bool WritableLockImpl::lockWriter(jsg::Lock& js, Controller& self, Writer& writer) { if (isLockedToWriter()) { return false; } auto closedPrp = js.newPromiseAndResolver(); closedPrp.promise.markAsHandled(js); auto readyPrp = js.newPromiseAndResolver(); readyPrp.promise.markAsHandled(js); auto lock = WriterLocked(writer, kj::mv(closedPrp.resolver), kj::mv(readyPrp.resolver)); if (self.state.template is()) { maybeResolvePromise(js, lock.getClosedFulfiller()); maybeResolvePromise(js, lock.getReadyFulfiller()); } else KJ_IF_SOME(errored, self.state.template tryGetUnsafe()) { maybeRejectPromise(js, lock.getClosedFulfiller(), errored.getHandle(js)); maybeRejectPromise(js, lock.getReadyFulfiller(), errored.getHandle(js)); } else { if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) { // Per spec (SetUpWritableStreamDefaultWriter step 4), the ready promise // is resolved when the stream is writable and not experiencing backpressure, // regardless of whether the start algorithm has completed. The backpressure // state is set synchronously during SetUpWritableStreamDefaultController. KJ_IF_SOME(erroring, self.isErroring(js)) { maybeRejectPromise(js, lock.getReadyFulfiller(), erroring); } else if (!self.hasBackpressure()) { maybeResolvePromise(js, lock.getReadyFulfiller()); } } else { if (self.isStarted()) { maybeResolvePromise(js, lock.getReadyFulfiller()); } } } state.template transitionTo(kj::mv(lock)); writer.attach(js, self, kj::mv(closedPrp.promise), kj::mv(readyPrp.promise)); return true; } template void WritableLockImpl::releaseWriter( Controller& self, Writer& writer, kj::Maybe maybeJs) { KJ_IF_SOME(locked, state.template tryGetUnsafe()) { KJ_ASSERT(&locked.getWriter() == &writer); KJ_IF_SOME(js, maybeJs) { KJ_SWITCH_ONEOF(self.state) { KJ_CASE_ONEOF(initial, typename Controller::Initial) {} KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(errored, StreamStates::Errored) {} KJ_CASE_ONEOF(controller, jsg::Ref) { controller->cancelPendingWrites( js, js.typeError("This WritableStream writer has been released."_kjc)); } } // Per spec (WritableStreamDefaultWriterRelease), both the ready and closed // promises must be rejected when the writer is released. auto releaseReason = js.v8TypeError("This WritableStream writer has been released."_kjc); if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) { if (locked.getReadyFulfiller() != kj::none) { maybeRejectPromise(js, locked.getReadyFulfiller(), releaseReason); } else { // The ready fulfiller was already consumed (promise was resolved). // Per spec (WritableStreamDefaultWriterEnsureReadyPromiseRejected), // we must replace it with a new rejected promise. auto prp = js.newPromiseAndResolver(); prp.promise.markAsHandled(js); prp.resolver.reject(js, releaseReason); locked.setReadyFulfiller(js, prp); } } else { maybeRejectPromise(js, locked.getReadyFulfiller(), releaseReason); } maybeRejectPromise(js, locked.getClosedFulfiller(), releaseReason); } locked.clear(); // When maybeJs is nullptr, that means releaseWriter was called when the writer is // being deconstructed and not as the result of explicitly calling releaseLock and // we do not have an isolate lock. In that case, we don't want to change the lock // state itself. Moving the lock above will free the lock state while keeping the // WritableStream marked as locked. if (maybeJs != kj::none) { state.template transitionTo(); } } } template bool WritableLockImpl::pipeLock( WritableStream& owner, jsg::Ref source, PipeToOptions& options) { if (isLockedToWriter()) { return false; } auto& sourceLock = KJ_ASSERT_NONNULL(source->getController().tryPipeLock()); state.template transitionTo(PipeLocked{ .source = sourceLock, .readableStreamRef = kj::mv(source), .maybeSignal = kj::mv(options.signal), .flags = { .preventAbort = options.preventAbort.orDefault(false), .preventCancel = options.preventCancel.orDefault(false), .preventClose = options.preventClose.orDefault(false), .pipeThrough = options.pipeThrough, }, }); return true; } template void WritableLockImpl::releasePipeLock() { if (state.template is()) { state.template transitionTo(); } } template void WritableLockImpl::visitForGc(jsg::GcVisitor& visitor) { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(locked, Unlocked) {} KJ_CASE_ONEOF(locked, Locked) {} KJ_CASE_ONEOF(locked, WriterLocked) { visitor.visit(locked); } KJ_CASE_ONEOF(locked, PipeLocked) { visitor.visit(locked.readableStreamRef); KJ_IF_SOME(signal, locked.maybeSignal) { visitor.visit(signal); } } } } template kj::Maybe> WritableLockImpl::PipeLocked::checkSignal( jsg::Lock& js, Controller& self) { KJ_IF_SOME(signal, maybeSignal) { if (signal->getAborted(js)) { auto reason = signal->getReason(js); if (!flags.preventCancel) { source.release(js, v8::Local(reason)); } else { source.release(js); } if (!flags.preventAbort) { return self.abort(js, reason).then(js, JSG_VISITABLE_LAMBDA((this, reason = reason.addRef(js), ref = self.addRef()), (reason, ref), (jsg::Lock& js) { return rejectedMaybeHandledPromise(js, reason.getHandle(js), flags.pipeThrough); })); } return rejectedMaybeHandledPromise(js, reason, flags.pipeThrough); } } return kj::none; } auto maybeAddFunctor(jsg::Lock& js, auto promise, auto onSuccess, auto onFailure) { KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { return promise.then( js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure))); } else { return promise.then(js, kj::mv(onSuccess), kj::mv(onFailure)); } } jsg::Promise maybeRunAlgorithm( jsg::Lock& js, auto& maybeAlgorithm, auto&& onSuccess, auto&& onFailure, auto&&... args) { // The algorithm is a JavaScript function mapped through jsg::Function. // It is expected to return a Promise mapped via jsg::Promise. If the // function returns synchronously, the jsg::Promise wrapper ensures // that it is properly mapped to a jsg::Promise, but if the Promise // throws synchronously, we have to convert that synchronous throw // into a proper rejected jsg::Promise. KJ_IF_SOME(algorithm, maybeAlgorithm) { // We need two layers of JSG_TRY here, unfortunately. The inner layer // covers the algorithm implementation itself and is our typical error // handling path. It ensures that if the algorithm throws an exception, // that is properly converted in to a rejected promise that is *then* // handled by the onFailure handler that is passed in. The outer JSG_TRY // handles the rare and generally unexpected failure of the calls to // .then() itself, which can throw JS exceptions synchronously in certain // rare cases. For those we return a rejected promise but do not call the // onFailure case since such errors are generally indicative of a fatal // condition in the isolate (e.g. out of memory, other fatal exception, etc). JSG_TRY(js) { KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { auto getInnerPromise = [&]() -> jsg::Promise { JSG_TRY(js) { return algorithm(js, kj::fwd(args)...); } JSG_CATCH(exception) { return js.rejectedPromise(kj::mv(exception)); } }; return getInnerPromise().then( js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure))); } else { auto getInnerPromise = [&]() -> jsg::Promise { JSG_TRY(js) { return algorithm(js, kj::fwd(args)...); } JSG_CATCH(exception) { return js.rejectedPromise(kj::mv(exception)); } }; return getInnerPromise().then(js, kj::mv(onSuccess), kj::mv(onFailure)); } } JSG_CATCH(exception) { return js.rejectedPromise(kj::mv(exception)); } } // If the algorithm does not exist, we just handle it as a success and move on. onSuccess(js); return js.resolvedPromise(); } jsg::Promise maybeRunAlgorithmAsync( jsg::Lock& js, auto& maybeAlgorithm, auto&& onSuccess, auto&& onFailure, auto&&... args) { // The algorithm is a JavaScript function mapped through jsg::Function. // It is expected to return a Promise mapped via jsg::Promise. If the // function returns synchronously, the jsg::Promise wrapper ensures // that it is properly mapped to a jsg::Promise, but if the Promise // throws synchronously, we have to convert that synchronous throw // into a proper rejected jsg::Promise. KJ_IF_SOME(algorithm, maybeAlgorithm) { // We need two layers of tryCatch here, unfortunately. The inner layer // covers the algorithm implementation itself and is our typical error // handling path. It ensures that if the algorithm throws an exception, // that is properly converted in to a rejected promise that is *then* // handled by the onFailure handler that is passed in. The outer tryCatch // handles the rare and generally unexpected failure of the calls to // .then() itself, which can throw JS exceptions synchronously in certain // rare cases. For those we return a rejected promise but do not call the // onFailure case since such errors are generally indicative of a fatal // condition in the isolate (e.g. out of memory, other fatal exception, etc). return js.tryCatch([&] { KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { return js .tryCatch([&] { return algorithm(js, kj::fwd(args)...); }, [&](jsg::Value&& exception) { return js.rejectedPromise(kj::mv(exception)); }) .then(js, ioContext.addFunctor(kj::mv(onSuccess)), ioContext.addFunctor(kj::mv(onFailure))); } else { return js .tryCatch([&] { return algorithm(js, kj::fwd(args)...); }, [&](jsg::Value&& exception) { return js.rejectedPromise(kj::mv(exception)); }).then(js, kj::mv(onSuccess), kj::mv(onFailure)); } }, [&](jsg::Value&& exception) { return js.rejectedPromise(kj::mv(exception)); }); } // If the algorithm does not exist, we handle it as a success but ensure // it runs asynchronously by scheduling via a resolved promise. KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { return js.resolvedPromise().then(js, ioContext.addFunctor(kj::mv(onSuccess))); } else { return js.resolvedPromise().then(js, kj::mv(onSuccess)); } } int getHighWaterMark( const UnderlyingSource& underlyingSource, const StreamQueuingStrategy& queuingStrategy) { bool isBytes = underlyingSource.type.map([](auto& s) { return s == "bytes"; }).orDefault(false); return queuingStrategy.highWaterMark.orDefault(isBytes ? 0 : 1); } } // namespace // It is possible for the controller state to be released synchronously while // we are in the middle of a read. When that happens we need to defer the actual // close/error state change until the read call is complete. deferControllerStateChange // handles this for us by using the state machine's operation tracking to defer // pending close/error transitions until the read is complete. template jsg::Promise deferControllerStateChange(jsg::Lock& js, Controller& controller, kj::FunctionParam()> readCallback) { bool endOperation = true; // The readCallback and the controller.doClose(..) and controller.doError(...) // methods, as well as the methods can trigger JavaScript errors to be thrown // synchronously in some cases. We want to make sure non-fatal errors cause the // stream to error and only fatal cases bubble up. return js.tryCatch([&] { controller.state.beginOperation(); auto result = readCallback(); endOperation = false; // endOperation() will automatically apply any pending state if this was the last operation. // Returns true if a pending state was applied. if (controller.state.endOperation()) { // A pending state was applied. Call the appropriate callback. // Skip callbacks if execution is being terminated (e.g., CPU time limit) since we can't // safely execute JavaScript in that state. if (!js.v8Isolate->IsExecutionTerminating()) { if (controller.state.template is()) { controller.lock.onClose(js); } else if (controller.state.template is()) { KJ_IF_SOME(err, controller.state.template tryGetUnsafe()) { controller.lock.onError(js, err.getHandle(js)); } } } } return kj::mv(result); }, [&](jsg::Value exception) -> jsg::Promise { if (endOperation) { // Clear any pending state since we're erroring controller.state.clearPendingState(); (void)controller.state.endOperation(); } controller.doError(js, exception.getHandle(js)); return js.rejectedPromise(kj::mv(exception)); }); } // The ReadableStreamJsController provides the implementation of custom // ReadableStreams backed by a user-code provided Underlying Source. The implementation // is fairly complicated and defined entirely by the streams specification. // // Another important thing to understand is that there are two types of JavaScript // backed ReadableStreams: value-oriented, and byte-oriented. // // When user code uses the `new ReadableStream(underlyingSource)` constructor, the // underlyingSource argument may have a `type` property, the value of which is either // `undefined`, the empty string, or the string value `'bytes'`. If the underlyingSource // argument is not given, the default value of `type` is `undefined`. If `type` is // `undefined` or the empty string, the ReadableStream is value-oriented. If `type` is // exactly equal to `'bytes'`, the ReadableStream is byte-oriented. // // For value-oriented streams, any JavaScript value can be pushed through the stream, // and the stream will only support use of the ReadableStreamDefaultReader to consume // the stream data. // // For byte-oriented streams, only byte data (as provided by `ArrayBufferView`s) can // be pushed through the stream. All byte-oriented streams support using both // ReadableStreamDefaultReader and ReadableStreamBYOBReader to consume the stream // data. // // When the ReadableStreamJsController::setup() method is called the type // of stream is determined, and the controller will create an instance of either // jsg::Ref or jsg::Ref. // These are the objects that are actually passed on to the user-code's Underlying Source // implementation. class ReadableStreamJsController final: public ReadableStreamController { public: using ReadableLockImpl = ReadableLockImpl; KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamJsController); explicit ReadableStreamJsController(); explicit ReadableStreamJsController(StreamStates::Closed closed); explicit ReadableStreamJsController(StreamStates::Errored errored); explicit ReadableStreamJsController(jsg::Lock& js, ValueReadable& consumer); explicit ReadableStreamJsController(jsg::Lock& js, ByteReadable& consumer); jsg::Ref addRef() override; void setup(jsg::Lock& js, jsg::Optional maybeUnderlyingSource, jsg::Optional maybeQueuingStrategy) override; // Signals that this ReadableStream is no longer interested in the underlying // data source. Whether this cancels the underlying data source also depends // on whether or not there are other ReadableStreams still attached to it. // This operation is terminal. Once called, even while the returned Promise // is still pending, the ReadableStream will be no longer usable and any // data still in the queue will be dropped. Pending read requests will be // rejected if a reason is given, or resolved with no data otherwise. jsg::Promise cancel(jsg::Lock& js, jsg::Optional> reason) override; void doClose(jsg::Lock& js); void doError(jsg::Lock& js, v8::Local reason); bool canCloseOrEnqueue(); bool hasBackpressure(); bool isByteOriented() const override; bool isDisturbed() override; bool isClosedOrErrored() const override; bool isClosed() const override; bool isLockedToReader() const override; bool lockReader(jsg::Lock& js, Reader& reader) override; kj::Maybe> isErrored(jsg::Lock& js); kj::Maybe getDesiredSize(); jsg::Promise pipeTo( jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) override; kj::Promise> pumpTo( jsg::Lock& js, kj::Own, bool end) override; kj::Maybe> read( jsg::Lock& js, kj::Maybe byobOptions) override; kj::Maybe> drainingRead( jsg::Lock& js, size_t maxRead = kj::maxValue) override; // See the comment for releaseReader in common.h for details on the use of maybeJs void releaseReader(Reader& reader, kj::Maybe maybeJs) override; void setOwnerRef(ReadableStream& stream) override; Tee tee(jsg::Lock& js) override; kj::Maybe tryPipeLock() override; void visitForGc(jsg::GcVisitor& visitor) override; kj::Maybe> getController(); 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::Own detach(jsg::Lock& js, bool ignoreDisturbed) override; void setPendingClosure() override { KJ_UNIMPLEMENTED("only implemented for WritableStreamInternalController"); } kj::StringPtr jsgGetMemoryName() const override; size_t jsgGetMemorySelfSize() const override; void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const override; private: // If the stream was created within the scope of a request, we want to treat it as I/O // and make sure it is not advanced from the scope of a different request. kj::Maybe ioContext; kj::Maybe owner; // Initial state before setup() is called. struct Initial { static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; }; // State machine for ReadableStreamJsController: // Initial is the default state before setup() is called // ValueReadable and ByteReadable are the active states (stream has data) // Closed and Errored are terminal states (stream is done) // Initial -> ValueReadable or ByteReadable (setup() called) // Initial -> Closed (constructed with Closed) // Initial -> Errored (constructed with Errored) // ValueReadable -> Closed (doClose() or cancel() called) // ValueReadable -> Errored (doError() called) // ByteReadable -> Closed (doClose() or cancel() called) // ByteReadable -> Errored (doError() called) // Note: No single ActiveState since there are two active variants. // PendingStates allows Closed/Errored transitions to be deferred during reads. using State = StateMachine, ErrorState, PendingStates, Initial, StreamStates::Closed, StreamStates::Errored, kj::Own, kj::Own>; State state = State::create(); kj::Maybe expectedLength = kj::none; bool canceling = false; // The lock state is separate because a closed or errored stream can still be locked. ReadableLockImpl lock; bool disturbed = false; template jsg::Promise readAll(jsg::Lock& js, uint64_t limit); friend ReadableLockImpl; friend ReadableLockImpl::PipeLocked; friend struct ValueReadable; friend struct ByteReadable; template friend jsg::Promise deferControllerStateChange(jsg::Lock& js, Controller& controller, kj::FunctionParam()> readCallback); }; // The WritableStreamJsController provides the implementation of custom // WritableStream's backed by a user-code provided Underlying Sink. The implementation // is fairly complicated and defined entirely by the streams specification. class WritableStreamJsController final: public WritableStreamController { public: using WritableLockImpl = WritableLockImpl; using Controller = jsg::Ref; explicit WritableStreamJsController(); explicit WritableStreamJsController(StreamStates::Closed closed); explicit WritableStreamJsController(StreamStates::Errored errored); ~WritableStreamJsController() noexcept(false); KJ_DISALLOW_COPY_AND_MOVE(WritableStreamJsController); jsg::Promise abort(jsg::Lock& js, jsg::Optional> reason) override; jsg::Ref addRef() override; jsg::Promise close(jsg::Lock& js, bool markAsHandled = false) override; jsg::Promise flush(jsg::Lock& js, bool markAsHandled = false) override { KJ_UNIMPLEMENTED("expected WritableStreamInternalController implementation to be enough"); } void doClose(jsg::Lock& js); void doError(jsg::Lock& js, v8::Local reason); // Error through the underlying controller if available, going through the proper // error transition (Erroring -> Errored). void errorIfNeeded(jsg::Lock& js, v8::Local reason); kj::Maybe getDesiredSize() override; kj::Maybe> isErroring(jsg::Lock& js) override; kj::Maybe> isErroredOrErroring(jsg::Lock& js); bool isLocked() const; bool isLockedToWriter() const override; bool isStarted(); bool hasBackpressure(); inline bool isWritable() const { return state.isActive(); } bool lockWriter(jsg::Lock& js, Writer& writer) override; void maybeRejectReadyPromise(jsg::Lock& js, v8::Local reason); void maybeResolveReadyPromise(jsg::Lock& js); // See the comment for releaseWriter in common.h for details on the use of maybeJs void releaseWriter(Writer& writer, kj::Maybe maybeJs) override; kj::Maybe> removeSink(jsg::Lock& js) override; void detach(jsg::Lock& js) override; void setOwnerRef(WritableStream& stream) override; void setup(jsg::Lock& js, jsg::Optional maybeUnderlyingSink, jsg::Optional maybeQueuingStrategy) override; kj::Maybe> tryPipeFrom( jsg::Lock& js, jsg::Ref source, PipeToOptions options) override; void updateBackpressure(jsg::Lock& js, bool backpressure); jsg::Promise write(jsg::Lock& js, jsg::Optional> value) override; void visitForGc(jsg::GcVisitor& visitor) override; bool isClosedOrClosing() override; bool isErrored() override; inline bool isByteOriented() const override { return false; } void setPendingClosure() override { KJ_UNIMPLEMENTED("only implemented for WritableStreamInternalController"); } kj::StringPtr jsgGetMemoryName() const override; size_t jsgGetMemorySelfSize() const override; void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override; private: jsg::Promise pipeLoop(jsg::Lock& js); kj::Maybe ioContext; kj::Maybe owner; // Initial state before setup() is called. struct Initial { static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; }; // State machine for WritableStreamJsController: // Initial is the default state before setup() is called // Controller is the active state (stream is writable) // Closed is terminal, Errored is implicitly terminal via ErrorState using State = StateMachine, ErrorState, ActiveState, Initial, StreamStates::Closed, StreamStates::Errored, Controller>; State state = State::create(); WritableLockImpl lock; kj::Maybe> maybeAbortPromise; friend WritableLockImpl; }; kj::Own newReadableStreamJsController() { return kj::heap(); } kj::Own newWritableStreamJsController() { return kj::heap(); } template ReadableImpl::ReadableImpl( UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy) : state(State::template create(getHighWaterMark(underlyingSource, queuingStrategy))), algorithms(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {} template void ReadableImpl::start(jsg::Lock& js, jsg::Ref self) { KJ_ASSERT(!flags.started && !flags.starting); flags.starting = true; // Per the streams spec, the size function should be called with `undefined` as `this`, // not as a method on the strategy object. KJ_IF_SOME(sizeFunc, algorithms.size) { sizeFunc.setReceiver(jsg::Value(js.v8Isolate, js.v8Undefined())); } auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { flags.started = true; flags.starting = false; pullIfNeeded(js, kj::mv(self)); }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { flags.started = true; flags.starting = false; doError(js, kj::mv(reason)); }); maybeRunAlgorithm(js, algorithms.start, kj::mv(onSuccess), kj::mv(onFailure), kj::mv(self)); algorithms.start = kj::none; } template size_t ReadableImpl::consumerCount() { return state.whenActiveOr([](Queue& q) { return q.getConsumerCount(); }, size_t{0}); } template jsg::Promise ReadableImpl::cancel( jsg::Lock& js, jsg::Ref self, v8::Local reason) { if (state.template is()) { // We are already closed. There's nothing to cancel. // This shouldn't happen but we handle the case anyway, just to be safe. return js.resolvedPromise(); } KJ_IF_SOME(errored, state.template tryGetUnsafe()) { // We are already errored. There's nothing to cancel. // This shouldn't happen but we handle the case anyway, just to be safe. return js.rejectedPromise(errored.getHandle(js)); } auto& queue = state.template getUnsafe(); size_t consumerCount = queue.getConsumerCount(); if (consumerCount > 1) { // If there is more than 1 consumer, then we just return here with an // immediately resolved promise. The consumer will remove itself, // canceling its interest in the underlying source but we do not yet // want to cancel the underlying source since there are still other // consumers that want data. return js.resolvedPromise(); } // Otherwise, there should be exactly one consumer at this point. KJ_ASSERT(consumerCount == 1); KJ_IF_SOME(pendingCancel, maybePendingCancel) { // If we're already waiting for cancel to complete, just return the // already existing pending promise. // This shouldn't happen but we handle the case anyway, just to be safe. return pendingCancel.promise.whenResolved(js); } auto prp = js.newPromiseAndResolver(); maybePendingCancel = PendingCancel{ .fulfiller = kj::mv(prp.resolver), .promise = kj::mv(prp.promise), }; auto promise = KJ_ASSERT_NONNULL(maybePendingCancel).promise.whenResolved(js); doCancel(js, kj::mv(self), reason); return kj::mv(promise); } template bool ReadableImpl::canCloseOrEnqueue() { return state.isActive(); } // doCancel() is triggered by cancel() being called, which is an explicit signal from // the ReadableStream that we don't care about the data this controller provides any // more. We don't need to notify the consumers because we presume they already know // that they called cancel. What we do want to do here, tho, is close the implementation // and trigger the cancel algorithm. template void ReadableImpl::doCancel(jsg::Lock& js, jsg::Ref self, v8::Local reason) { state.template transitionTo(); auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { doClose(js); KJ_IF_SOME(pendingCancel, maybePendingCancel) { maybeResolvePromise(js, pendingCancel.fulfiller); } else { // Else block to avert dangling else compiler warning. } }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { // We do not call doError() here because there's really no point. Everything // that cares about the state of this controller impl has signaled that it // no longer cares and has gone away. doClose(js); KJ_IF_SOME(pendingCancel, maybePendingCancel) { maybeRejectPromise(js, pendingCancel.fulfiller, reason.getHandle(js)); } else { // Else block to avert dangling else compiler warning. } }); maybeRunAlgorithm(js, algorithms.cancel, kj::mv(onSuccess), kj::mv(onFailure), reason); } template void ReadableImpl::enqueue(jsg::Lock& js, kj::Rc entry, jsg::Ref self) { JSG_REQUIRE(canCloseOrEnqueue(), TypeError, "This ReadableStream is closed."); KJ_DEFER(pullIfNeeded(js, kj::mv(self))); auto& queue = state.template getUnsafe(); queue.push(js, kj::mv(entry)); } template void ReadableImpl::close(jsg::Lock& js) { JSG_REQUIRE(canCloseOrEnqueue(), TypeError, "This ReadableStream is closed."); auto& queue = state.template getUnsafe(); if (queue.hasPartiallyFulfilledRead()) { auto error = js.v8Ref(js.v8TypeError("This ReadableStream was closed with a partial read pending.")); doError(js, error.addRef(js)); js.throwException(kj::mv(error)); return; } queue.close(js); state.template transitionTo(); doClose(js); } template void ReadableImpl::doClose(jsg::Lock& js) { // The state should have already been set to closed. KJ_ASSERT(state.template is()); algorithms.clear(); } template void ReadableImpl::doError(jsg::Lock& js, jsg::Value reason) { // If already closed or errored, do nothing if (state.isInactive()) { return; } auto& queue = state.template getUnsafe(); queue.error(js, reason.addRef(js)); state.template transitionTo(kj::mv(reason)); algorithms.clear(); } template kj::Maybe ReadableImpl::getDesiredSize() { if (state.template is()) { return 0; } if (state.template is()) { return kj::none; } return state.template getUnsafe().desiredSize(); } // We should call pull if any of the consumers known to the queue have read requests or // we haven't yet signalled backpressure. template bool ReadableImpl::shouldCallPull() { return state.whenActiveOr( [this](Queue& q) { return q.wantsRead() || getDesiredSize().orDefault(0) > 0; }, false); } template void ReadableImpl::pullIfNeeded(jsg::Lock& js, jsg::Ref self) { // Determining if we need to pull is fairly complicated. All of the following // must hold true: if (!shouldCallPull()) { return; } if (flags.pulling) { flags.pullAgain = true; return; } KJ_ASSERT(!flags.pullAgain); flags.pulling = true; auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { flags.pulling = false; if (flags.pullAgain) { flags.pullAgain = false; pullIfNeeded(js, kj::mv(self)); } }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { flags.pulling = false; doError(js, kj::mv(reason)); }); maybeRunAlgorithm(js, algorithms.pull, kj::mv(onSuccess), kj::mv(onFailure), self.addRef()); } template void ReadableImpl::forcePullIfNeeded(jsg::Lock& js, jsg::Ref self) { // Like pullIfNeeded but bypasses the shouldCallPull() check. Used for draining reads // which need to pull all available data regardless of backpressure settings. if (!canCloseOrEnqueue()) { return; } if (flags.pulling) { flags.pullAgain = true; return; } KJ_ASSERT(!flags.pullAgain); flags.pulling = true; auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { flags.pulling = false; if (flags.pullAgain) { flags.pullAgain = false; // After a force pull, we go back to normal pullIfNeeded behavior. pullIfNeeded(js, kj::mv(self)); } }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { flags.pulling = false; doError(js, kj::mv(reason)); }); maybeRunAlgorithm(js, algorithms.pull, kj::mv(onSuccess), kj::mv(onFailure), self.addRef()); } template void ReadableImpl::visitForGc(jsg::GcVisitor& visitor) { state.visitForGc(visitor); KJ_IF_SOME(pendingCancel, maybePendingCancel) { visitor.visit(pendingCancel.fulfiller, pendingCancel.promise); } visitor.visit(algorithms); } template kj::Own::Consumer> ReadableImpl::getConsumer( kj::Maybe::StateListener&> listener) { auto& queue = state.template getUnsafe(); return kj::heap::Consumer>(queue, listener); } // ====================================================================================== template WritableImpl::WritableImpl( jsg::Lock& js, WritableStream& owner, jsg::Ref abortSignal) : owner(owner.addWeakRef()), signal(kj::mv(abortSignal)) { flags.pedanticWpt = FeatureFlags::get(js).getPedanticWpt(); } template jsg::Promise WritableImpl::abort( jsg::Lock& js, jsg::Ref self, v8::Local reason) { // Per the spec, the signal.reason should be a DOMException with name 'AbortError' // when no reason is provided, but the stored error should remain as the original reason. auto signalReason = [&]() -> jsg::JsValue { if (reason->IsUndefined() && FeatureFlags::get(js).getPedanticWpt()) { auto ex = js.domException( kj::str("AbortError"), kj::str("This writable stream has been aborted."), kj::none); return jsg::JsValue(KJ_ASSERT_NONNULL(ex.tryGetHandle(js))); } return jsg::JsValue(reason); }(); signal->triggerAbort(js, signalReason); // We have to check this again after the AbortSignal is triggered. if (state.isTerminal()) { return js.resolvedPromise(); } KJ_IF_SOME(pendingAbort, maybePendingAbort) { // Notice here that, per the spec, the reason given in this call of abort is // intentionally ignored if there is already an abort pending. return pendingAbort->whenResolved(js); } bool wasAlreadyErroring = false; if (state.template is()) { wasAlreadyErroring = true; reason = js.v8Undefined(); } KJ_DEFER(if (!wasAlreadyErroring) { startErroring(js, kj::mv(self), reason); }); maybePendingAbort = kj::heap(js, reason, wasAlreadyErroring); return KJ_ASSERT_NONNULL(maybePendingAbort)->whenResolved(js); } template kj::Maybe WritableImpl::tryGetOwner() { KJ_IF_SOME(o, owner) { return o->tryGet().map([](WritableStream& owner) -> WritableStreamJsController& { return static_cast(owner.getController()); }); } return kj::none; } template ssize_t WritableImpl::getDesiredSize() { return highWaterMark - amountBuffered; } template void WritableImpl::advanceQueueIfNeeded(jsg::Lock& js, jsg::Ref self) { if (!flags.started || inFlightWrite != kj::none) { return; } KJ_ASSERT(isWritable() || state.template is()); if (state.template is()) { return finishErroring(js, kj::mv(self)); } if (writeRequests.empty()) { if (closeRequest != kj::none) { KJ_ASSERT(inFlightClose == kj::none); KJ_ASSERT_NONNULL(closeRequest); inFlightClose = kj::mv(closeRequest); auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { finishInFlightClose(js, kj::mv(self)); }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { finishInFlightClose(js, kj::mv(self), reason.getHandle(js)); }); // Per the spec, the close algorithm should always run asynchronously, even if // there's no user-provided close handler. This ensures that releaseLock() can // reject the closed promise before the close completes. // The original maybeRunAlgorithm would call the onSuccess continuation // synchronously if algorithms.close is not specified. maybeRunAlgorithmAsync // always defers to a microtask. if (FeatureFlags::get(js).getPedanticWpt()) { maybeRunAlgorithmAsync(js, algorithms.close, kj::mv(onSuccess), kj::mv(onFailure)); } else { maybeRunAlgorithm(js, algorithms.close, kj::mv(onSuccess), kj::mv(onFailure)); } } return; } KJ_ASSERT(inFlightWrite == kj::none); auto req = dequeueWriteRequest(); auto value = req.value.addRef(js); auto size = req.size; inFlightWrite = kj::mv(req); auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef(), size), (self), (jsg::Lock& js) { amountBuffered -= size; finishInFlightWrite(js, self.addRef()); KJ_ASSERT(isWritable() || state.template is()); if (!isCloseQueuedOrInFlight() && isWritable()) { updateBackpressure(js); } if (state.template is() || writeRequests.empty()) { // In this case, we know advanceQueueIfNeeded won't recurse further, so we can // avoid the extra microtask hop. advanceQueueIfNeeded(js, kj::mv(self)); return js.resolvedPromise(); } // Here, however, let's avoid potentially deep recursion by hopping to a new // microtask to continue processing the queue. return js.resolvedPromise().then( js, JSG_VISITABLE_LAMBDA((this, self = kj::mv(self)), (self), (jsg::Lock & js) mutable { if (isWritable() || state.template is()) { advanceQueueIfNeeded(js, kj::mv(self)); } })); }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef(), size), (self), (jsg::Lock& js, jsg::Value reason) { amountBuffered -= size; finishInFlightWrite(js, kj::mv(self), reason.getHandle(js)); return js.resolvedPromise(); }); // Per the spec, the write algorithm should always run asynchronously, even if // there's no user-provided write handler. This ensures that backpressure changes // from the write don't resolve the ready promise synchronously, preserving correct // microtask ordering (e.g., ready rejects before closed on releaseLock). if (FeatureFlags::get(js).getPedanticWpt()) { maybeRunAlgorithmAsync(js, algorithms.write, kj::mv(onSuccess), kj::mv(onFailure), value.getHandle(js), self.addRef()); } else { maybeRunAlgorithm(js, algorithms.write, kj::mv(onSuccess), kj::mv(onFailure), value.getHandle(js), self.addRef()); } } template jsg::Promise WritableImpl::close(jsg::Lock& js, jsg::Ref self) { if (state.template is()) { return js.rejectedPromise(js.v8TypeError("This WritableStream has been closed."_kj)); } KJ_IF_SOME(errored, state.template tryGetUnsafe()) { return js.rejectedPromise(errored.addRef(js)); } KJ_ASSERT(isWritable() || state.template is()); JSG_REQUIRE( !isCloseQueuedOrInFlight(), TypeError, "Cannot close a writer that is already being closed"); auto prp = js.newPromiseAndResolver(); closeRequest = kj::mv(prp.resolver); if (flags.backpressure && isWritable()) { KJ_IF_SOME(owner, tryGetOwner()) { owner.maybeResolveReadyPromise(js); } } advanceQueueIfNeeded(js, kj::mv(self)); return kj::mv(prp.promise); } template void WritableImpl::dealWithRejection( jsg::Lock& js, jsg::Ref self, v8::Local reason) { if (isWritable()) { return startErroring(js, kj::mv(self), reason); } KJ_ASSERT(state.template is()); finishErroring(js, kj::mv(self)); } template WritableImpl::WriteRequest WritableImpl::dequeueWriteRequest() { auto write = kj::mv(writeRequests.front()); writeRequests.pop_front(); return kj::mv(write); } template void WritableImpl::doClose(jsg::Lock& js) { KJ_ASSERT(closeRequest == kj::none); KJ_ASSERT(inFlightClose == kj::none); KJ_ASSERT(inFlightWrite == kj::none); KJ_ASSERT(maybePendingAbort == kj::none); KJ_ASSERT(writeRequests.empty()); // State should have already been transitioned to Closed KJ_ASSERT(state.template is()); algorithms.clear(); KJ_IF_SOME(owner, tryGetOwner()) { owner.doClose(js); } } template void WritableImpl::doError(jsg::Lock& js, v8::Local reason) { KJ_ASSERT(closeRequest == kj::none); KJ_ASSERT(inFlightClose == kj::none); KJ_ASSERT(inFlightWrite == kj::none); KJ_ASSERT(maybePendingAbort == kj::none); KJ_ASSERT(writeRequests.empty()); // State should have already been transitioned to Errored KJ_ASSERT(state.template is()); algorithms.clear(); KJ_IF_SOME(owner, tryGetOwner()) { owner.doError(js, reason); } } template void WritableImpl::error(jsg::Lock& js, jsg::Ref self, v8::Local reason) { if (isWritable()) { algorithms.clear(); startErroring(js, kj::mv(self), reason); } } template void WritableImpl::finishErroring(jsg::Lock& js, jsg::Ref self) { auto erroring = kj::mv(KJ_ASSERT_NONNULL(state.template tryGetUnsafe())); auto reason = erroring.reason.getHandle(js); KJ_ASSERT(inFlightWrite == kj::none); KJ_ASSERT(inFlightClose == kj::none); state.template transitionTo(kj::mv(erroring.reason)); while (!writeRequests.empty()) { dequeueWriteRequest().resolver.reject(js, reason); } KJ_ASSERT(writeRequests.empty()); KJ_IF_SOME(pendingAbort, maybePendingAbort) { if (pendingAbort->reject) { pendingAbort->fail(js, reason); return rejectCloseAndClosedPromiseIfNeeded(js); } auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { auto& pendingAbort = KJ_ASSERT_NONNULL(maybePendingAbort); pendingAbort->reject = false; pendingAbort->complete(js); rejectCloseAndClosedPromiseIfNeeded(js); }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { auto& pendingAbort = KJ_ASSERT_NONNULL(maybePendingAbort); pendingAbort->fail(js, reason.getHandle(js)); rejectCloseAndClosedPromiseIfNeeded(js); }); maybeRunAlgorithm(js, algorithms.abort, kj::mv(onSuccess), kj::mv(onFailure), reason); return; } rejectCloseAndClosedPromiseIfNeeded(js); } template void WritableImpl::finishInFlightClose( jsg::Lock& js, jsg::Ref self, kj::Maybe> maybeReason) { algorithms.clear(); KJ_ASSERT_NONNULL(inFlightClose); KJ_ASSERT(isWritable() || state.template is()); KJ_IF_SOME(reason, maybeReason) { maybeRejectPromise(js, inFlightClose, reason); KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) { pendingAbort->fail(js, reason); } return dealWithRejection(js, kj::mv(self), reason); } maybeResolvePromise(js, inFlightClose); if (state.template is()) { KJ_IF_SOME(pendingAbort, PendingAbort::dequeue(maybePendingAbort)) { pendingAbort->reject = false; pendingAbort->complete(js); } } KJ_ASSERT(maybePendingAbort == kj::none); state.template transitionTo(); doClose(js); } template void WritableImpl::finishInFlightWrite( jsg::Lock& js, jsg::Ref self, kj::Maybe> maybeReason) { auto& write = KJ_ASSERT_NONNULL(inFlightWrite); KJ_IF_SOME(reason, maybeReason) { write.resolver.reject(js, reason); inFlightWrite = kj::none; KJ_ASSERT(isWritable() || state.template is()); return dealWithRejection(js, kj::mv(self), reason); } write.resolver.resolve(js); inFlightWrite = kj::none; } template bool WritableImpl::isCloseQueuedOrInFlight() { return closeRequest != kj::none || inFlightClose != kj::none; } template void WritableImpl::rejectCloseAndClosedPromiseIfNeeded(jsg::Lock& js) { algorithms.clear(); auto reason = KJ_ASSERT_NONNULL(state.template tryGetUnsafe()).getHandle(js); maybeRejectPromise(js, closeRequest, reason); PendingAbort::dequeue(maybePendingAbort); doError(js, reason); } template void WritableImpl::setup(jsg::Lock& js, jsg::Ref self, UnderlyingSink underlyingSink, StreamQueuingStrategy queuingStrategy) { KJ_ASSERT(!flags.started && !flags.starting); flags.starting = true; highWaterMark = queuingStrategy.highWaterMark.orDefault(1); auto startAlgorithm = kj::mv(underlyingSink.start); algorithms.write = kj::mv(underlyingSink.write); algorithms.close = kj::mv(underlyingSink.close); algorithms.abort = kj::mv(underlyingSink.abort); algorithms.size = kj::mv(queuingStrategy.size); // Per the streams spec, the size function should be called with `undefined` as `this`, // not as a method on the strategy object. KJ_IF_SOME(sizeFunc, algorithms.size) { sizeFunc.setReceiver(jsg::Value(js.v8Isolate, js.v8Undefined())); } auto onSuccess = JSG_VISITABLE_LAMBDA((this, self = self.addRef()), (self), (jsg::Lock& js) { KJ_ASSERT(isWritable() || state.template is()); if (isWritable()) { // Only resolve the ready promise if an abort is not pending. // It will have been rejected already. KJ_IF_SOME(owner, tryGetOwner()) { owner.maybeResolveReadyPromise(js); } else { // Else block to avert dangling else compiler warning. } } flags.started = true; flags.starting = false; advanceQueueIfNeeded(js, kj::mv(self)); }); auto onFailure = JSG_VISITABLE_LAMBDA( (this, self = self.addRef()), (self), (jsg::Lock& js, jsg::Value reason) { auto handle = reason.getHandle(js); KJ_ASSERT(isWritable() || state.template is()); KJ_IF_SOME(owner, tryGetOwner()) { owner.maybeRejectReadyPromise(js, handle); } else { // Else block to avert dangling else compiler warning. } flags.started = true; flags.starting = false; dealWithRejection(js, kj::mv(self), handle); }); flags.backpressure = getDesiredSize() <= 0; maybeRunAlgorithm(js, startAlgorithm, kj::mv(onSuccess), kj::mv(onFailure), self.addRef()); } template void WritableImpl::startErroring( jsg::Lock& js, jsg::Ref self, v8::Local reason) { KJ_ASSERT(isWritable()); KJ_IF_SOME(owner, tryGetOwner()) { owner.maybeRejectReadyPromise(js, reason); } state.template transitionTo(js.v8Ref(reason)); if (inFlightWrite == kj::none && inFlightClose == kj::none && flags.started) { finishErroring(js, kj::mv(self)); } } template void WritableImpl::updateBackpressure(jsg::Lock& js) { KJ_ASSERT(isWritable()); KJ_ASSERT(!isCloseQueuedOrInFlight()); bool bp = getDesiredSize() <= 0; if (bp != flags.backpressure) { flags.backpressure = bp; KJ_IF_SOME(owner, tryGetOwner()) { owner.updateBackpressure(js, flags.backpressure); } } } template jsg::Promise WritableImpl::write( jsg::Lock& js, jsg::Ref self, v8::Local value) { size_t size = 1; KJ_IF_SOME(sizeFunc, algorithms.size) { kj::Maybe failure; JSG_TRY(js) { size = sizeFunc(js, value); } JSG_CATCH(exception) { startErroring(js, self.addRef(), exception.getHandle(js)); failure = kj::mv(exception); } KJ_IF_SOME(exception, failure) { return js.rejectedPromise(kj::mv(exception)); } } // Per spec (WritableStreamDefaultWriterWrite step 5), after calling the size // algorithm, re-check that the stream is still locked to a writer. If // releaseLock() was called from within strategy.size(), the write must be // rejected. This check must occur before any state checks, as the stream // state may still appear writable even after the writer was released. if (FeatureFlags::get(js).getWritableStreamSpecCompliantWriter()) { KJ_IF_SOME(owner, tryGetOwner()) { if (!owner.isLockedToWriter()) { return js.rejectedPromise( js.v8TypeError("This WritableStream writer has been released."_kjc)); } } } KJ_IF_SOME(error, state.template tryGetUnsafe()) { return js.rejectedPromise(error.addRef(js)); } if (isCloseQueuedOrInFlight() || state.template is()) { return js.rejectedPromise(js.v8TypeError("This ReadableStream is closed."_kj)); } KJ_IF_SOME(erroring, state.template tryGetUnsafe()) { return js.rejectedPromise(erroring.reason.addRef(js)); } KJ_ASSERT(isWritable()); auto prp = js.newPromiseAndResolver(); writeRequests.push_back(WriteRequest{ .resolver = kj::mv(prp.resolver), .value = js.v8Ref(value), .size = size, }); amountBuffered += size; updateBackpressure(js); advanceQueueIfNeeded(js, kj::mv(self)); return kj::mv(prp.promise); } template void WritableImpl::visitForGc(jsg::GcVisitor& visitor) { state.visitForGc(visitor); visitor.visit(inFlightWrite, inFlightClose, closeRequest, algorithms, signal); KJ_IF_SOME(pendingAbort, maybePendingAbort) { visitor.visit(*pendingAbort); } visitor.visitAll(writeRequests); } template bool WritableImpl::isWritable() const { return state.isActive(); } template void WritableImpl::cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason) { for (auto& write: writeRequests) { write.resolver.reject(js, reason); } writeRequests.clear(); } // ====================================================================================== namespace { template struct ReadableState { Controller controller; kj::Own consumer; ReadableStreamJsController& owner; ReadableState(Controller controller, kj::Own consumer, ReadableStreamJsController& owner) : controller(kj::mv(controller)), consumer(kj::mv(consumer)), owner(owner) {} ReadableState(Controller controller, Queue::ConsumerImpl::StateListener& listener, ReadableStreamJsController& owner) : ReadableState(controller.addRef(), controller->getConsumer(listener), owner) {} ReadableState clone(jsg::Lock& js, Queue::ConsumerImpl::StateListener& listener, ReadableStreamJsController& owner) { return ReadableState(controller.addRef(), consumer->clone(js, listener), owner); } }; struct ValueReadable final: private api::ValueQueue::ConsumerImpl::StateListener { using State = ReadableState; kj::Maybe state; bool reading = false; bool pendingCancel = false; JSG_MEMORY_INFO(ValueReadable) { KJ_IF_SOME(s, state) { tracker.trackField("controller", s.controller); tracker.trackField("consumer", s.consumer); } } void visitForGc(jsg::GcVisitor& visitor) { KJ_IF_SOME(s, state) { visitor.visit(s.controller, *s.consumer); } } ValueReadable(DefaultController controller, ReadableStreamJsController& owner) : state(State(kj::mv(controller), *this, owner)) {} ValueReadable(jsg::Lock& js, ReadableStreamJsController& owner, ValueReadable& other) : state(KJ_ASSERT_NONNULL(other.state).clone(js, *this, owner)) {} KJ_DISALLOW_COPY_AND_MOVE(ValueReadable); void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) { KJ_IF_SOME(s, state) { s.consumer->cancelPendingReads(js, reason); } } kj::Own clone(jsg::Lock& js, ReadableStreamJsController& owner) { // A single ReadableStreamDefaultController can have multiple consumers. // When the ValueReadable constructor is used, the new consumer is added // and starts to receive new data that becomes enqueued. When clone // is used, any state currently held by this consumer is copied to the // new consumer. return kj::heap(js, owner, *this); } jsg::Promise read(jsg::Lock& js) { KJ_IF_SOME(s, state) { auto prp = js.newPromiseAndResolver(); reading = true; s.consumer->read(js, ValueQueue::ReadRequest{ .resolver = kj::mv(prp.resolver), }); reading = false; if (pendingCancel) { // If we were canceled while reading, we need to drop our state now. state = kj::none; pendingCancel = false; } return kj::mv(prp.promise); } // We are canceled! There's nothing to do. return js.resolvedPromise(ReadResult{.done = true}); } jsg::Promise drainingRead(jsg::Lock& js, size_t maxRead) { KJ_IF_SOME(s, state) { // Note: We do NOT call beginOperation()/endOperation() here. The caller // (ReadableStreamJsController::drainingRead) manages the operation scope // around both this call and the returned promise's lifetime. If we added // our own beginOperation/endOperation here, the endOperation would fire // before the caller's wrapDrainingRead could set up its .then() callbacks, // potentially destroying the Consumer while the returned promise still has // dangling this-capturing callbacks from consumer->drainingRead(). return s.consumer->drainingRead(js, maxRead); } // We are canceled! Return done with empty chunks. return js.resolvedPromise(DrainingReadResult{ .chunks = kj::Array>(), .done = true, }); } jsg::Promise cancel(jsg::Lock& js, jsg::Optional> maybeReason) { // When a ReadableStream is canceled, the expected behavior is that the underlying // controller is notified and the cancel algorithm on the underlying source is // called. When there are multiple ReadableStreams sharing consumption of a // controller, however, it should act as a shared pointer of sorts, canceling // the underlying controller only when the last reader is canceled. // Here, we rely on the controller implementing the correct behavior since it owns // the queue that knows about all of the attached consumers. if (pendingCancel) return js.resolvedPromise(); KJ_IF_SOME(s, state) { // Check if there's a pending draining read before calling cancel, since cancel // will resolve the pending read and we need to know if we should defer destruction. bool hasPendingDrainingRead = s.consumer->hasPendingDrainingRead(); s.consumer->cancel(js, maybeReason); auto promise = s.controller->cancel(js, kj::mv(maybeReason)); // If we're currently in a read (sync or draining), we need to wait for that to // finish before dropping our state. For draining reads, the promise callbacks // capture 'this' (the Consumer) to clear hasPendingDrainingRead. If we destroy // the state now, those callbacks will UAF. if (reading || hasPendingDrainingRead) { pendingCancel = true; } else { state = kj::none; } return kj::mv(promise); } return js.resolvedPromise(); } void onConsumerClose(jsg::Lock& js) override { // Called by the consumer when a state change to closed happens. // We need to notify the owner. Note that the owner may drop this // readable in doClose so it is not safe to access anything on this // after calling doClose. KJ_IF_SOME(s, state) { s.owner.doClose(js); } } void onConsumerError(jsg::Lock& js, jsg::Value reason) override { // Called by the consumer when a state change to errored happens. // We need to notify the owner. Note that the owner may drop this // readable in doClose so it is not safe to access anything on this // after calling doError. KJ_IF_SOME(s, state) { s.owner.doError(js, reason.getHandle(js)); } } bool onConsumerWantsData(jsg::Lock& js) override { // Called by the consumer when it has a queued pending read and needs // data to be provided to fulfill it. We need to notify the controller // to initiate pulling to provide the data. // Returns true if the pull completed synchronously (meaning more pumping // might yield additional synchronous data), false otherwise. KJ_IF_SOME(s, state) { // Save a reference to the owner before calling pull. The pull callback // may trigger close/error which could destroy this ValueReadable. By // using beginOperation(), we ensure doClose/doError defers the // actual destruction until after we return. ReadableStreamJsController& owner = s.owner; owner.state.beginOperation(); // For draining reads, use forcePull to bypass backpressure checks. // This ensures we pull all available data regardless of highWaterMark. if (s.consumer->hasPendingDrainingRead()) { s.controller->forcePull(js); } else { s.controller->pull(js); } // Check if state is still valid BEFORE calling endOperation(), // because that call may destroy this ValueReadable if close was deferred. bool result = state.map([](State& s2) { return !s2.controller->isPulling(); }).orDefault(false); // Process any deferred close/error. This may destroy this ValueReadable. if (owner.state.endOperation()) { // A pending state was applied. Call the appropriate callback. if (owner.state.template is()) { owner.lock.onClose(js); } else if (owner.state.template is()) { KJ_IF_SOME(err, owner.state.template tryGetUnsafe()) { owner.lock.onError(js, err.getHandle(js)); } } } return result; } return false; } kj::Maybe getDesiredSize() { KJ_IF_SOME(s, state) { return s.controller->getDesiredSize(); } return kj::none; } bool canCloseOrEnqueue() { return state.map([](State& s) { return s.controller->canCloseOrEnqueue(); }).orDefault(false); } kj::Maybe getControllerRef() { return state.map([](State& s) { return s.controller.addRef(); }); } }; struct ByteReadable final: private api::ByteQueue::ConsumerImpl::StateListener { using State = ReadableState; kj::Maybe state; kj::Maybe autoAllocateChunkSize; bool pendingCancel = false; JSG_MEMORY_INFO(ByteReadable) { KJ_IF_SOME(s, state) { tracker.trackField("controller", s.controller); tracker.trackField("consumer", s.consumer); } } void visitForGc(jsg::GcVisitor& visitor) { KJ_IF_SOME(s, state) { visitor.visit(s.controller, *s.consumer); } } ByteReadable(ByobController controller, ReadableStreamJsController& owner, kj::Maybe autoAllocateChunkSize) : state(State(kj::mv(controller), *this, owner)), autoAllocateChunkSize(autoAllocateChunkSize) {} ByteReadable(jsg::Lock& js, ReadableStreamJsController& owner, ByteReadable& other) : state(KJ_ASSERT_NONNULL(other.state).clone(js, *this, owner)), autoAllocateChunkSize(other.autoAllocateChunkSize) {} KJ_DISALLOW_COPY_AND_MOVE(ByteReadable); void cancelPendingReads(jsg::Lock& js, jsg::JsValue reason) { KJ_IF_SOME(s, state) { s.consumer->cancelPendingReads(js, reason); } } // A single ReadableByteStreamController can have multiple consumers. // When the ByteReadable constructor is used, the new consumer is added // and starts to receive new data that becomes enqueued. When clone // is used, any state currently held by this consumer is copied to the // new consumer. kj::Own clone(jsg::Lock& js, ReadableStreamJsController& owner) { return kj::heap(js, owner, *this); } jsg::Promise read( jsg::Lock& js, kj::Maybe byobOptions) { KJ_IF_SOME(s, state) { auto prp = js.newPromiseAndResolver(); KJ_IF_SOME(byob, byobOptions) { jsg::BufferSource source(js, byob.bufferView.getHandle(js)); // If atLeast is not given, then by default it is the element size of the view // that we were given. If atLeast is given, we make sure that it is aligned // with the element size. No matter what, atLeast cannot be less than 1. auto atLeast = kj::max(source.getElementSize(), byob.atLeast.orDefault(1)); atLeast = kj::max(1, atLeast - (atLeast % source.getElementSize())); s.consumer->read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = jsg::BufferSource(js, source.detach(js)), .atLeast = atLeast, .type = ByteQueue::ReadRequest::Type::BYOB, })); } else KJ_IF_SOME(chunkSize, autoAllocateChunkSize) { // autoAllocateChunkSize is set, so we allocate a buffer and do a BYOB read. // This makes the buffer available to the underlying source via controller.byobRequest. KJ_IF_SOME(store, jsg::BufferSource::tryAlloc(js, chunkSize)) { // Ensure that the handle is created here so that the size of the buffer // is accounted for in the isolate memory tracking. s.consumer->read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = kj::mv(store), .type = ByteQueue::ReadRequest::Type::BYOB, })); } else { prp.resolver.reject(js, js.v8Error("Failed to allocate buffer for read.")); } } else { // autoAllocateChunkSize is not set. Per spec, we do a DEFAULT read which means // the underlying source's pull method won't get a byobRequest. It must use // controller.enqueue() to provide data instead. constexpr size_t kDefaultReadSize = 16384; // 16KB default buffer KJ_IF_SOME(store, jsg::BufferSource::tryAlloc(js, kDefaultReadSize)) { s.consumer->read(js, ByteQueue::ReadRequest(kj::mv(prp.resolver), { .store = kj::mv(store), .type = ByteQueue::ReadRequest::Type::DEFAULT, })); } else { prp.resolver.reject(js, js.v8Error("Failed to allocate buffer for read.")); } } return kj::mv(prp.promise); } // We are canceled! There's nothing else to do. KJ_IF_SOME(byob, byobOptions) { // If a BYOB buffer was given, we need to give it back wrapped in a TypedArray // whose size is set to zero. jsg::BufferSource source(js, byob.bufferView.getHandle(js)); auto store = source.detach(js); store.consume(store.size()); return js.resolvedPromise(ReadResult{ .value = js.v8Ref(store.createHandle(js)), .done = true, }); } else { return js.resolvedPromise(ReadResult{.done = true}); } } jsg::Promise drainingRead(jsg::Lock& js, size_t maxRead) { KJ_IF_SOME(s, state) { // Note: We do NOT call beginOperation()/endOperation() here. The caller // (ReadableStreamJsController::drainingRead) manages the operation scope // around both this call and the returned promise's lifetime. See the // comment in ValueReadable::drainingRead for the detailed explanation. return s.consumer->drainingRead(js, maxRead); } // We are canceled! Return done with empty chunks. return js.resolvedPromise(DrainingReadResult{ .chunks = kj::Array>(), .done = true, }); } // When a ReadableStream is canceled, the expected behavior is that the underlying // controller is notified and the cancel algorithm on the underlying source is // called. When there are multiple ReadableStreams sharing consumption of a // controller, however, it should act as a shared pointer of sorts, canceling // the underlying controller only when the last reader is canceled. // Here, we rely on the controller implementing the correct behavior since it owns // the queue that knows about all of the attached consumers. jsg::Promise cancel(jsg::Lock& js, jsg::Optional> maybeReason) { if (pendingCancel) return js.resolvedPromise(); KJ_IF_SOME(s, state) { // Check if there's a pending draining read before calling cancel, since cancel // will resolve the pending read and we need to know if we should defer destruction. bool hasPendingDrainingRead = s.consumer->hasPendingDrainingRead(); s.consumer->cancel(js, maybeReason); auto promise = s.controller->cancel(js, kj::mv(maybeReason)); // If there's a pending draining read, we need to wait for it to finish before // dropping our state. The draining read's promise callbacks capture 'this' (the // Consumer) to clear hasPendingDrainingRead. If we destroy the state now, those // callbacks will UAF. if (hasPendingDrainingRead) { pendingCancel = true; } else { state = kj::none; } return kj::mv(promise); } return js.resolvedPromise(); } void onConsumerClose(jsg::Lock& js) override { // Note that the owner may drop this readable in doClose so it // is not safe to access anything on this after calling doClose. KJ_IF_SOME(s, state) { s.owner.doClose(js); } } void onConsumerError(jsg::Lock& js, jsg::Value reason) override { // Note that the owner may drop this readable in doClose so it // is not safe to access anything on this after calling doError. KJ_IF_SOME(s, state) { s.owner.doError(js, reason.getHandle(js)); }; } // Called by the consumer when it has a queued pending read and needs // data to be provided to fulfill it. We need to notify the controller // to initiate pulling to provide the data. // Returns true if the pull completed synchronously (meaning more pumping // might yield additional synchronous data), false otherwise. bool onConsumerWantsData(jsg::Lock& js) override { KJ_IF_SOME(s, state) { // Save a reference to the owner before calling pull. The pull callback // may trigger close/error which could destroy this ByteReadable. By // using beginOperation(), we ensure doClose/doError defers the // actual destruction until after we return. ReadableStreamJsController& owner = s.owner; owner.state.beginOperation(); // For draining reads, use forcePull to bypass backpressure checks. // This ensures we pull all available data regardless of highWaterMark. if (s.consumer->hasPendingDrainingRead()) { s.controller->forcePull(js); } else { s.controller->pull(js); } // Check if state is still valid BEFORE calling endOperation(), // because that call may destroy this ByteReadable if close was deferred. bool result = state.map([](State& s2) { return !s2.controller->isPulling(); }).orDefault(false); // Process any deferred close/error. This may destroy this ByteReadable. if (owner.state.endOperation()) { // A pending state was applied. Call the appropriate callback. if (owner.state.template is()) { owner.lock.onClose(js); } else if (owner.state.template is()) { KJ_IF_SOME(err, owner.state.template tryGetUnsafe()) { owner.lock.onError(js, err.getHandle(js)); } } } return result; } return false; } kj::Maybe getDesiredSize() { KJ_IF_SOME(s, state) { return s.controller->getDesiredSize(); } return kj::none; } bool canCloseOrEnqueue() { return state.map([](State& s) { return s.controller->canCloseOrEnqueue(); }).orDefault(false); } kj::Maybe getControllerRef() { return state.map([](State& state) { return state.controller.addRef(); }); } }; } // namespace // ======================================================================================= ReadableStreamDefaultController::ReadableStreamDefaultController( UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy) : ioContext(tryGetIoContext()), impl(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {} kj::Maybe ReadableStreamDefaultController::getMaybeErrorState( jsg::Lock& js) { KJ_IF_SOME(errored, impl.state.tryGetUnsafe()) { return errored.addRef(js); } return kj::none; } void ReadableStreamDefaultController::start(jsg::Lock& js) { impl.start(js, JSG_THIS); } bool ReadableStreamDefaultController::canCloseOrEnqueue() { return impl.canCloseOrEnqueue(); } bool ReadableStreamDefaultController::hasBackpressure() { return !impl.shouldCallPull(); } kj::Maybe ReadableStreamDefaultController::getDesiredSize() { return impl.getDesiredSize(); } void ReadableStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) { visitor.visit(impl); } jsg::Promise ReadableStreamDefaultController::cancel( jsg::Lock& js, jsg::Optional> maybeReason) { return impl.cancel(js, JSG_THIS, maybeReason.orDefault([&] { return js.v8Undefined(); })); } void ReadableStreamDefaultController::close(jsg::Lock& js) { impl.close(js); } void ReadableStreamDefaultController::enqueue( jsg::Lock& js, jsg::Optional> chunk) { // Hold a strong reference to prevent this controller from being freed if the // user-provided size algorithm (below) re-enters JS and errors the controller // through a side-channel (e.g. TransformStreamDefaultController::error() // dropping all external jsg::Refs to this controller). auto self = JSG_THIS; auto value = chunk.orDefault(js.undefined()); JSG_REQUIRE(impl.canCloseOrEnqueue(), TypeError, "Unable to enqueue"); size_t size = 1; bool errored = false; KJ_IF_SOME(sizeFunc, impl.algorithms.size) { js.tryCatch([&] { size = sizeFunc(js, value); }, [&](jsg::Value exception) { impl.doError(js, kj::mv(exception)); errored = true; }); } // Re-check canCloseOrEnqueue: the size callback may have errored us without // throwing (e.g. by calling transformController.error()), in which case // `errored` is still false but the impl state has transitioned to Errored. if (!errored && impl.canCloseOrEnqueue()) { impl.enqueue(js, kj::rc(js.v8Ref(value), size), kj::mv(self)); } } void ReadableStreamDefaultController::error(jsg::Lock& js, v8::Local reason) { impl.doError(js, js.v8Ref(reason)); } // When a consumer receives a read request, but does not have the data available to // fulfill the request, the consumer will call pull on the controller to pull that // data if needed. void ReadableStreamDefaultController::pull(jsg::Lock& js) { impl.pullIfNeeded(js, JSG_THIS); } void ReadableStreamDefaultController::forcePull(jsg::Lock& js) { impl.forcePullIfNeeded(js, JSG_THIS); } kj::Own ReadableStreamDefaultController::getConsumer( kj::Maybe stateListener) { return impl.getConsumer(stateListener); } // ====================================================================================== ReadableStreamBYOBRequest::Impl::Impl(jsg::Lock& js, kj::Own readRequest, kj::Rc> controller) : readRequest(kj::mv(readRequest)), controller(kj::mv(controller)), view(js.v8Ref(this->readRequest->getView(js))), originalBufferByteLength(this->readRequest->getOriginalBufferByteLength(js)), originalByteOffsetPlusBytesFilled(this->readRequest->getOriginalByteOffsetPlusBytesFilled()) { } void ReadableStreamBYOBRequest::Impl::updateView(jsg::Lock& js) { jsg::check(view.getHandle(js)->Buffer()->Detach(v8::Local())); view = js.v8Ref(readRequest->getView(js)); } void ReadableStreamBYOBRequest::visitForGc(jsg::GcVisitor& visitor) { KJ_IF_SOME(impl, maybeImpl) { visitor.visit(impl.view); } } ReadableStreamBYOBRequest::ReadableStreamBYOBRequest(jsg::Lock& js, kj::Own readRequest, kj::Rc> controller) : ioContext(tryGetIoContext()), maybeImpl(Impl(js, kj::mv(readRequest), kj::mv(controller))) {} kj::Maybe ReadableStreamBYOBRequest::getAtLeast() { KJ_IF_SOME(impl, maybeImpl) { return impl.readRequest->getAtLeast(); } return kj::none; } kj::Maybe> ReadableStreamBYOBRequest::getView(jsg::Lock& js) { KJ_IF_SOME(impl, maybeImpl) { return impl.view.addRef(js); } return kj::none; } void ReadableStreamBYOBRequest::invalidate(jsg::Lock& js) { KJ_IF_SOME(impl, maybeImpl) { // If the user code happened to have retained a reference to the view or // the buffer, we need to detach it so that those references cannot be used // to modify or observe modifications. jsg::check(impl.view.getHandle(js)->Buffer()->Detach(v8::Local())); impl.controller->runIfAlive( [](ReadableByteStreamController& controller) { controller.maybeByobRequest = kj::none; }); } maybeImpl = kj::none; } void ReadableStreamBYOBRequest::respond(jsg::Lock& js, int bytesWritten) { auto& impl = JSG_REQUIRE_NONNULL( maybeImpl, TypeError, "This ReadableStreamBYOBRequest has been invalidated."); JSG_REQUIRE(impl.controller->isValid(), Error, "The ReadableStreamBYOBRequest is invalid."); JSG_REQUIRE(impl.view.getHandle(js)->ByteLength() > 0, TypeError, "Cannot respond with a zero-length or detached view"); impl.controller->runIfAlive([&](ReadableByteStreamController& controller) { if (!controller.canCloseOrEnqueue()) { JSG_REQUIRE(bytesWritten == 0, TypeError, "The bytesWritten must be zero after the stream is closed."); KJ_ASSERT(impl.readRequest->isInvalidated()); invalidate(js); } else { bool shouldInvalidate = false; if (impl.readRequest->isInvalidated() && controller.impl.consumerCount() >= 1) { // While this particular request may be invalidated, there are still // other branches we can push the data to. Let's do so. jsg::BufferSource source(js, impl.view.getHandle(js)); auto entry = kj::rc(jsg::BufferSource(js, source.detach(js))); controller.impl.enqueue(js, kj::mv(entry), controller.getSelf()); } else { JSG_REQUIRE(bytesWritten > 0, TypeError, "The bytesWritten must be more than zero while the stream is open."); if (impl.readRequest->respond(js, bytesWritten)) { // The read request was fulfilled, we need to invalidate. shouldInvalidate = true; } else { // The response did not fulfill the minimum requirements of the read. // We do not want to invalidate the read request and we need to update the // view so that on the next read the view will be properly adjusted. impl.updateView(js); } } controller.pull(js); if (shouldInvalidate) { invalidate(js); } } }); } void ReadableStreamBYOBRequest::respondWithNewView(jsg::Lock& js, jsg::BufferSource view) { auto& impl = JSG_REQUIRE_NONNULL( maybeImpl, TypeError, "This ReadableStreamBYOBRequest has been invalidated."); JSG_REQUIRE(impl.controller->isValid(), Error, "The ReadableStreamBYOBRequest is invalid."); impl.controller->runIfAlive([&](ReadableByteStreamController& controller) { if (!controller.canCloseOrEnqueue()) { JSG_REQUIRE(view.size() == 0, TypeError, "The view byte length must be zero after the stream is closed."); if (FeatureFlags::get(js).getPedanticWpt()) { // Per the spec, when the stream is closed: // 1. The view byte length must be zero (TypeError if not) // 2. The underlying buffer must not be detached (TypeError) // 3. The buffer byte length must not be zero (RangeError) // 4. The buffer byte length must match the original (RangeError) auto handle = view.getHandle(js); auto buffer = handle->IsArrayBuffer() ? handle.As() : handle.As()->Buffer(); JSG_REQUIRE( !buffer->WasDetached(), TypeError, "The underlying ArrayBuffer has been detached."); JSG_REQUIRE(view.canDetach(js), TypeError, "Unable to use non-detachable ArrayBuffer."); // Use the stored values since the ByobRequest may have been invalidated during close. auto actualBufferByteLength = buffer->ByteLength(); JSG_REQUIRE( actualBufferByteLength != 0, RangeError, "The underlying ArrayBuffer is zero-length."); JSG_REQUIRE(actualBufferByteLength == impl.originalBufferByteLength, RangeError, "The underlying ArrayBuffer is not the correct length."); // The view's byte offset must match the original byte offset plus bytes filled. auto viewByteOffset = handle->IsArrayBuffer() ? 0 : handle.As()->ByteOffset(); JSG_REQUIRE(viewByteOffset == impl.originalByteOffsetPlusBytesFilled, RangeError, "The view has an invalid byte offset."); } else { KJ_ASSERT(impl.readRequest->isInvalidated()); } invalidate(js); } else { bool shouldInvalidate = false; if (impl.readRequest->isInvalidated() && controller.impl.consumerCount() >= 1) { // While this particular request may be invalidated, there are still // other branches we can push the data to. Let's do so. auto entry = kj::rc(jsg::BufferSource(js, view.detach(js))); controller.impl.enqueue(js, kj::mv(entry), controller.getSelf()); } else { JSG_REQUIRE(view.size() > 0, TypeError, "The view byte length must be more than zero while the stream is open."); if (impl.readRequest->respondWithNewView(js, kj::mv(view))) { // The read request was fulfilled, we need to invalidate. shouldInvalidate = true; } else { // The response did not fulfill the minimum requirements of the read. // We do not want to invalidate the read request and we need to update the // view so that on the next read the view will be properly adjusted. impl.updateView(js); } } controller.pull(js); if (shouldInvalidate) { invalidate(js); } } }); } bool ReadableStreamBYOBRequest::isPartiallyFulfilled() { KJ_IF_SOME(impl, maybeImpl) { return impl.readRequest->isPartiallyFulfilled(); } return false; } // ====================================================================================== ReadableByteStreamController::ReadableByteStreamController( UnderlyingSource underlyingSource, StreamQueuingStrategy queuingStrategy) : weakSelf(kj::rc>( kj::Badge{}, *this)), ioContext(tryGetIoContext()), impl(kj::mv(underlyingSource), kj::mv(queuingStrategy)) {} ReadableByteStreamController::~ReadableByteStreamController() noexcept(false) { weakSelf->invalidate(); } void ReadableByteStreamController::start(jsg::Lock& js) { impl.start(js, JSG_THIS); } bool ReadableByteStreamController::canCloseOrEnqueue() { return impl.canCloseOrEnqueue(); } bool ReadableByteStreamController::hasBackpressure() { return !impl.shouldCallPull(); } kj::Maybe ReadableByteStreamController::getDesiredSize() { return impl.getDesiredSize(); } void ReadableByteStreamController::visitForGc(jsg::GcVisitor& visitor) { visitor.visit(maybeByobRequest, impl); } jsg::Promise ReadableByteStreamController::cancel( jsg::Lock& js, jsg::Optional> maybeReason) { KJ_IF_SOME(byobRequest, maybeByobRequest) { if (impl.consumerCount() == 1) { byobRequest->invalidate(js); } } return impl.cancel(js, JSG_THIS, maybeReason.orDefault(js.undefined())); } void ReadableByteStreamController::close(jsg::Lock& js) { KJ_IF_SOME(byobRequest, maybeByobRequest) { JSG_REQUIRE(!byobRequest->isPartiallyFulfilled(), TypeError, "This ReadableStream was closed with a partial read pending."); } else if (FeatureFlags::get(js).getPedanticWpt()) { // If maybeByobRequest is not set, check if there's a pending byob request. // If so, materialize it before closing so it remains accessible after // the state changes to Closed. This is required by the spec for proper // respondWithNewView() error handling in the closed state. // Only do this if the queue doesn't have a partially fulfilled read. KJ_IF_SOME(queue, impl.state.tryGetUnsafe()) { if (!queue.hasPartiallyFulfilledRead()) { getByobRequest(js); } } } impl.close(js); } void ReadableByteStreamController::enqueue(jsg::Lock& js, jsg::BufferSource chunk) { // Hold a strong reference up front. Operations below (invalidate, detach) touch // the JS heap and C++ argument evaluation order is unspecified, so JSG_THIS as a // function argument would not reliably precede chunk.detach(js). auto self = JSG_THIS; JSG_REQUIRE(chunk.size() > 0, TypeError, "Cannot enqueue a zero-length ArrayBuffer."); JSG_REQUIRE(chunk.canDetach(js), TypeError, "The provided ArrayBuffer must be detachable."); JSG_REQUIRE(impl.canCloseOrEnqueue(), TypeError, "This ReadableByteStreamController is closed."); KJ_IF_SOME(byobRequest, maybeByobRequest) { KJ_IF_SOME(view, byobRequest->getView(js)) { JSG_REQUIRE(view.getHandle(js)->ByteLength() > 0, TypeError, "The byobRequest.view is zero-length or was detached"); } byobRequest->invalidate(js); } impl.enqueue(js, kj::rc(jsg::BufferSource(js, chunk.detach(js))), kj::mv(self)); } void ReadableByteStreamController::error(jsg::Lock& js, v8::Local reason) { impl.doError(js, js.v8Ref(reason)); } kj::Maybe> ReadableByteStreamController::getByobRequest( jsg::Lock& js) { if (maybeByobRequest == kj::none) { KJ_IF_SOME(queue, impl.state.tryGetUnsafe()) { KJ_IF_SOME(pendingByob, queue.nextPendingByobReadRequest()) { maybeByobRequest = js.alloc(js, kj::mv(pendingByob), weakSelf.addRef()); } } else { return kj::none; } } return maybeByobRequest.map( [&](jsg::Ref& req) { return req.addRef(); }); } // When a consumer receives a read request, but does not have the data available to // fulfill the request, the consumer will call pull on the controller to pull that // data if needed. void ReadableByteStreamController::pull(jsg::Lock& js) { impl.pullIfNeeded(js, JSG_THIS); } void ReadableByteStreamController::forcePull(jsg::Lock& js) { impl.forcePullIfNeeded(js, JSG_THIS); } kj::Own ReadableByteStreamController::getConsumer( kj::Maybe stateListener) { return impl.getConsumer(stateListener); } // ====================================================================================== ReadableStreamJsController::ReadableStreamJsController(): ioContext(tryGetIoContext()) {} ReadableStreamJsController::ReadableStreamJsController(StreamStates::Closed closed) : ioContext(tryGetIoContext()) { state.transitionTo(); } ReadableStreamJsController::ReadableStreamJsController(StreamStates::Errored errored) : ioContext(tryGetIoContext()) { state.transitionTo(kj::mv(errored)); } ReadableStreamJsController::ReadableStreamJsController(jsg::Lock& js, ValueReadable& consumer) : ioContext(tryGetIoContext()) { state.transitionTo>(consumer.clone(js, *this)); } ReadableStreamJsController::ReadableStreamJsController(jsg::Lock& js, ByteReadable& consumer) : ioContext(tryGetIoContext()) { state.transitionTo>(consumer.clone(js, *this)); } jsg::Ref ReadableStreamJsController::addRef() { return KJ_REQUIRE_NONNULL(owner).addRef(); } jsg::Promise ReadableStreamJsController::cancel( jsg::Lock& js, jsg::Optional> maybeReason) { disturbed = true; const auto doCancel = [&](auto& consumer) { auto reason = js.v8Ref(maybeReason.orDefault([&] { return js.v8Undefined(); })); KJ_DEFER(doClose(js)); return consumer->cancel(js, reason.getHandle(js)); }; // Check for pending state first (deferred close/error during a read operation) if (state.pendingStateIs()) { return js.resolvedPromise(); } KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { return js.rejectedPromise(pendingError.addRef(js)); } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream not yet set up, treat as closed. return js.resolvedPromise(); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return js.resolvedPromise(); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.rejectedPromise(errored.addRef(js)); } KJ_CASE_ONEOF(consumer, kj::Own) { if (canceling) return js.resolvedPromise(); canceling = true; return doCancel(consumer); } KJ_CASE_ONEOF(consumer, kj::Own) { if (canceling) return js.resolvedPromise(); canceling = true; return doCancel(consumer); } } KJ_UNREACHABLE; } // Finalizes the closed state of this ReadableStream. The connection to the underlying // controller is released with no further action. Importantly, this method is triggered // by the underlying controller as a result of that controller closing or being canceled. // We detach ourselves from the underlying controller by releasing the ValueReadable or // ByteReadable in the state and changing that to closed. // We also clean up other state here. void ReadableStreamJsController::doClose(jsg::Lock& js) { // If already in a terminal state, nothing to do. if (state.isTerminal()) return; // deferTransitionTo will defer if an operation is in progress, otherwise transition immediately. // Returns true if transition happened immediately. if (state.deferTransitionTo()) { lock.onClose(js); } // If deferred, lock.onClose will be called when the pending state is applied // via applyPendingState in deferControllerStateChange. } // As with doClose(), doError() finalizes the error state of this ReadableStream. // The connection to the underlying controller is released with no further action. // This method is triggered by the underlying controller as a result of that controller // erroring. We detach ourselves from the underlying controller by releasing the ValueReadable // or ByteReadable in the state and changing that to errored. // We also clean up other state here. void ReadableStreamJsController::doError(jsg::Lock& js, v8::Local reason) { // If already in a terminal state, nothing to do. if (state.isTerminal()) return; // deferTransitionTo will defer if an operation is in progress, otherwise transition immediately. // Returns true if transition happened immediately. if (state.deferTransitionTo(js.v8Ref(reason))) { lock.onError(js, reason); } // If deferred, lock.onError will be called when the pending state is applied // via applyPendingState in deferControllerStateChange. } bool ReadableStreamJsController::isByteOriented() const { return state.is>(); } bool ReadableStreamJsController::isClosedOrErrored() const { // Check if we're in a terminal state or have one pending return state.isTerminal() || state.hasPendingState(); } bool ReadableStreamJsController::isClosed() const { // Check current state first, then pending state if (state.is()) return true; return state.pendingStateIs(); } bool ReadableStreamJsController::isDisturbed() { return disturbed; } bool ReadableStreamJsController::isLockedToReader() const { return lock.isLockedToReader(); } bool ReadableStreamJsController::lockReader(jsg::Lock& js, Reader& reader) { return lock.lockReader(js, *this, reader); } jsg::Promise ReadableStreamJsController::pipeTo( jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) { KJ_DASSERT(!isLockedToReader()); KJ_DASSERT(!destination.isLockedToWriter()); disturbed = true; KJ_IF_SOME(promise, destination.tryPipeFrom(js, addRef(), kj::mv(options))) { return kj::mv(promise); } return js.rejectedPromise( js.v8TypeError("This ReadableStream cannot be piped to this WritableStream"_kj)); } kj::Maybe> ReadableStreamJsController::read( jsg::Lock& js, kj::Maybe maybeByobOptions) { disturbed = true; KJ_IF_SOME(byobOptions, maybeByobOptions) { byobOptions.detachBuffer = true; auto view = byobOptions.bufferView.getHandle(js); if (!view->Buffer()->IsDetachable()) { return js.rejectedPromise( js.v8TypeError("Unabled to use non-detachable ArrayBuffer."_kj)); } if (view->ByteLength() == 0 || view->Buffer()->ByteLength() == 0) { return js.rejectedPromise( js.v8TypeError("Unable to use a zero-length ArrayBuffer."_kj)); } // Check for pending error first (deferred error during a prior read operation) KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { return js.rejectedPromise(pendingError.addRef(js)); } if (state.is() || state.pendingStateIs()) { // If it is a BYOB read, then the spec requires that we return an empty // view of the same type provided, that uses the same backing memory // as that provided, but with zero-length. auto source = jsg::BufferSource(js, byobOptions.bufferView.getHandle(js)); auto store = source.detach(js); store.consume(store.size()); return js.resolvedPromise(ReadResult{ .value = js.v8Ref(store.createHandle(js)), .done = true, }); } } // Check for pending state (deferred close/error during a prior read operation) if (state.pendingStateIs()) { // The closed state for BYOB reads is handled in the maybeByobOptions check above. KJ_ASSERT(maybeByobOptions == kj::none); return js.resolvedPromise(ReadResult{.done = true}); } KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { return js.rejectedPromise(pendingError.addRef(js)); } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream not yet set up, treat as closed. KJ_ASSERT(maybeByobOptions == kj::none); return js.resolvedPromise(ReadResult{.done = true}); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { // The closed state for BYOB reads is handled in the maybeByobOptions check above. KJ_ASSERT(maybeByobOptions == kj::none); return js.resolvedPromise(ReadResult{.done = true}); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.rejectedPromise(errored.addRef(js)); } KJ_CASE_ONEOF(consumer, kj::Own) { // The ReadableStreamDefaultController does not support ByobOptions. // It should never happen, but let's make sure. KJ_ASSERT(maybeByobOptions == kj::none); return deferControllerStateChange(js, *this, [&]() mutable { return consumer->read(js); }); } KJ_CASE_ONEOF(consumer, kj::Own) { return deferControllerStateChange( js, *this, [&]() mutable { return consumer->read(js, kj::mv(maybeByobOptions)); }); } } KJ_UNREACHABLE; } kj::Maybe> ReadableStreamJsController::drainingRead( jsg::Lock& js, size_t maxRead) { disturbed = true; // Check for pending state first (deferred close/error during a prior read operation) if (state.pendingStateIs()) { return js.resolvedPromise(DrainingReadResult{ .chunks = kj::Array>(), .done = true, }); } KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { return js.rejectedPromise(pendingError.addRef(js)); } // Like deferControllerStateChange for regular reads, we need to prevent the controller // state from being destroyed while a draining read's promise callbacks are pending. // The drainingRead implementation captures `this` (the Consumer) in promise lambdas to // clear hasPendingDrainingRead. If the state is changed (destroying the Consumer) before // those callbacks run, we get a use-after-free. // // CRITICAL: state.beginOperation() MUST be called BEFORE consumer->drainingRead(), not // after. The consumer->drainingRead() call may trigger onConsumerWantsData -> forcePull // -> close/error, which calls deferTransitionTo. If no operation is in progress at that // point, the transition fires immediately, destroying the Consumer while we're still // inside its method and before the returned promise's .then() callbacks are set up. // The endOperation() happens in the .then() callbacks below, ensuring the deferred // state change only fires after the promise resolves/rejects and the Consumer's // this-capturing callbacks have already run. auto wrapDrainingRead = [this](jsg::Lock& js, jsg::Promise promise) -> jsg::Promise { return promise.then(js, [this](jsg::Lock& js, DrainingReadResult result) { if (state.endOperation()) { // A pending state was applied. Call the appropriate callback. if (state.template is()) { lock.onClose(js); } else if (state.template is()) { KJ_IF_SOME(err, state.template tryGetUnsafe()) { lock.onError(js, err.getHandle(js)); // The error was applied during this operation — the data we collected // may be invalid. Discard it and propagate the error rather than // silently returning possibly-corrupt data. js.throwException(err.addRef(js)); } } } return kj::mv(result); }, [this](jsg::Lock& js, jsg::Value exception) -> DrainingReadResult { state.clearPendingState(); (void)state.endOperation(); js.throwException(kj::mv(exception)); }); }; KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream not yet set up, treat as closed. return js.resolvedPromise(DrainingReadResult{ .chunks = kj::Array>(), .done = true, }); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return js.resolvedPromise(DrainingReadResult{ .chunks = kj::Array>(), .done = true, }); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.rejectedPromise(errored.addRef(js)); } KJ_CASE_ONEOF(consumer, kj::Own) { // beginOperation MUST be before consumer->drainingRead() — see comment above. state.beginOperation(); JSG_TRY(js) { return wrapDrainingRead(js, consumer->drainingRead(js, maxRead)); } JSG_CATCH(exception) { state.clearPendingState(); (void)state.endOperation(); doError(js, exception.getHandle(js)); return js.rejectedPromise(kj::mv(exception)); }; } KJ_CASE_ONEOF(consumer, kj::Own) { // beginOperation MUST be before consumer->drainingRead() — see comment above. state.beginOperation(); JSG_TRY(js) { return wrapDrainingRead(js, consumer->drainingRead(js, maxRead)); } JSG_CATCH(exception) { state.clearPendingState(); (void)state.endOperation(); doError(js, exception.getHandle(js)); return js.rejectedPromise(kj::mv(exception)); }; } } KJ_UNREACHABLE; } void ReadableStreamJsController::releaseReader(Reader& reader, kj::Maybe maybeJs) { lock.releaseReader(*this, reader, maybeJs); } ReadableStreamController::Tee ReadableStreamJsController::tee(jsg::Lock& js) { JSG_REQUIRE(!isLockedToReader(), TypeError, "This ReadableStream is locked to a reader."); lock.state.transitionTo(); disturbed = true; // This will leave this stream locked, disturbed, and closed. // Check for pending state first (deferred close/error during a prior read operation) if (state.pendingStateIs()) { return Tee{ .branch1 = js.alloc(kj::heap(StreamStates::Closed())), .branch2 = js.alloc(kj::heap(StreamStates::Closed())), }; } KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { return Tee{ .branch1 = js.alloc(kj::heap(pendingError.addRef(js))), .branch2 = js.alloc(kj::heap(pendingError.addRef(js))), }; } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream not yet set up, treat as closed. return Tee{ .branch1 = js.alloc(kj::heap(StreamStates::Closed())), .branch2 = js.alloc(kj::heap(StreamStates::Closed())), }; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return Tee{ .branch1 = js.alloc(kj::heap(StreamStates::Closed())), .branch2 = js.alloc(kj::heap(StreamStates::Closed())), }; } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return Tee{ .branch1 = js.alloc(kj::heap(errored.addRef(js))), .branch2 = js.alloc(kj::heap(errored.addRef(js))), }; } KJ_CASE_ONEOF(consumer, kj::Own) { KJ_DEFER(state.transitionTo()); // We create two additional streams that clone this stream's consumer state, // then close this stream's consumer. return Tee{ .branch1 = js.alloc(kj::heap(js, *consumer)), .branch2 = js.alloc(kj::heap(js, *consumer)), }; } KJ_CASE_ONEOF(consumer, kj::Own) { KJ_DEFER(state.transitionTo()); // We create two additional streams that clone this stream's consumer state, // then close this stream's consumer. return Tee{ .branch1 = js.alloc(kj::heap(js, *consumer)), .branch2 = js.alloc(kj::heap(js, *consumer)), }; } } KJ_UNREACHABLE; } void ReadableStreamJsController::setOwnerRef(ReadableStream& stream) { KJ_ASSERT(owner == kj::none); owner = &stream; } void ReadableStreamJsController::setup(jsg::Lock& js, jsg::Optional maybeUnderlyingSource, jsg::Optional maybeQueuingStrategy) { auto underlyingSource = kj::mv(maybeUnderlyingSource).orDefault({}); auto queuingStrategy = kj::mv(maybeQueuingStrategy).orDefault({}); auto type = underlyingSource.type.map([](kj::StringPtr s) { return s; }).orDefault(""_kj); expectedLength = underlyingSource.expectedLength; if (type == "bytes") { // Per spec, autoAllocateChunkSize should only be set if the user explicitly provides it. // If not set, the underlying source's pull method won't receive a byobRequest for // non-BYOB reads and must use controller.enqueue() instead. // // However, our original implementation always defaulted to 4096, so we need a compat flag // to control this behavior. Default to legacy behavior if flags aren't available. bool useSpecCompliantBehavior = false; KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) { useSpecCompliantBehavior = flags.getNoAutoAllocateChunkSize(); } kj::Maybe autoAllocateChunkSize; if (useSpecCompliantBehavior) { // Spec-compliant: only set if user explicitly provides it autoAllocateChunkSize = underlyingSource.autoAllocateChunkSize.map([](int size) { return size; }); } else { // Legacy behavior: apply a default autoAllocateChunkSize if not provided. auto defaultChunkSize = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE; if (util::Autogate::isEnabled(util::AutogateKey::UPDATED_AUTO_ALLOCATE_CHUNK_SIZE)) { defaultChunkSize = UnderlyingSource::DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2; } autoAllocateChunkSize = underlyingSource.autoAllocateChunkSize.orDefault(defaultChunkSize); } auto controller = js.alloc(kj::mv(underlyingSource), kj::mv(queuingStrategy)); KJ_IF_SOME(chunkSize, autoAllocateChunkSize) { JSG_REQUIRE(chunkSize > 0, TypeError, "The autoAllocateChunkSize option cannot be zero."); } // We account for the memory usage of the ByteReadable and its controller together because // their lifetimes are identical (in practice) and memory accounting itself has a memory // overhead. The same applies to ValueReadable below. state.transitionTo>( kj::heap(controller.addRef(), *this, autoAllocateChunkSize) .attach(js.getExternalMemoryAdjustment( sizeof(ByteReadable) + sizeof(ReadableByteStreamController)))); controller->start(js); } else { JSG_REQUIRE( type == "", TypeError, kj::str("\"", type, "\" is not a valid type of ReadableStream.")); auto controller = js.alloc( kj::mv(underlyingSource), kj::mv(queuingStrategy)); state.transitionTo>( kj::heap(controller.addRef(), *this) .attach(js.getExternalMemoryAdjustment( sizeof(ValueReadable) + sizeof(ReadableStreamDefaultController)))); controller->start(js); } } kj::Maybe ReadableStreamJsController::tryPipeLock() { return lock.tryPipeLock(*this); } void ReadableStreamJsController::visitForGc(jsg::GcVisitor& visitor) { // Visit pending state if it's an error (Closed has no GC-traceable content) KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { visitor.visit(pendingError); } // Note: We cannot use state.visitForGc(visitor) here because the state machine's // visitForGc passes kj::Own& to visitor.visit(), but GcVisitor expects T& for // types with visitForGc methods. We must dereference kj::Own manually. KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) {} KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(error, StreamStates::Errored) { visitor.visit(error); } KJ_CASE_ONEOF(consumer, kj::Own) { visitor.visit(*consumer); } KJ_CASE_ONEOF(consumer, kj::Own) { visitor.visit(*consumer); } } visitor.visit(lock); } kj::Maybe ReadableStreamJsController::getDesiredSize() { // If there's a pending state transition, return none if (state.hasPendingState()) { return kj::none; } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return kj::none; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return kj::none; } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return kj::none; } KJ_CASE_ONEOF(consumer, kj::Own) { return consumer->getDesiredSize(); } KJ_CASE_ONEOF(consumer, kj::Own) { return consumer->getDesiredSize(); } } KJ_UNREACHABLE; } kj::Maybe> ReadableStreamJsController::isErrored(jsg::Lock& js) { // Check for pending error first KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { return pendingError.getHandle(js); } // Pending Closed means not errored, so we can just check current state return state.tryGetUnsafe().map( [&](jsg::Value& reason) { return reason.getHandle(js); }); } bool ReadableStreamJsController::canCloseOrEnqueue() { // If there's a pending state transition, can't close or enqueue if (state.hasPendingState()) { return false; } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return false; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return false; } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return false; } KJ_CASE_ONEOF(consumer, kj::Own) { return consumer->canCloseOrEnqueue(); } KJ_CASE_ONEOF(consumer, kj::Own) { return consumer->canCloseOrEnqueue(); } } KJ_UNREACHABLE; } bool ReadableStreamJsController::hasBackpressure() { KJ_IF_SOME(size, getDesiredSize()) { return size <= 0; } return false; } kj::Maybe> ReadableStreamJsController:: getController() { // If there's a pending state transition, return none if (state.hasPendingState()) { return kj::none; } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return kj::none; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return kj::none; } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return kj::none; } KJ_CASE_ONEOF(consumer, kj::Own) { return consumer->getControllerRef(); } KJ_CASE_ONEOF(consumer, kj::Own) { return consumer->getControllerRef(); } } KJ_UNREACHABLE; } namespace { // Consumes all bytes from a stream, buffering in memory, with the purpose // of producing either a single concatenated kj::Array or kj::String. class AllReader { public: using PartList = kj::Array>; AllReader(jsg::Ref stream, uint64_t limit) : state(State::create>(kj::mv(stream))), limit(limit) {} KJ_DISALLOW_COPY_AND_MOVE(AllReader); jsg::Promise allBytes(jsg::Lock& js) { return loop(js).then(js, [this](auto& js, PartList&& partPtrs) -> jsg::BufferSource { auto out = jsg::BackingStore::alloc(js, runningTotal); copyInto(out.asArrayPtr(), partPtrs.asPtr()); return jsg::BufferSource(js, kj::mv(out)); }); } jsg::Promise allText( jsg::Lock& js, ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE) { return loop(js).then(js, [this, option](auto& js, PartList&& partPtrs) { // Strip UTF-8 BOM if requested if ((option & ReadAllTextOption::STRIP_BOM) && partPtrs.size() > 0 && hasUtf8Bom(partPtrs[0])) { partPtrs[0] = partPtrs[0].slice(UTF8_BOM_SIZE); runningTotal -= UTF8_BOM_SIZE; } JSG_REQUIRE(runningTotal <= v8::String::kMaxLength, RangeError, "String length exceeds v8::String::kMaxLength."); auto out = kj::heapArray(runningTotal + 1); copyInto(out.first(out.size() - 1).asBytes(), partPtrs.asPtr()); out.back() = '\0'; return kj::String(kj::mv(out)); }); } void visitForGc(jsg::GcVisitor& visitor) { state.visitForGc(visitor); } private: // State machine for AllReader: // Closed is terminal, Errored is implicitly terminal via ErrorState. // jsg::Ref is the active state (still reading). using State = StateMachine, ErrorState, ActiveState>, StreamStates::Closed, StreamStates::Errored, jsg::Ref>; State state; uint64_t limit; kj::Vector parts; uint64_t runningTotal = 0; jsg::Promise loop(jsg::Lock& js) { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(closed, StreamStates::Closed) { return js.resolvedPromise(KJ_MAP(p, parts) { return p.asArrayPtr(); }); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.template rejectedPromise(errored.getHandle(js)); } KJ_CASE_ONEOF(readable, jsg::Ref) { // Note that these nested lambda retain references to `this` and `readable` // and are passed into to promise returned by this method. It is the responsibility // of the caller to ensure that the AllReader instance is kept alive until the // promise is settled. auto onSuccess = JSG_VISITABLE_LAMBDA((this, readable = readable.addRef()), (readable), (jsg::Lock & js, ReadResult result) mutable->jsg::Promise { if (result.done) { state.template transitionTo(); return loop(js); } // If we're not done, the result value must be interpretable as // bytes for the read to make any sense. auto handle = KJ_ASSERT_NONNULL(result.value).getHandle(js); if (!handle->IsArrayBufferView() && !handle->IsArrayBuffer()) { auto error = js.v8TypeError("This ReadableStream did not return bytes."); state.template transitionTo(js.v8Ref(error)); return readable->getController().cancel(js, error).then( js, [&](jsg::Lock& js) { return loop(js); }); } jsg::BufferSource bufferSource(js, handle); if (bufferSource.size() == 0) { // Weird but allowed, we'll skip it. return loop(js); } if ((runningTotal + bufferSource.size()) > limit) { auto error = js.v8TypeError("Memory limit exceeded before EOF."); state.template transitionTo(js.v8Ref(error)); return readable->getController().cancel(js, error).then( js, [&](jsg::Lock& js) { return loop(js); }); } runningTotal += bufferSource.size(); parts.add(bufferSource.copy(js)); return loop(js); }); auto onFailure = [this](auto& js, jsg::Value exception) -> jsg::Promise { // In this case the stream should already be errored. state.template transitionTo(js.v8Ref(exception.getHandle(js))); return loop(js); }; return maybeAddFunctor(js, KJ_ASSERT_NONNULL(readable->getController().read(js, kj::none)), kj::mv(onSuccess), kj::mv(onFailure)); } } KJ_UNREACHABLE; } void copyInto(kj::ArrayPtr out, kj::ArrayPtr> in) { for (auto& part: in) { KJ_ASSERT(part.size() <= out.size()); out.first(part.size()).copyFrom(part); out = out.slice(part.size()); } } }; // PumpToReader implements the original JS promise-loop approach to pumping data from // a ReadableStream to a WritableStreamSink. It reads one chunk at a time using the // standard read() API, writes each chunk to the sink, and loops until done or errored. // This is the fallback path used when the ENABLE_DRAINING_READ_ON_STANDARD_STREAMS // autogate is not enabled. class PumpToReader { public: PumpToReader(jsg::Ref stream, kj::Own sink, bool end) : ioContext(IoContext::current()), state(State::create>(kj::mv(stream))), sink(kj::mv(sink)), self(kj::refcounted>(kj::Badge{}, *this)), end(end) {} KJ_DISALLOW_COPY_AND_MOVE(PumpToReader); ~PumpToReader() noexcept(false) { self->invalidate(); // Ensure that if a write promise is pending it is proactively canceled. canceler.cancel("PumpToReader was destroyed"); } kj::Promise pumpTo(jsg::Lock& js) { ioContext.requireCurrentOrThrowJs(); KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(stream, jsg::Ref) { auto readable = stream.addRef(); state.template transitionTo(); return ioContext.awaitJs( js, pumpLoop(js, ioContext, kj::mv(readable), ioContext.addObject(self->addRef()))); } KJ_CASE_ONEOF(pumping, Pumping) { return KJ_EXCEPTION(FAILED, "pumping is already in progress"); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return KJ_EXCEPTION(FAILED, "stream has already been consumed"); } KJ_CASE_ONEOF(errored, kj::Exception) { return errored.clone(); } } KJ_UNREACHABLE; } private: struct Pumping { static constexpr kj::StringPtr NAME KJ_UNUSED = "pumping"_kj; }; IoContext& ioContext; using State = StateMachine, ErrorState, Pumping, StreamStates::Closed, kj::Exception, jsg::Ref>; State state; kj::Own sink; kj::Own> self; kj::Canceler canceler; bool end; bool isErroredOrClosed() { return state.isTerminal(); } jsg::Promise pumpLoop(jsg::Lock& js, IoContext& ioContext, jsg::Ref readable, IoOwn> pumpToReader) { ioContext.requireCurrentOrThrowJs(); KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(ready, jsg::Ref) { KJ_UNREACHABLE; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return end ? ioContext.awaitIoLegacy(js, sink->end().attach(kj::mv(sink))) : js.resolvedPromise(); } KJ_CASE_ONEOF(errored, kj::Exception) { if (end) { sink->abort(errored.clone()); } return js.rejectedPromise(errored.clone()); } KJ_CASE_ONEOF(pumping, Pumping) { using Result = kj::OneOf, StreamStates::Closed, jsg::Value>; return KJ_ASSERT_NONNULL(readable->getController().read(js, kj::none)) .then(js, ioContext.addFunctor([byteStream = readable->getController().isByteOriented()]( auto& js, ReadResult result) mutable -> Result { if (result.done) { return StreamStates::Closed(); } auto handle = KJ_ASSERT_NONNULL(result.value).getHandle(js); if (!handle->IsArrayBufferView() && !handle->IsArrayBuffer()) { return js.v8Ref(js.v8TypeError("This ReadableStream did not return bytes.")); } jsg::BufferSource bufferSource(js, handle); if (bufferSource.size() == 0) { return Pumping{}; } if (byteStream) { jsg::BackingStore backing = bufferSource.detach(js); return backing.asArrayPtr().attach(kj::mv(backing)); } return bufferSource.asArrayPtr().attach(kj::mv(bufferSource)); }), [](auto& js, jsg::Value exception) mutable -> Result { return kj::mv(exception); }) .then(js, ioContext.addFunctor( JSG_VISITABLE_LAMBDA((readable = kj::mv(readable), pumpToReader = kj::mv(pumpToReader)), (readable), (jsg::Lock & js, Result result) mutable { KJ_IF_SOME(reader, pumpToReader->tryGet()) { reader.ioContext.requireCurrentOrThrowJs(); auto& ioContext = IoContext::current(); KJ_SWITCH_ONEOF(result) { KJ_CASE_ONEOF(bytes, kj::Array) { auto promise = reader.sink->write(bytes).attach(kj::mv(bytes)); return ioContext.awaitIo(js, reader.canceler.wrap(kj::mv(promise))) .then(js, [](jsg::Lock& js) -> kj::Maybe { return kj::Maybe(kj::none); }, [](jsg::Lock& js, jsg::Value exception) mutable -> kj::Maybe { return kj::mv(exception); }) .then(js, ioContext.addFunctor(JSG_VISITABLE_LAMBDA( (readable = readable.addRef(), pumpToReader = kj::mv(pumpToReader)), (readable), (jsg::Lock & js, kj::Maybe maybeException) mutable { KJ_IF_SOME(reader, pumpToReader->tryGet()) { auto& ioContext = reader.ioContext; ioContext.requireCurrentOrThrowJs(); KJ_IF_SOME(exception, maybeException) { if (!reader.isErroredOrClosed()) { reader.state.transitionTo( js.exceptionToKj(kj::mv(exception))); } } else { // Else block to avert dangling else compiler warning. } return reader.pumpLoop( js, ioContext, readable.addRef(), kj::mv(pumpToReader)); } else { return readable->getController().cancel(js, maybeException.map( [&](jsg::Value& ex) { return ex.getHandle(js); })); } }))); } KJ_CASE_ONEOF(pumping, Pumping) {} KJ_CASE_ONEOF(closed, StreamStates::Closed) { if (!reader.isErroredOrClosed()) { reader.state.transitionTo(); } } KJ_CASE_ONEOF(exception, jsg::Value) { if (!reader.isErroredOrClosed()) { reader.state.transitionTo(js.exceptionToKj(kj::mv(exception))); } } } return reader.pumpLoop(js, ioContext, readable.addRef(), kj::mv(pumpToReader)); } else { KJ_SWITCH_ONEOF(result) { KJ_CASE_ONEOF(bytes, kj::Array) { return readable->getController().cancel(js, kj::none); } KJ_CASE_ONEOF(pumping, Pumping) { return readable->getController().cancel(js, kj::none); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return js.resolvedPromise(); } KJ_CASE_ONEOF(exception, jsg::Value) { return readable->getController().cancel(js, exception.getHandle(js)); } } } KJ_UNREACHABLE; }))); } } KJ_UNREACHABLE; } }; // pumpToCoroutine uses a DrainingReader to efficiently pull all synchronously available // data from the stream in each iteration, then writes it to the sink using vectored // I/O. This minimizes isolate lock acquisitions by batching: each time the lock is // held, the stream's internal queue is fully drained and the JS pull callback is // pumped synchronously as many times as possible. // // The pump loop is a kj coroutine. Dropping the returned kj::Promise drops the // coroutine frame, which destroys the DrainingReader (releasing the stream lock) // and the sink. No WeakRef/IoOwn dance is needed because ownership is clear. // The coroutine that implements the pump loop takes ownership of the DrainingReader // and sink. The jsg::Ref is not passed into the coroutine because // jsg::Ref is disallowed in coroutine parameters; instead, the DrainingReader holds // a reference to the stream internally. kj::Promise pumpToImpl(IoContext& ioContext, kj::Own reader, kj::Own sink, bool end) { bool writeFailed = false; KJ_TRY { while (true) { // Perform a draining read to get all synchronously available data if possible // or fall back to a regular read if not. DrainingReadResult result = co_await ioContext.run([&reader](jsg::Lock& js) mutable { auto& ioContext = IoContext::current(); // Use a 256KB limit to allow periodic yielding to the event loop, // preventing a fast producer from monopolizing the thread. constexpr size_t kMaxReadPerCycle = 256 * 1024; return ioContext.awaitJs(js, reader->read(js, kMaxReadPerCycle)); }); // Write all the chunks we received using vectored write for efficiency. if (result.chunks.size() > 0) { KJ_ON_SCOPE_FAILURE(writeFailed = true); auto pieces = KJ_MAP(chunk, result.chunks) -> kj::ArrayPtr { return chunk.asPtr(); }; co_await sink->write(pieces); } // If the stream is done, end the output if needed and exit. if (result.done) { KJ_ON_SCOPE_FAILURE(writeFailed = true); if (end) { co_await sink->end(); } co_return; } } } KJ_CATCH(exception) { if (!writeFailed) { sink->abort(exception.clone()); } co_await ioContext.run([&reader, ex = exception.clone()](jsg::Lock& js) mutable { auto& ioContext = IoContext::current(); auto error = js.exceptionToJsValue(kj::mv(ex)); return ioContext.awaitJs(js, reader->cancel(js, error.getHandle(js))); }); kj::throwFatalException(kj::mv(exception)); } } } // namespace template jsg::Promise ReadableStreamJsController::readAll(jsg::Lock& js, uint64_t limit) { if (isLockedToReader()) { return js.rejectedPromise(KJ_EXCEPTION( FAILED, "jsg.TypeError: This ReadableStream is currently locked to a reader.")); } disturbed = true; bool stripBom = false; KJ_IF_SOME(flags, FeatureFlags::tryGet(js)) { stripBom = flags.getStripBomInReadAllText(); } // This operation leaves the stream locked and disturbed. The loop will read until // the stream is closed or errored. If the limit is reached, the loop will error. const auto readAll = [this, limit, stripBom](auto& js) -> jsg::Promise { KJ_ASSERT(lock.lock()); // The AllReader will hold a traceable reference to the ReadableStream. auto reader = kj::heap(addRef(), limit); auto promise = ([&js, &reader, stripBom]() -> jsg::Promise { if constexpr (kj::isSameType()) { (void)stripBom; // Unused in this branch. return reader->allBytes(js); } else { auto option = ReadAllTextOption::NULL_TERMINATE; if (stripBom) { option |= ReadAllTextOption::STRIP_BOM; } return reader->allText(js, option); } })(); return maybeAddFunctor(js, kj::mv(promise), // reader is a GC visitable type that holds a reference to either the stream // or an error. Accordingly, we wrap it in a visitable lambda attached as a // continuation on the promise to ensure that it is GC visited and kept alive until // the promise settles. JSG_VISITABLE_LAMBDA((reader = kj::mv(reader)), (reader), (jsg::Lock & js, T result)->jsg::Promise { return js.resolvedPromise(kj::mv(result)); }), [](jsg::Lock& js, jsg::Value exception) -> jsg::Promise { return js.rejectedPromise(kj::mv(exception)); }); }; KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream not yet set up, treat as closed. if constexpr (kj::isSameType()) { auto backing = jsg::BackingStore::alloc(js, 0); return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); } else { return js.resolvedPromise(T()); } } KJ_CASE_ONEOF(closed, StreamStates::Closed) { if constexpr (kj::isSameType()) { auto backing = jsg::BackingStore::alloc(js, 0); return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); } else { return js.resolvedPromise(T()); } } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.rejectedPromise(errored.addRef(js)); } KJ_CASE_ONEOF(valueReadable, kj::Own) { return readAll(js); } KJ_CASE_ONEOF(byteReadable, kj::Own) { return readAll(js); } } KJ_UNREACHABLE; } jsg::Promise ReadableStreamJsController::readAllBytes( jsg::Lock& js, uint64_t limit) { return readAll(js, limit); } jsg::Promise ReadableStreamJsController::readAllText(jsg::Lock& js, uint64_t limit) { return readAll(js, limit); } kj::Own ReadableStreamJsController::detach( jsg::Lock& js, bool ignored /* unused */) { KJ_ASSERT(!isLockedToReader()); KJ_ASSERT(!isDisturbed()); KJ_ASSERT(!state.hasOperationInProgress(), "Unable to detach with read pending"); auto controller = kj::heap(); controller->expectedLength = expectedLength; disturbed = true; // Clones this streams state into a new ReadableStreamController, leaving this stream // locked, disturbed, and closed. // The controller starts in Initial state by default, so we can use regular transitionTo. KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Still in initial state, transition to closed controller->state.transitionTo(); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { controller->state.transitionTo(); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { controller->state.transitionTo(errored.addRef(js)); } KJ_CASE_ONEOF(readable, kj::Own) { KJ_ASSERT(lock.lock()); controller->state.transitionTo>(readable->clone(js, *controller)); state.transitionTo(); lock.onClose(js); } KJ_CASE_ONEOF(readable, kj::Own) { KJ_ASSERT(lock.lock()); controller->state.transitionTo>(readable->clone(js, *controller)); state.transitionTo(); lock.onClose(js); } } return kj::mv(controller); } kj::Maybe ReadableStreamJsController::tryGetLength(StreamEncoding encoding) { return expectedLength; } kj::Promise> ReadableStreamJsController::pumpTo( jsg::Lock& js, kj::Own sink, bool end) { KJ_ASSERT(IoContext::hasCurrent(), "Unable to consume this ReadableStream outside of a request"); KJ_REQUIRE(!isLockedToReader(), "This ReadableStream is currently locked to a reader."); disturbed = true; // This operation will leave the ReadableStream locked and disturbed. It will consume // the stream until it either closed or errors. // // When the ENABLE_DRAINING_READ_ON_STANDARD_STREAMS autogate is enabled, uses the new // pumpToImpl coroutine with DrainingReader for batched reads and vectored writes. // Otherwise, falls back to the original PumpToReader JS promise loop that reads one // chunk at a time. const auto handlePump = [&] { if (util::Autogate::isEnabled(util::AutogateKey::ENABLE_DRAINING_READ_ON_STANDARD_STREAMS)) { auto reader = KJ_ASSERT_NONNULL(DrainingReader::create(js, *this->addRef()), "Failed to create DrainingReader — stream should not be locked"); auto& ioContext = IoContext::current(); return addNoopDeferredProxy(pumpToImpl(ioContext, kj::mv(reader), kj::mv(sink), end)); } else { KJ_ASSERT(lock.lock()); auto reader = kj::heap(addRef(), kj::mv(sink), end); return addNoopDeferredProxy(reader->pumpTo(js).attach(kj::mv(reader))); } }; KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream not yet set up, treat as closed. return addNoopDeferredProxy(sink->end().attach(kj::mv(sink))); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return addNoopDeferredProxy(sink->end().attach(kj::mv(sink))); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.exceptionToKj(errored.addRef(js)); } KJ_CASE_ONEOF(readable, kj::Own) { return handlePump(); } KJ_CASE_ONEOF(readable, kj::Own) { return handlePump(); } } KJ_UNREACHABLE; } // ====================================================================================== WritableStreamDefaultController::WritableStreamDefaultController( jsg::Lock& js, WritableStream& owner, jsg::Ref abortSignal) : ioContext(tryGetIoContext()), impl(js, owner, kj::mv(abortSignal)) {} jsg::Promise WritableStreamDefaultController::abort( jsg::Lock& js, v8::Local reason) { return impl.abort(js, JSG_THIS, reason); } void WritableStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) { visitor.visit(impl); } jsg::Promise WritableStreamDefaultController::close(jsg::Lock& js) { return impl.close(js, JSG_THIS); } void WritableStreamDefaultController::error( jsg::Lock& js, jsg::Optional> reason) { impl.error(js, JSG_THIS, reason.orDefault(js.undefined())); } kj::Maybe WritableStreamDefaultController::getDesiredSize() { // Per the spec, desiredSize should be null when the stream is erroring. if (impl.flags.pedanticWpt && isErroring()) { return kj::none; } return impl.getDesiredSize(); } jsg::Ref WritableStreamDefaultController::getSignal() { return impl.signal.addRef(); } kj::Maybe> WritableStreamDefaultController::isErroring(jsg::Lock& js) { KJ_IF_SOME(erroring, impl.state.tryGetUnsafe()) { return erroring.reason.getHandle(js); } return kj::none; } void WritableStreamDefaultController::setup( jsg::Lock& js, UnderlyingSink underlyingSink, StreamQueuingStrategy queuingStrategy) { impl.setup(js, JSG_THIS, kj::mv(underlyingSink), kj::mv(queuingStrategy)); } jsg::Promise WritableStreamDefaultController::write( jsg::Lock& js, v8::Local value) { return impl.write(js, JSG_THIS, value); } void WritableStreamDefaultController::cancelPendingWrites(jsg::Lock& js, jsg::JsValue reason) { impl.cancelPendingWrites(js, reason); } void WritableStreamDefaultController::clearAlgorithms() { impl.algorithms.clear(); } WritableStreamDefaultController::~WritableStreamDefaultController() noexcept(false) { // Clear algorithms in destructor to break circular references clearAlgorithms(); } // ====================================================================================== WritableStreamJsController::WritableStreamJsController(): ioContext(tryGetIoContext()) {} WritableStreamJsController::~WritableStreamJsController() noexcept(false) { // Clear algorithms to break circular references during destruction KJ_IF_SOME(controller, state.tryGetUnsafe()) { controller->clearAlgorithms(); } // Clear the state to break the circular reference to the controller. // During destruction, we force the transition since the current state doesn't matter. state.forceTransitionTo(); // Clear owner reference owner = kj::none; // Clear any pending abort promise maybeAbortPromise = kj::none; } WritableStreamJsController::WritableStreamJsController(StreamStates::Closed closed) : ioContext(tryGetIoContext()) { state.transitionTo(); } WritableStreamJsController::WritableStreamJsController(StreamStates::Errored errored) : ioContext(tryGetIoContext()) { state.transitionTo(kj::mv(errored)); } jsg::Promise WritableStreamJsController::abort( jsg::Lock& js, jsg::Optional> reason) { // The spec requires that if abort is called multiple times, it is supposed to return the same // promise each time. That's a bit cumbersome here with jsg::Promise so we intentionally just // return a continuation branch off the same promise. KJ_IF_SOME(abortPromise, maybeAbortPromise) { return abortPromise.whenResolved(js); } KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { // Stream hasn't been set up yet - treat like closed for abort purposes maybeAbortPromise = js.resolvedPromise(); return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { maybeAbortPromise = js.resolvedPromise(); return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { // Per the spec, if the stream is errored, we are to return a resolved promise. maybeAbortPromise = js.resolvedPromise(); return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); } KJ_CASE_ONEOF(controller, Controller) { maybeAbortPromise = controller->abort(js, reason.orDefault(js.undefined())); return KJ_ASSERT_NONNULL(maybeAbortPromise).whenResolved(js); } } KJ_UNREACHABLE; } jsg::Ref WritableStreamJsController::addRef() { return KJ_ASSERT_NONNULL(owner).addRef(); } bool WritableStreamJsController::isClosedOrClosing() { return state.is(); } bool WritableStreamJsController::isErrored() { return state.isErrored(); } jsg::Promise WritableStreamJsController::close(jsg::Lock& js, bool markAsHandled) { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return rejectedMaybeHandledPromise( js, js.v8TypeError("This WritableStream has been closed."_kj), markAsHandled); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return rejectedMaybeHandledPromise( js, js.v8TypeError("This WritableStream has been closed."_kj), markAsHandled); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { if (FeatureFlags::get(js).getPedanticWpt()) { return rejectedMaybeHandledPromise( js, js.v8TypeError("This WritableStream has been errored."_kj), markAsHandled); } return rejectedMaybeHandledPromise(js, errored.getHandle(js), markAsHandled); } KJ_CASE_ONEOF(controller, Controller) { return controller->close(js); } } KJ_UNREACHABLE; } void WritableStreamJsController::doClose(jsg::Lock& js) { // If already in a terminal state, nothing to do. if (state.isTerminal()) return; // Clear algorithms to break circular references before changing state KJ_IF_SOME(controller, state.tryGetUnsafe()) { controller->clearAlgorithms(); } state.transitionTo(); KJ_IF_SOME(locked, lock.state.tryGetUnsafe()) { maybeResolvePromise(js, locked.getClosedFulfiller()); maybeResolvePromise(js, locked.getReadyFulfiller()); } else { (void)lock.state.transitionFromTo(); } } void WritableStreamJsController::doError(jsg::Lock& js, v8::Local reason) { // If already in a terminal state, nothing to do. if (state.isTerminal()) return; // Clear algorithms to break circular references before changing state KJ_IF_SOME(controller, state.tryGetUnsafe()) { controller->clearAlgorithms(); } state.transitionTo(js.v8Ref(reason)); KJ_IF_SOME(locked, lock.state.tryGetUnsafe()) { maybeRejectPromise(js, locked.getClosedFulfiller(), reason); maybeResolvePromise(js, locked.getReadyFulfiller()); } else KJ_IF_SOME(pipeLocked, lock.state.tryGetUnsafe()) { // When the writable side of a pipe errors, we need to release the source stream. // The pipeLoop may be waiting on a read from the source that will never complete, // so we need to proactively release the source here. if (!pipeLocked.flags.preventCancel) { pipeLocked.source.release(js, reason); } else { pipeLocked.source.release(js); } lock.state.transitionTo(); } } void WritableStreamJsController::errorIfNeeded(jsg::Lock& js, v8::Local reason) { // Error through the underlying controller if available, which goes through the proper // error transition (Erroring -> Errored). This allows close() to be called while the // stream is "erroring" and reject with the stored error. KJ_IF_SOME(controller, state.tryGetUnsafe()) { controller->error(js, reason); } // If state is not Controller (already Closed or Errored), this is a no-op. } kj::Maybe WritableStreamJsController::getDesiredSize() { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return 0; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return 0; } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return kj::none; } KJ_CASE_ONEOF(controller, Controller) { return controller->getDesiredSize().map([](ssize_t size) -> int { return size; }); } } KJ_UNREACHABLE; } kj::Maybe> WritableStreamJsController::isErroring(jsg::Lock& js) { KJ_IF_SOME(controller, state.tryGetUnsafe()) { return controller->isErroring(js); } return kj::none; } bool WritableStreamDefaultController::isErroring() const { return impl.state.is(); } kj::Maybe> WritableStreamJsController::isErroredOrErroring(jsg::Lock& js) { KJ_IF_SOME(err, state.tryGetErrorUnsafe()) { return err.getHandle(js); } return isErroring(js); } bool WritableStreamJsController::isStarted() { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return false; } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return true; } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return true; } KJ_CASE_ONEOF(controller, Controller) { return controller->isStarted(); } } KJ_UNREACHABLE; } bool WritableStreamJsController::hasBackpressure() { KJ_IF_SOME(controller, state.tryGetUnsafe()) { return controller->hasBackpressure(); } return false; } bool WritableStreamJsController::isLocked() const { return isLockedToWriter(); } bool WritableStreamJsController::isLockedToWriter() const { return !lock.state.is(); } bool WritableStreamJsController::lockWriter(jsg::Lock& js, Writer& writer) { return lock.lockWriter(js, *this, writer); } void WritableStreamJsController::maybeRejectReadyPromise( jsg::Lock& js, v8::Local reason) { KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe()) { if (writerLock.getReadyFulfiller() != kj::none) { maybeRejectPromise(js, writerLock.getReadyFulfiller(), reason); } else { auto prp = js.newPromiseAndResolver(); prp.promise.markAsHandled(js); prp.resolver.reject(js, reason); writerLock.setReadyFulfiller(js, prp); } } } void WritableStreamJsController::maybeResolveReadyPromise(jsg::Lock& js) { KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe()) { maybeResolvePromise(js, writerLock.getReadyFulfiller()); } } void WritableStreamJsController::releaseWriter(Writer& writer, kj::Maybe maybeJs) { lock.releaseWriter(*this, writer, maybeJs); } kj::Maybe> WritableStreamJsController::removeSink(jsg::Lock& js) { return kj::none; } void WritableStreamJsController::detach(jsg::Lock& js) { KJ_UNIMPLEMENTED("WritableStreamJsController::detach is not implemented"); } void WritableStreamJsController::setOwnerRef(WritableStream& stream) { owner = stream; } void WritableStreamJsController::setup(jsg::Lock& js, jsg::Optional maybeUnderlyingSink, jsg::Optional maybeQueuingStrategy) { auto underlyingSink = kj::mv(maybeUnderlyingSink).orDefault({}); auto queuingStrategy = kj::mv(maybeQueuingStrategy).orDefault({}); if (FeatureFlags::get(js).getPedanticWpt()) { // Per the spec, the type property for WritableStream's underlying sink must be undefined. // If it's anything else, throw a RangeError. JSG_REQUIRE(underlyingSink.type == kj::none, RangeError, "Invalid underlying sink type. Only undefined is valid."); } // We account for the memory usage of the WritableStreamDefaultController and AbortSignal together // because their lifetimes are identical and memory accounting itself has a memory overhead. auto controller = js.allocAccounted( sizeof(WritableStreamDefaultController) + sizeof(AbortSignal), js, KJ_ASSERT_NONNULL(owner), js.alloc()); auto& controllerRef = *controller; state.transitionTo(kj::mv(controller)); controllerRef.setup(js, kj::mv(underlyingSink), kj::mv(queuingStrategy)); } kj::Maybe> WritableStreamJsController::tryPipeFrom( jsg::Lock& js, jsg::Ref source, PipeToOptions options) { JSG_REQUIRE_NONNULL( ioContext, Error, "Unable to pipe to a WritableStream created outside of a request"); // The ReadableStream source here can be either a JavaScript-backed ReadableStream // or ReadableStreamSource-backed. In either case, however, this WritableStream is // JavaScript-based and must use a JavaScript promise-based data flow for piping data. // We'll treat all ReadableStreams as if they are JavaScript-backed. // // This method will return a JavaScript promise that is resolved when the pipe operation // completes, or is rejected if the pipe operation is aborted or errored. // Let's also acquire the destination pipe lock. lock.pipeLock(KJ_ASSERT_NONNULL(owner), kj::mv(source), options); return pipeLoop(js).then(js, JSG_VISITABLE_LAMBDA((ref = addRef()), (ref), (auto& js){})); } jsg::Promise WritableStreamJsController::pipeLoop(jsg::Lock& js) { auto maybePipeLock = lock.tryGetPipe(); if (maybePipeLock == kj::none) return js.resolvedPromise(); auto& pipeLock = KJ_REQUIRE_NONNULL(maybePipeLock); auto preventAbort = pipeLock.flags.preventAbort; auto preventCancel = pipeLock.flags.preventCancel; auto preventClose = pipeLock.flags.preventClose; auto pipeThrough = pipeLock.flags.pipeThrough; auto& source = pipeLock.source; // At the start of each pipe step, we check to see if either the source or // the destination has closed or errored and propagate that on to the other. KJ_IF_SOME(promise, pipeLock.checkSignal(js, *this)) { lock.releasePipeLock(); return kj::mv(promise); } KJ_IF_SOME(errored, pipeLock.source.tryGetErrored(js)) { source.release(js); lock.releasePipeLock(); if (!preventAbort) { auto onSuccess = JSG_VISITABLE_LAMBDA( (pipeThrough, reason = js.v8Ref(errored)), (reason), (jsg::Lock& js) { return rejectedMaybeHandledPromise(js, reason.getHandle(js), pipeThrough); }); auto promise = abort(js, errored); KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { return promise.then(js, ioContext.addFunctor(kj::mv(onSuccess))); } else { return promise.then(js, kj::mv(onSuccess)); } } return rejectedMaybeHandledPromise(js, errored, pipeThrough); } KJ_IF_SOME(errored, state.tryGetUnsafe()) { lock.releasePipeLock(); auto reason = errored.getHandle(js); if (!preventCancel) { source.release(js, reason); } else { source.release(js); } return rejectedMaybeHandledPromise(js, reason, pipeThrough); } KJ_IF_SOME(erroring, isErroring(js)) { lock.releasePipeLock(); if (!preventCancel) { source.release(js, erroring); } else { source.release(js); } return rejectedMaybeHandledPromise(js, erroring, pipeThrough); } if (source.isClosed()) { source.release(js); lock.releasePipeLock(); if (!preventClose) { auto promise = close(js); if (pipeThrough) { promise.markAsHandled(js); } return kj::mv(promise); } return js.resolvedPromise(); } if (state.is()) { lock.releasePipeLock(); auto reason = js.v8TypeError("This destination writable stream is closed."_kj); if (!preventCancel) { source.release(js, reason); } else { source.release(js); } return rejectedMaybeHandledPromise(js, reason, pipeThrough); } // Assuming we get by that, we perform a read on the source. If the read errors, // we propagate the error to the destination, depending on options and reject // the pipe promise. If the read is successful then we'll get a ReadResult // back. If the ReadResult indicates done, then we close the destination // depending on options and resolve the pipe promise. If the ReadResult is // not done, we write the value on to the destination. If the write operation // fails, we reject the pipe promise and propagate the error back to the // source (again, depending on options). If the write operation is successful, // we call pipeLoop again to move on to the next iteration. auto onSuccess = JSG_VISITABLE_LAMBDA((this, ref = addRef(), preventCancel, pipeThrough), (ref), (jsg::Lock & js, ReadResult result)->jsg::Promise { auto maybePipeLock = lock.tryGetPipe(); if (maybePipeLock == kj::none) return js.resolvedPromise(); auto& pipeLock = KJ_REQUIRE_NONNULL(maybePipeLock); KJ_IF_SOME(promise, pipeLock.checkSignal(js, *this)) { lock.releasePipeLock(); return kj::mv(promise); } else { } // Trailing else() is squash compiler warning if (result.done) { // We'll handle the close at the start of the next iteration. return pipeLoop(js); } auto onSuccess = JSG_VISITABLE_LAMBDA( (this, ref=addRef()), (ref) , (jsg::Lock& js) { return pipeLoop(js); } ); auto onFailure = JSG_VISITABLE_LAMBDA( (this, ref=addRef(), preventCancel, pipeThrough), (ref) , (jsg::Lock& js, jsg::Value value) { // The write failed. We need to release the source if the pipe lock still exists. auto reason = value.getHandle(js); KJ_IF_SOME(pipeLock, lock.tryGetPipe()) { if (!preventCancel) { pipeLock.source.release(js, reason); } else { pipeLock.source.release(js); } } else {} // Trailing else() to squash compiler warning return rejectedMaybeHandledPromise(js, reason, pipeThrough); } ); auto promise = write(js, result.value.map([&](jsg::Value& value) { return value.getHandle(js); })); return maybeAddFunctor(js, kj::mv(promise), kj::mv(onSuccess), kj::mv(onFailure)); }); auto onFailure = JSG_VISITABLE_LAMBDA((this, ref = addRef()), (ref), (jsg::Lock& js, jsg::Value value) { // The read failed. We will handle the error at the start of the next iteration. return pipeLoop(js); }); return maybeAddFunctor(js, pipeLock.source.read(js), kj::mv(onSuccess), kj::mv(onFailure)); } void WritableStreamJsController::updateBackpressure(jsg::Lock& js, bool backpressure) { KJ_IF_SOME(writerLock, lock.state.tryGetUnsafe()) { if (backpressure) { // Per the spec, when backpressure is updated and is true, we replace the existing // ready promise on the writer with a new pending promise, regardless of whether // the existing one is resolved or not. auto prp = js.newPromiseAndResolver(); prp.promise.markAsHandled(js); return writerLock.setReadyFulfiller(js, prp); } // When backpressure is updated and is false, we resolve the ready promise on the writer maybeResolvePromise(js, writerLock.getReadyFulfiller()); } } jsg::Promise WritableStreamJsController::write( jsg::Lock& js, jsg::Optional> value) { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) { return js.rejectedPromise(js.v8TypeError("This WritableStream has been closed."_kj)); } KJ_CASE_ONEOF(closed, StreamStates::Closed) { return js.rejectedPromise(js.v8TypeError("This WritableStream has been closed."_kj)); } KJ_CASE_ONEOF(errored, StreamStates::Errored) { return js.rejectedPromise(errored.addRef(js)); } KJ_CASE_ONEOF(controller, Controller) { return controller->write(js, value.orDefault([&] { return js.undefined(); })); } } KJ_UNREACHABLE; } void WritableStreamJsController::visitForGc(jsg::GcVisitor& visitor) { state.visitForGc(visitor); visitor.visit(maybeAbortPromise, lock); } // ======================================================================================= TransformStreamDefaultController::TransformStreamDefaultController(jsg::Lock& js) : ioContext(tryGetIoContext()), startPromise(js.newPromiseAndResolver()) {} kj::Maybe TransformStreamDefaultController::getDesiredSize() { KJ_IF_SOME(readableController, tryGetReadableController()) { return readableController.getDesiredSize(); } return kj::none; } void TransformStreamDefaultController::enqueue(jsg::Lock& js, v8::Local chunk) { auto& readableController = JSG_REQUIRE_NONNULL(tryGetReadableController(), TypeError, "The readable side of this TransformStream is no longer readable."); // Hold a strong reference to the readable controller for the duration of this // method. The readableController.enqueue() call below invokes the user-provided // size algorithm, which can re-enter JS and call error() on this transform // controller, dropping the jsg::Ref held by this->readable and the one held by // the ReadableStreamJsController's ValueReadable. Without this ref the // ReadableStreamDefaultController would be freed while its enqueue() method is // still on the stack. auto readableControllerRef = kj::addRef(readableController); JSG_REQUIRE(readableController.canCloseOrEnqueue(), TypeError, "The readable side of this TransformStream is no longer readable."); js.tryCatch([&] { readableController.enqueue(js, chunk); }, [&](jsg::Value exception) { errorWritableAndUnblockWrite(js, exception.getHandle(js)); js.throwException(kj::mv(exception)); }); // If the controller was errored during the enqueue (e.g. by the size callback // calling error()), skip the backpressure update — the stream is already torn down. if (!readableController.canCloseOrEnqueue()) { return; } bool newBackpressure = readableController.hasBackpressure(); if (newBackpressure != backpressure) { KJ_ASSERT(newBackpressure); // Unfortunately the original implementation forgot to actually set the backpressure // here so the backpressure signaling failed to work correctly. This is unfortunate // because applying the backpressure here could break existing code, so we need to // put the fix behind a compat flag. Doh! if (FeatureFlags::get(js).getFixupTransformStreamBackpressure()) { setBackpressure(js, true); } } } void TransformStreamDefaultController::error(jsg::Lock& js, v8::Local reason) { KJ_IF_SOME(readableController, tryGetReadableController()) { readableController.error(js, reason); readable = kj::none; } errorWritableAndUnblockWrite(js, reason); } void TransformStreamDefaultController::terminate(jsg::Lock& js) { KJ_IF_SOME(readableController, tryGetReadableController()) { readableController.close(js); readable = kj::none; } errorWritableAndUnblockWrite(js, js.v8TypeError("The transform stream has been terminated"_kj)); } jsg::Promise TransformStreamDefaultController::write( jsg::Lock& js, v8::Local chunk) { KJ_IF_SOME(writableController, tryGetWritableController()) { KJ_IF_SOME(error, writableController.isErroredOrErroring(js)) { return js.rejectedPromise(error); } KJ_ASSERT(writableController.isWritable()); if (backpressure) { auto chunkRef = js.v8Ref(chunk); return KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.whenResolved(js).then(js, JSG_VISITABLE_LAMBDA((chunkRef = kj::mv(chunkRef), ref=JSG_THIS), (chunkRef, ref), (jsg::Lock& js) mutable -> jsg::Promise { KJ_IF_SOME(writableController, ref->tryGetWritableController()) { KJ_IF_SOME(error, writableController.isErroring(js)) { return js.rejectedPromise(error); } else { // Else block to avert dangling else compiler warning. } } else { // Else block to avert dangling else compiler warning. } return ref->performTransform(js, chunkRef.getHandle(js)); })); } return performTransform(js, chunk); } else { return js.rejectedPromise( KJ_EXCEPTION(FAILED, "jsg.TypeError: Writing to the TransformStream failed.")); } } jsg::Promise TransformStreamDefaultController::abort( jsg::Lock& js, v8::Local reason) { if (FeatureFlags::get(js).getPedanticWpt()) { // If a finish operation is already in progress, return the existing promise // or handle the case where we're being called synchronously from within another // finish operation. if (algorithms.finishStarted) { KJ_IF_SOME(finish, algorithms.maybeFinish) { return finish.whenResolved(js); } // finishStarted is true but maybeFinish is not set yet - this means we're being // called synchronously from within another finish operation (like cancel). // We need to error the stream with the abort reason so that both the current // operation and this abort reject with the abort reason. error(js, reason); return js.rejectedPromise(js.v8Ref(reason)); } // Mark that we're starting a finish operation before running the algorithm. algorithms.finishStarted = true; } else { KJ_IF_SOME(finish, algorithms.maybeFinish) { return finish.whenResolved(js); } } return algorithms.maybeFinish .emplace(maybeRunAlgorithm(js, algorithms.cancel, JSG_VISITABLE_LAMBDA( (this, ref = JSG_THIS, reason = jsg::JsRef(js, jsg::JsValue(reason))), (ref, reason), (jsg::Lock & js)->jsg::Promise { // If the readable side is errored, return a rejected promise with the stored error { KJ_IF_SOME(err, getReadableErrorState(js)) { return js.rejectedPromise(kj::mv(err)); } else { // Else block to avert dangling else compiler warning. } } // Otherwise... error with the given reason and resolve the abort promise error(js, reason.getHandle(js)); return js.resolvedPromise(); }), JSG_VISITABLE_LAMBDA((this, ref = JSG_THIS), (ref), (jsg::Lock & js, jsg::Value reason)->jsg::Promise { error(js, reason.getHandle(js)); return js.rejectedPromise(kj::mv(reason)); }), jsg::JsValue(reason))) .whenResolved(js); } jsg::Promise TransformStreamDefaultController::close(jsg::Lock& js) { auto flags = FeatureFlags::get(js); if (flags.getPedanticWpt()) { // If a finish operation is already in progress (e.g., from cancel or abort), // we should not run flush. Per the WHATWG streams spec, close/flush should // coordinate with cancel to avoid calling both. if (algorithms.finishStarted) { KJ_IF_SOME(finish, algorithms.maybeFinish) { return finish.whenResolved(js); } // finishStarted is true but maybeFinish is not set yet - this means we're being // called synchronously from within another finish operation. If the stream was // errored during that operation, return a rejected promise with the error. KJ_IF_SOME(writableController, tryGetWritableController()) { KJ_IF_SOME(err, writableController.isErroredOrErroring(js)) { return js.rejectedPromise(err); } } KJ_IF_SOME(err, getReadableErrorState(js)) { return js.rejectedPromise(kj::mv(err)); } return js.resolvedPromise(); } // Mark that we're starting a finish operation before running the algorithm, // since the algorithm may synchronously call other finish operations. algorithms.finishStarted = true; } auto onSuccess = JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), (jsg::Lock & js)->jsg::Promise { // If the stream was errored during the flush algorithm (e.g., by controller.error() // or by a parallel cancel() calling abort()), we should reject with that error. if (FeatureFlags::get(js).getPedanticWpt()) { KJ_IF_SOME(err, ref->getReadableErrorState(js)) { return js.rejectedPromise(kj::mv(err)); } else { // Else block to avert dangling else compiler warning. } } // Allows for a graceful close of the readable side. Close will // complete once all of the queued data is read or the stream // errors. Only close if the stream can still be closed (e.g., // it wasn't closed by a cancel operation from within flush). { KJ_IF_SOME(readableController, ref->tryGetReadableController()) { if (readableController.canCloseOrEnqueue()) { readableController.close(js); } } else { // Else block to avert dangling else compiler warning. } } return js.resolvedPromise(); }); auto onFailure = JSG_VISITABLE_LAMBDA( (ref = JSG_THIS), (ref), (jsg::Lock & js, jsg::Value reason)->jsg::Promise { ref->error(js, reason.getHandle(js)); return js.rejectedPromise(kj::mv(reason)); }); if (flags.getPedanticWpt()) { return algorithms.maybeFinish .emplace( maybeRunAlgorithm(js, algorithms.flush, kj::mv(onSuccess), kj::mv(onFailure), JSG_THIS)) .whenResolved(js); } return maybeRunAlgorithm(js, algorithms.flush, kj::mv(onSuccess), kj::mv(onFailure), JSG_THIS); } jsg::Promise TransformStreamDefaultController::pull(jsg::Lock& js) { KJ_ASSERT(backpressure); setBackpressure(js, false); return KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.whenResolved(js); } jsg::Promise TransformStreamDefaultController::cancel( jsg::Lock& js, v8::Local reason) { if (FeatureFlags::get(js).getPedanticWpt()) { // If a finish operation is already in progress, return the existing promise // or check for errors if we're being called synchronously from within another // finish operation. if (algorithms.finishStarted) { KJ_IF_SOME(finish, algorithms.maybeFinish) { return finish.whenResolved(js); } // finishStarted is true but maybeFinish is not set yet - check if the stream // was errored during that operation. KJ_IF_SOME(err, getReadableErrorState(js)) { return js.rejectedPromise(kj::mv(err)); } return js.resolvedPromise(); } // Mark that we're starting a finish operation before running the algorithm. algorithms.finishStarted = true; } return algorithms.maybeFinish .emplace(maybeRunAlgorithm(js, algorithms.cancel, JSG_VISITABLE_LAMBDA( (this, ref = JSG_THIS, reason = jsg::JsRef(js, jsg::JsValue(reason))), (ref, reason), (jsg::Lock & js)->jsg::Promise { // If the stream was errored during the cancel algorithm (e.g., by controller.error() // or by a parallel abort()), we should reject with that error. if (FeatureFlags::get(js).getPedanticWpt()) { KJ_IF_SOME(err, getReadableErrorState(js)) { readable = kj::none; errorWritableAndUnblockWrite(js, reason.getHandle(js)); return js.rejectedPromise(kj::mv(err)); } else { // Else block to avert dangling else compiler warning. } } readable = kj::none; errorWritableAndUnblockWrite(js, reason.getHandle(js)); return js.resolvedPromise(); }), JSG_VISITABLE_LAMBDA((this, ref = JSG_THIS), (ref), (jsg::Lock & js, jsg::Value reason)->jsg::Promise { readable = kj::none; errorWritableAndUnblockWrite(js, reason.getHandle(js)); return js.rejectedPromise(kj::mv(reason)); }), jsg::JsValue(reason))) .whenResolved(js); } jsg::Promise TransformStreamDefaultController::performTransform( jsg::Lock& js, v8::Local chunk) { if (algorithms.transform != kj::none) { return maybeRunAlgorithm(js, algorithms.transform, [](jsg::Lock& js) -> jsg::Promise { return js.resolvedPromise(); }, JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), (jsg::Lock & js, jsg::Value reason)->jsg::Promise { ref->error(js, reason.getHandle(js)); return js.rejectedPromise(kj::mv(reason)); }), chunk, JSG_THIS); } // If we got here, there is no transform algorithm. Per the spec, the default // behavior then is to just pass along the value untransformed. return js.tryCatch([&] { enqueue(js, chunk); return js.resolvedPromise(); }, [&](jsg::Value exception) { return js.rejectedPromise(kj::mv(exception)); }); } void TransformStreamDefaultController::setBackpressure(jsg::Lock& js, bool newBackpressure) { KJ_ASSERT(newBackpressure != backpressure); KJ_IF_SOME(prp, maybeBackpressureChange) { prp.resolver.resolve(js); } maybeBackpressureChange = js.newPromiseAndResolver(); KJ_ASSERT_NONNULL(maybeBackpressureChange).promise.markAsHandled(js); backpressure = newBackpressure; } void TransformStreamDefaultController::errorWritableAndUnblockWrite( jsg::Lock& js, v8::Local reason) { algorithms.clear(); KJ_IF_SOME(writableController, tryGetWritableController()) { if (FeatureFlags::get(js).getPedanticWpt()) { // Use errorIfNeeded which goes through the proper error transition (Erroring -> Errored). // This allows close() to be called while the stream is "erroring" and reject with the // stored error, which is the expected behavior per the WHATWG streams spec. writableController.errorIfNeeded(js, reason); } else if (writableController.isWritable()) { writableController.doError(js, reason); } writable = kj::none; } if (backpressure) { setBackpressure(js, false); } } void TransformStreamDefaultController::visitForGc(jsg::GcVisitor& visitor) { KJ_IF_SOME(backpressureChange, maybeBackpressureChange) { visitor.visit(backpressureChange.promise, backpressureChange.resolver); } visitor.visit(writable, readable, startPromise.resolver, startPromise.promise, algorithms); } void TransformStreamDefaultController::init(jsg::Lock& js, jsg::Ref& readable, jsg::Ref& writable, jsg::Optional maybeTransformer) { KJ_ASSERT(this->readable == kj::none); KJ_ASSERT(this->writable == kj::none); this->writable = writable.addRef(); // The TransformStreamDefaultController needs to have a reference to the underlying controller // and not just the readable because if the readable is teed, or passed off to source, etc, // the TransformStream has to make sure that it can continue to interface with the controller // to push data into it. auto& readableController = static_cast(readable->getController()); auto readableRef = KJ_ASSERT_NONNULL(readableController.getController()); this->readable = KJ_ASSERT_NONNULL(readableRef.tryGet()).addRef(); auto transformer = kj::mv(maybeTransformer).orDefault({}); // TODO(someday): The stream standard includes placeholders for supporting byte-oriented // TransformStreams but does not yet define them. For now, we are limiting our implementation // here to only support value-based transforms. JSG_REQUIRE(transformer.readableType == kj::none, TypeError, "transformer.readableType must be undefined."); JSG_REQUIRE(transformer.writableType == kj::none, TypeError, "transformer.writableType must be undefined."); KJ_IF_SOME(transform, transformer.transform) { algorithms.transform = kj::mv(transform); } KJ_IF_SOME(flush, transformer.flush) { algorithms.flush = kj::mv(flush); } KJ_IF_SOME(cancel, transformer.cancel) { algorithms.cancel = kj::mv(cancel); } setBackpressure(js, true); maybeRunAlgorithm(js, transformer.start, JSG_VISITABLE_LAMBDA( (ref = JSG_THIS), (ref), (jsg::Lock& js) { ref->startPromise.resolver.resolve(js); }), JSG_VISITABLE_LAMBDA((ref = JSG_THIS), (ref), (jsg::Lock& js, jsg::Value reason) { ref->startPromise.resolver.reject(js, reason.getHandle(js)); }), JSG_THIS); } kj::Maybe TransformStreamDefaultController:: tryGetReadableController() { KJ_IF_SOME(controller, readable) { return *controller; } return kj::none; } kj::Maybe TransformStreamDefaultController:: tryGetWritableController() { KJ_IF_SOME(w, writable) { return static_cast(w->getController()); } return kj::none; } kj::Maybe TransformStreamDefaultController::getReadableErrorState(jsg::Lock& js) { KJ_IF_SOME(controller, tryGetReadableController()) { return controller.getMaybeErrorState(js); } return kj::none; } template kj::StringPtr WritableImpl::jsgGetMemoryName() const { return "WritableImpl"_kjc; } template size_t WritableImpl::jsgGetMemorySelfSize() const { return sizeof(WritableImpl); } template void WritableImpl::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("signal", signal); KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(error, StreamStates::Errored) { tracker.trackField("error", error); } KJ_CASE_ONEOF(erroring, StreamStates::Erroring) { tracker.trackField("erroring", erroring.reason); } KJ_CASE_ONEOF(writable, Writable) {} } tracker.trackField("abortAlgorithm", algorithms.abort); tracker.trackField("closeAlgorithm", algorithms.close); tracker.trackField("writeAlgorithm", algorithms.write); tracker.trackField("sizeAlgorithm", algorithms.size); for (auto& request: writeRequests) { tracker.trackField("pendingWrite", request); } tracker.trackField("inFlightWrite", inFlightWrite); tracker.trackField("inFlightClose", inFlightClose); tracker.trackField("closeRequest", closeRequest); tracker.trackField("maybePendingAbort", maybePendingAbort); } kj::StringPtr WritableStreamJsController::jsgGetMemoryName() const { return "WritableStreamJsController"_kjc; } size_t WritableStreamJsController::jsgGetMemorySelfSize() const { return sizeof(WritableStreamJsController); } void WritableStreamJsController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) {} KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(error, StreamStates::Errored) { tracker.trackField("error", error); } KJ_CASE_ONEOF(controller, Controller) { tracker.trackField("controller", controller); } } tracker.trackField("lock", lock); tracker.trackField("maybeAbortPromise", maybeAbortPromise); } void WritableStreamDefaultController::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("impl", impl); } kj::StringPtr ReadableStreamJsController::jsgGetMemoryName() const { return "ReadableStreamJsController"_kjc; } size_t ReadableStreamJsController::jsgGetMemorySelfSize() const { return sizeof(ReadableStreamJsController); } void ReadableStreamJsController::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(initial, Initial) {} KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(error, StreamStates::Errored) { tracker.trackField("error", error); } KJ_CASE_ONEOF(readable, kj::Own) { tracker.trackField("readable", readable); } KJ_CASE_ONEOF(readable, kj::Own) { tracker.trackField("readable", readable); } } tracker.trackField("lock", lock); // Track pending error state if present (Closed has no trackable content) KJ_IF_SOME(pendingError, state.tryGetPendingStateUnsafe()) { tracker.trackField("pendingError", pendingError); } } template kj::StringPtr ReadableImpl::jsgGetMemoryName() const { return "ReadableImpl"_kjc; } template size_t ReadableImpl::jsgGetMemorySelfSize() const { return sizeof(ReadableImpl); } template void ReadableImpl::jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const { KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(closed, StreamStates::Closed) {} KJ_CASE_ONEOF(error, StreamStates::Errored) { tracker.trackField("error", error); } KJ_CASE_ONEOF(queue, Queue) { tracker.trackField("queue", queue); } } tracker.trackField("startAlgorithm", algorithms.start); tracker.trackField("pullAlgorithm", algorithms.pull); tracker.trackField("cancelAlgorithm", algorithms.cancel); tracker.trackField("sizeAlgorithm", algorithms.size); tracker.trackField("pendingCancel", maybePendingCancel); } void ReadableStreamBYOBRequest::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { KJ_IF_SOME(impl, maybeImpl) { tracker.trackField("readRequest", impl.readRequest); tracker.trackField("view", impl.view); } } void TransformStreamDefaultController::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("startPromise", startPromise); tracker.trackField("maybeBackpressureChange", maybeBackpressureChange); tracker.trackField("transformAlgorithm", algorithms.transform); tracker.trackField("flushAlgorithm", algorithms.flush); tracker.trackField("writable", writable); tracker.trackField("readable", readable); } // ====================================================================================== jsg::Ref ReadableStream::from( jsg::Lock& js, jsg::AsyncGenerator generator) { // AsyncGenerator is not a refcounted type, so we need to wrap it in a refcounted // struct so that we can keep it alive through the various promise branches below. auto rcGenerator = kj::rc>>(kj::mv(generator)); // clang-format off return constructor(js, UnderlyingSource{ .pull = [generator = rcGenerator.addRef()](jsg::Lock& js, auto controller) mutable { auto& c = controller.template get(); return generator->getWrapped().next(js).then(js, JSG_VISITABLE_LAMBDA((controller = c.addRef(), generator = generator.addRef()), (controller), (jsg::Lock& js, kj::Maybe value) { KJ_IF_SOME(v, value) { auto handle = v.getHandle(js); // Per the ReadableStream.from spec, if the value is a promise, // the stream should wait for it to resolve and enqueue the // resolved value... // ... yes, this means that ReadableStream.from where the inputs // are promises will be slow, but that's the spec. if (handle->IsPromise()) { return js.toPromise(handle.As()).then(js, JSG_VISITABLE_LAMBDA( (controller=controller.addRef()), (controller), (jsg::Lock& js, jsg::Value val) mutable { controller->enqueue(js, val.getHandle(js)); return js.resolvedPromise(); })); } controller->enqueue(js, v.getHandle(js)); } else { controller->close(js); } return js.resolvedPromise(); }), JSG_VISITABLE_LAMBDA((controller = c.addRef(), generator = generator.addRef()), (controller), (jsg::Lock& js, jsg::Value reason) { controller->error(js, reason.getHandle(js)); return js.rejectedPromise(kj::mv(reason)); })); }, .cancel = [generator = rcGenerator.addRef()](jsg::Lock& js, auto reason) mutable { return generator->getWrapped().return_(js, js.v8Ref(reason)) .then(js, [generator = kj::mv(generator)](auto& lock, auto) { // The generator might produce a value on return and might even want to continue, // but the stream has been canceled at this point, so we stop here. }); }, }, StreamQueuingStrategy{ .highWaterMark = 0 }); // clang-format on } } // namespace workerd::api