// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include "common.h" #include #include namespace workerd::api { class WritableStreamDefaultWriter: public jsg::Object, public WritableStreamController::Writer { public: explicit WritableStreamDefaultWriter(); ~WritableStreamDefaultWriter() noexcept(false) override; // JavaScript API static jsg::Ref constructor( jsg::Lock& js, jsg::Ref stream); jsg::MemoizedIdentity>& getClosed(); jsg::MemoizedIdentity>& getReady(); kj::Maybe getDesiredSize(); jsg::Promise abort(jsg::Lock& js, jsg::Optional> reason); // Closes the stream. All present write requests will complete, but future write requests will // be rejected with a TypeError to the effect of "This writable stream has been closed." // `reason` will be passed to the underlying sink's close algorithm -- if this writable stream // is one side of a transform stream, then its close algorithm causes the transform's readable // side to become closed. // // Note: According to my reading of the Streams spec, if `writer.close()` is called on a // transform stream while the readable side has readable chunks in its queue, those chunks get // lost. This seems like a bug to me. Why would we wait for all present write requests to // complete on this side if we don't care that they're actually read? jsg::Promise close(jsg::Lock& js); jsg::Promise write(jsg::Lock& js, jsg::Optional> chunk); void releaseLock(jsg::Lock& js); JSG_RESOURCE_TYPE(WritableStreamDefaultWriter, CompatibilityFlags::Reader flags) { if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed); JSG_READONLY_PROTOTYPE_PROPERTY(ready, getReady); JSG_READONLY_PROTOTYPE_PROPERTY(desiredSize, getDesiredSize); } else { JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed); JSG_READONLY_INSTANCE_PROPERTY(ready, getReady); JSG_READONLY_INSTANCE_PROPERTY(desiredSize, getDesiredSize); } JSG_METHOD(abort); JSG_METHOD(close); JSG_METHOD(write); JSG_METHOD(releaseLock); JSG_TS_OVERRIDE( { write(chunk?: W): Promise; }); } // Internal API void attach(jsg::Lock& js, WritableStreamController& controller, jsg::Promise closedPromise, jsg::Promise readyPromise) override; void detach() override; void lockToStream(jsg::Lock& js, WritableStream& stream); void replaceReadyPromise(jsg::Lock& js, jsg::Promise readyPromise) override; void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; kj::Maybe> isReady(jsg::Lock& js); private: struct Initial { static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; }; // While a Writer is attached to a WritableStream, it holds a strong reference to the // WritableStream to prevent it from being GC'ed so long as the Writer is available. // Once the writer is closed, released, or GC'ed the reference to the WritableStream // is cleared and the WritableStream can be GC'ed if there are no other references to // it being held anywhere. If the writer is still attached to the WritableStream when // it is destroyed, the WritableStream's reference to the writer is cleared but the // WritableStream remains in the "writer locked" state, per the spec. struct Attached { static constexpr kj::StringPtr NAME KJ_UNUSED = "attached"_kj; jsg::Ref stream; }; // Released: The user explicitly called releaseLock() to detach the writer from the stream. // The stream remains usable and can be locked by a new writer. struct Released { static constexpr kj::StringPtr NAME KJ_UNUSED = "released"_kj; }; // Closed: The underlying stream ended (closed or errored) while the writer was attached. // The stream is no longer usable. struct Closed { static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; }; // State machine for WritableStreamDefaultWriter: // Initial -> Attached (attach() called) // Attached -> Closed (detach() called when stream closes) // Attached -> Released (releaseLock() called) // Closed and Released are terminal states. // Initial is not terminal but most methods assert if called in this state. using WriterState = StateMachine, ActiveState, Initial, Attached, Closed, Released>; kj::Maybe ioContext; WriterState state; inline void assertAttachedOrTerminal() const { KJ_ASSERT(!state.is(), "this writer was never attached"); } kj::Maybe>> closedPromise; kj::Maybe>> readyPromise; kj::Maybe> readyPromisePending; void visitForGc(jsg::GcVisitor& visitor); }; class WritableStream: public jsg::Object { public: explicit WritableStream(IoContext& ioContext, kj::Own sink, kj::Maybe> observer, kj::Maybe maybeHighWaterMark = kj::none, kj::Maybe> maybeClosureWaitable = kj::none); explicit WritableStream(kj::Own controller); ~WritableStream() noexcept(false) { weakRef->invalidate(); } WritableStreamController& getController(); jsg::Ref addRef(); // Remove and return the underlying implementation of this WritableStream. Throw a TypeError if // this WritableStream is locked or closed, otherwise this WritableStream becomes immediately // locked and closed. If this writable stream is errored, throw the stored error. // TODO(cleanup): There are a couple of places where we need to convert to using detach() // or the inner removeSink (on WritableStreamController) before we can remove this method. virtual KJ_DEPRECATED("Use detach() instead") kj::Own removeSink( jsg::Lock& js); virtual void detach(jsg::Lock& js); // --------------------------------------------------------------------------- // JS interface static jsg::Ref constructor(jsg::Lock& js, jsg::Optional underlyingSink, jsg::Optional queuingStrategy); bool isLocked(); // Errors the stream. All present and future read requests are rejected with a TypeError to the // effect of "This writable stream has been requested to abort." `reason` will be passed to the // underlying sink's abort algorithm -- if this writable stream is one side of a transform stream, // then its abort algorithm causes the transform's readable side to become errored with `reason`. jsg::Promise abort(jsg::Lock& js, jsg::Optional> reason); jsg::Promise close(jsg::Lock& js); jsg::Promise flush(jsg::Lock& js); jsg::Ref getWriter(jsg::Lock& js); jsg::JsString inspectState(jsg::Lock& js); bool inspectExpectsBytes(); JSG_RESOURCE_TYPE(WritableStream, CompatibilityFlags::Reader flags) { if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_READONLY_PROTOTYPE_PROPERTY(locked, isLocked); } else { JSG_READONLY_INSTANCE_PROPERTY(locked, isLocked); } JSG_METHOD(abort); JSG_METHOD(close); JSG_METHOD(getWriter); JSG_INSPECT_PROPERTY(state, inspectState); JSG_INSPECT_PROPERTY(expectsBytes, inspectExpectsBytes); JSG_TS_OVERRIDE( { getWriter(): WritableStreamDefaultWriter; }); } void serialize(jsg::Lock& js, jsg::Serializer& serializer); static jsg::Ref deserialize( jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer); JSG_SERIALIZABLE(rpc::SerializationTag::WRITABLE_STREAM); void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; private: kj::Maybe ioContext; kj::Own controller; kj::Own> weakRef = kj::refcounted>(kj::Badge(), *this); kj::Own> addWeakRef() { return weakRef->addRef(); } void visitForGc(jsg::GcVisitor& visitor); template friend class WritableImpl; }; } // namespace workerd::api