Skip to content
File

Blob: src/workerd/api/streams/writable.h

cpp228 lines
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 
12namespace workerd::api {
13 
14class 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 
136class 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