File
Blob: src/workerd/api/streams/writable.h
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #pragma once |
| 6 | |
| 7 | #include "common.h" |
| 8 | |
| 9 | #include <workerd/util/state-machine.h> |
| 10 | #include <workerd/util/weak-refs.h> |
| 11 | |
| 12 | namespace workerd::api { |
| 13 | |
| 14 | class WritableStreamDefaultWriter: public jsg::Object, public WritableStreamController::Writer { |
| 15 | public: |
| 16 | explicit WritableStreamDefaultWriter(); |
| 17 | |
| 18 | ~WritableStreamDefaultWriter() noexcept(false) override; |
| 19 | |
| 20 | // JavaScript API |
| 21 | |
| 22 | static jsg::Ref<WritableStreamDefaultWriter> constructor( |
| 23 | jsg::Lock& js, jsg::Ref<WritableStream> stream); |
| 24 | |
| 25 | jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed(); |
| 26 | jsg::MemoizedIdentity<jsg::Promise<void>>& getReady(); |
| 27 | kj::Maybe<int> getDesiredSize(); |
| 28 | |
| 29 | jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason); |
| 30 | |
| 31 | // Closes the stream. All present write requests will complete, but future write requests will |
| 32 | // be rejected with a TypeError to the effect of "This writable stream has been closed." |
| 33 | // `reason` will be passed to the underlying sink's close algorithm -- if this writable stream |
| 34 | // is one side of a transform stream, then its close algorithm causes the transform's readable |
| 35 | // side to become closed. |
| 36 | // |
| 37 | // Note: According to my reading of the Streams spec, if `writer.close()` is called on a |
| 38 | // transform stream while the readable side has readable chunks in its queue, those chunks get |
| 39 | // lost. This seems like a bug to me. Why would we wait for all present write requests to |
| 40 | // complete on this side if we don't care that they're actually read? |
| 41 | jsg::Promise<void> close(jsg::Lock& js); |
| 42 | |
| 43 | jsg::Promise<void> write(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> chunk); |
| 44 | void releaseLock(jsg::Lock& js); |
| 45 | |
| 46 | JSG_RESOURCE_TYPE(WritableStreamDefaultWriter, CompatibilityFlags::Reader flags) { |
| 47 | if (flags.getJsgPropertyOnPrototypeTemplate()) { |
| 48 | JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed); |
| 49 | JSG_READONLY_PROTOTYPE_PROPERTY(ready, getReady); |
| 50 | JSG_READONLY_PROTOTYPE_PROPERTY(desiredSize, getDesiredSize); |
| 51 | } else { |
| 52 | JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed); |
| 53 | JSG_READONLY_INSTANCE_PROPERTY(ready, getReady); |
| 54 | JSG_READONLY_INSTANCE_PROPERTY(desiredSize, getDesiredSize); |
| 55 | } |
| 56 | JSG_METHOD(abort); |
| 57 | JSG_METHOD(close); |
| 58 | JSG_METHOD(write); |
| 59 | JSG_METHOD(releaseLock); |
| 60 | |
| 61 | JSG_TS_OVERRIDE(<W = any> { |
| 62 | write(chunk?: W): Promise<void>; |
| 63 | }); |
| 64 | } |
| 65 | |
| 66 | // Internal API |
| 67 | |
| 68 | void attach(jsg::Lock& js, |
| 69 | WritableStreamController& controller, |
| 70 | jsg::Promise<void> closedPromise, |
| 71 | jsg::Promise<void> readyPromise) override; |
| 72 | |
| 73 | void detach() override; |
| 74 | |
| 75 | void lockToStream(jsg::Lock& js, WritableStream& stream); |
| 76 | |
| 77 | void replaceReadyPromise(jsg::Lock& js, jsg::Promise<void> readyPromise) override; |
| 78 | |
| 79 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 80 | |
| 81 | kj::Maybe<jsg::Promise<void>> isReady(jsg::Lock& js); |
| 82 | |
| 83 | private: |
| 84 | struct Initial { |
| 85 | static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; |
| 86 | }; |
| 87 | // While a Writer is attached to a WritableStream, it holds a strong reference to the |
| 88 | // WritableStream to prevent it from being GC'ed so long as the Writer is available. |
| 89 | // Once the writer is closed, released, or GC'ed the reference to the WritableStream |
| 90 | // is cleared and the WritableStream can be GC'ed if there are no other references to |
| 91 | // it being held anywhere. If the writer is still attached to the WritableStream when |
| 92 | // it is destroyed, the WritableStream's reference to the writer is cleared but the |
| 93 | // WritableStream remains in the "writer locked" state, per the spec. |
| 94 | struct Attached { |
| 95 | static constexpr kj::StringPtr NAME KJ_UNUSED = "attached"_kj; |
| 96 | jsg::Ref<WritableStream> stream; |
| 97 | }; |
| 98 | // Released: The user explicitly called releaseLock() to detach the writer from the stream. |
| 99 | // The stream remains usable and can be locked by a new writer. |
| 100 | struct Released { |
| 101 | static constexpr kj::StringPtr NAME KJ_UNUSED = "released"_kj; |
| 102 | }; |
| 103 | // Closed: The underlying stream ended (closed or errored) while the writer was attached. |
| 104 | // The stream is no longer usable. |
| 105 | struct Closed { |
| 106 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 107 | }; |
| 108 | |
| 109 | // State machine for WritableStreamDefaultWriter: |
| 110 | // Initial -> Attached (attach() called) |
| 111 | // Attached -> Closed (detach() called when stream closes) |
| 112 | // Attached -> Released (releaseLock() called) |
| 113 | // Closed and Released are terminal states. |
| 114 | // Initial is not terminal but most methods assert if called in this state. |
| 115 | using WriterState = StateMachine<TerminalStates<Closed, Released>, |
| 116 | ActiveState<Attached>, |
| 117 | Initial, |
| 118 | Attached, |
| 119 | Closed, |
| 120 | Released>; |
| 121 | |
| 122 | kj::Maybe<IoContext&> ioContext; |
| 123 | WriterState state; |
| 124 | |
| 125 | inline void assertAttachedOrTerminal() const { |
| 126 | KJ_ASSERT(!state.is<Initial>(), "this writer was never attached"); |
| 127 | } |
| 128 | |
| 129 | kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<void>>> closedPromise; |
| 130 | kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<void>>> readyPromise; |
| 131 | kj::Maybe<jsg::Promise<void>> readyPromisePending; |
| 132 | |
| 133 | void visitForGc(jsg::GcVisitor& visitor); |
| 134 | }; |
| 135 | |
| 136 | class WritableStream: public jsg::Object { |
| 137 | public: |
| 138 | explicit WritableStream(IoContext& ioContext, |
| 139 | kj::Own<WritableStreamSink> sink, |
| 140 | kj::Maybe<kj::Own<ByteStreamObserver>> observer, |
| 141 | kj::Maybe<uint64_t> maybeHighWaterMark = kj::none, |
| 142 | kj::Maybe<jsg::Promise<void>> maybeClosureWaitable = kj::none); |
| 143 | |
| 144 | explicit WritableStream(kj::Own<WritableStreamController> controller); |
| 145 | ~WritableStream() noexcept(false) { |
| 146 | weakRef->invalidate(); |
| 147 | } |
| 148 | |
| 149 | WritableStreamController& getController(); |
| 150 | |
| 151 | jsg::Ref<WritableStream> addRef(); |
| 152 | |
| 153 | // Remove and return the underlying implementation of this WritableStream. Throw a TypeError if |
| 154 | // this WritableStream is locked or closed, otherwise this WritableStream becomes immediately |
| 155 | // locked and closed. If this writable stream is errored, throw the stored error. |
| 156 | // TODO(cleanup): There are a couple of places where we need to convert to using detach() |
| 157 | // or the inner removeSink (on WritableStreamController) before we can remove this method. |
| 158 | virtual KJ_DEPRECATED("Use detach() instead") kj::Own<WritableStreamSink> removeSink( |
| 159 | jsg::Lock& js); |
| 160 | virtual void detach(jsg::Lock& js); |
| 161 | |
| 162 | // --------------------------------------------------------------------------- |
| 163 | // JS interface |
| 164 | |
| 165 | static jsg::Ref<WritableStream> constructor(jsg::Lock& js, |
| 166 | jsg::Optional<UnderlyingSink> underlyingSink, |
| 167 | jsg::Optional<StreamQueuingStrategy> queuingStrategy); |
| 168 | |
| 169 | bool isLocked(); |
| 170 | |
| 171 | // Errors the stream. All present and future read requests are rejected with a TypeError to the |
| 172 | // effect of "This writable stream has been requested to abort." `reason` will be passed to the |
| 173 | // underlying sink's abort algorithm -- if this writable stream is one side of a transform stream, |
| 174 | // then its abort algorithm causes the transform's readable side to become errored with `reason`. |
| 175 | jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason); |
| 176 | |
| 177 | jsg::Promise<void> close(jsg::Lock& js); |
| 178 | jsg::Promise<void> flush(jsg::Lock& js); |
| 179 | |
| 180 | jsg::Ref<WritableStreamDefaultWriter> getWriter(jsg::Lock& js); |
| 181 | |
| 182 | jsg::JsString inspectState(jsg::Lock& js); |
| 183 | bool inspectExpectsBytes(); |
| 184 | |
| 185 | JSG_RESOURCE_TYPE(WritableStream, CompatibilityFlags::Reader flags) { |
| 186 | if (flags.getJsgPropertyOnPrototypeTemplate()) { |
| 187 | JSG_READONLY_PROTOTYPE_PROPERTY(locked, isLocked); |
| 188 | } else { |
| 189 | JSG_READONLY_INSTANCE_PROPERTY(locked, isLocked); |
| 190 | } |
| 191 | JSG_METHOD(abort); |
| 192 | JSG_METHOD(close); |
| 193 | JSG_METHOD(getWriter); |
| 194 | |
| 195 | JSG_INSPECT_PROPERTY(state, inspectState); |
| 196 | JSG_INSPECT_PROPERTY(expectsBytes, inspectExpectsBytes); |
| 197 | |
| 198 | JSG_TS_OVERRIDE(<W = any> { |
| 199 | getWriter(): WritableStreamDefaultWriter<W>; |
| 200 | }); |
| 201 | } |
| 202 | |
| 203 | void serialize(jsg::Lock& js, jsg::Serializer& serializer); |
| 204 | static jsg::Ref<WritableStream> deserialize( |
| 205 | jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer); |
| 206 | |
| 207 | JSG_SERIALIZABLE(rpc::SerializationTag::WRITABLE_STREAM); |
| 208 | |
| 209 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 210 | |
| 211 | private: |
| 212 | kj::Maybe<IoContext&> ioContext; |
| 213 | kj::Own<WritableStreamController> controller; |
| 214 | kj::Own<WeakRef<WritableStream>> weakRef = |
| 215 | kj::refcounted<WeakRef<WritableStream>>(kj::Badge<WritableStream>(), *this); |
| 216 | |
| 217 | kj::Own<WeakRef<WritableStream>> addWeakRef() { |
| 218 | return weakRef->addRef(); |
| 219 | } |
| 220 | |
| 221 | void visitForGc(jsg::GcVisitor& visitor); |
| 222 | |
| 223 | template <typename T> |
| 224 | friend class WritableImpl; |
| 225 | }; |
| 226 | |
| 227 | } // namespace workerd::api |