Skip to content
File

Blob: src/workerd/api/streams/writable-sink-adapter.h

cpp486 lines
1#include "common.h"
2#include "writable-sink.h"
3 
4#include <workerd/util/state-machine.h>
5#include <workerd/util/weak-refs.h>
6 
7namespace workerd::api::streams {
8// Wraps a WritableStreamSink with a more JS-friendly interface that implements
9// queued writes and backpressure signaling. This is arguably what WritableStreamSink
10// should have been in the first place. Eventually we might be able to replace
11// WritableStreamSink with this class directly, but for now we need to keep both.
12//
13// Instances of WritableStreamSinkJsAdapter are meant to be used from within the
14// isolate lock, when you have need to write data to a kj stream from JavaScript.
15// As such, it is not a jsg::Object itself, nor is it a kj I/O object, but it
16// sits between the two worlds. Internally it holds the WritableStreamSink within
17// an IoOwn so that correct IoContext usage is enforced. But the kj::Own for the
18// adapter itself is meant to be held in JS land.
19//
20// Once created, the adapter owns the underlying WritableStreamSink. It is not
21// possible to extract the sink from the adapter. This is because the adapter
22// needs to be able to enforce its own state machine and queued write mechanism.
23//
24// The adapter implements backpressure signaling based on a high water mark
25// configured at construction time. When the number of bytes in flight exceeds
26// the high water mark, we signal backpressure by causing the ready promise
27// to be reset to a new pending promise. When backpressure is released again,
28// the ready promise is resolved. The identity of the ready promise changes
29// whenever the backpressure state changes.
30//
31// The adapter also implements flush signaling. Flushing signals are checkpoints
32// that are inserted into the write queue, essentially like a no-op write. They
33// can be used as synchronization points to ensure that all prior writes have
34// completed. Flush signals do not affect backpressure or stream state.
35//
36// Dropping the adapter will cancel any in-flight and pending operations
37// immediately. Dropping the IoContext while the adapter is still active
38// will also cancel any in-flight and pending operations and cause the
39// adapter to be invalidated (the Active state is held with an IoOwn).
40//
41// ┌───────────────────────────────────────────┐
42// │ JavaScript Code │
43// │ │
44// │ • write(data) → Promise<void> │
45// │ • flush() → Promise<void> │
46// │ • end() → Promise<void> │
47// │ • abort(reason) │
48// │ • getReady() → Promise<void> │
49// └───────────────────────────────────────────┘
50// │
51// ▼
52// ┌───────────────────────────────────────────┐
53// │ WritableStreamSinkJsAdapter │
54// │ │
55// │ ┌─────────────────────────────────────┐ │
56// │ │ JavaScript API │ │
57// │ │ │ │
58// │ │ • write(data) → Promise<void> │ │
59// │ │ • flush() → Promise<void> │ │
60// │ │ • end() → Promise<void> │ │
61// │ │ • abort(reason) │ │
62// │ │ • getReady() → Promise<void> │ │
63// │ │ • getDesiredSize() → number │ │
64// │ └─────────────────────────────────────┘ │
65// │ │ │
66// │ ▼ │
67// │ ┌─────────────────────────────────────┐ │
68// │ │ Backpressure Management │ │
69// │ │ │ │
70// │ │ • High water mark (16KB default) │ │
71// │ │ • Bytes in flight tracking │ │
72// │ │ • Ready promise signaling │ │
73// │ │ • Queue depth management │ │
74// │ └─────────────────────────────────────┘ │
75// │ │ │
76// │ ▼ │
77// │ ┌─────────────────────────────────────┐ │
78// │ │ Write Queue Management │ │
79// │ │ │ │
80// │ │ • Queued writes with ordering │ │
81// │ │ • Flush checkpoints │ │
82// │ │ • Single in-flight write │ │
83// │ │ • Error propagation │ │
84// │ └─────────────────────────────────────┘ │
85// │ │ │
86// │ ▼ │
87// │ ┌─────────────────────────────────────┐ │
88// │ │ KJ Integration │ │
89// │ │ │ │
90// │ │ IoOwn<WritableStreamSink> │ │
91// │ │ WeakRef for safe references │ │
92// │ │ IoContext-aware operations │ │
93// │ └─────────────────────────────────────┘ │
94// └───────────────────────────────────────────┘
95// │
96// ▼
97// ┌───────────────────────────────────────────┐
98// │ WritableStreamSink │
99// │ (KJ Native Sink) │
100// │ │
101// │ • write(buffer) → Promise<void> │
102// │ • end() → Promise<void> │
103// │ • abort(reason) │
104// └───────────────────────────────────────────┘
105//
106class WritableStreamSinkJsAdapter final {
107 public:
108 struct Options {
109 // While the WritableStreamSink interface, and kj streams in general, do
110 // not have a notion of backpressure, and instead generally require only
111 // one write to be in flight at a time, it's better for performance for
112 // us to be able to buffer a bit more data in flight. So we will implement
113 // a simple high water mark mechanism. The default is 16KB.
114 size_t highWaterMark = 16384;
115 
116 // When detachOnWrite is true, and a write() is made with an ArrayBuffer,
117 // or ArrayBufferView, we will attempt to detach the underlying buffer
118 // before writing it to the sink. Detaching is required by the
119 // streams spec but our original implementation does not detach
120 // and it turns out there are old workers depending on that behavior.
121 bool detachOnWrite = false;
122 };
123 
124 WritableStreamSinkJsAdapter(jsg::Lock& js,
125 IoContext& ioContext,
126 kj::Own<WritableSink> sink,
127 kj::Maybe<Options> options = kj::none);
128 WritableStreamSinkJsAdapter(jsg::Lock& js,
129 IoContext& ioContext,
130 kj::Own<kj::AsyncOutputStream> stream,
131 StreamEncoding encoding,
132 kj::Maybe<Options> options = kj::none);
133 KJ_DISALLOW_COPY_AND_MOVE(WritableStreamSinkJsAdapter);
134 ~WritableStreamSinkJsAdapter() noexcept(false);
135 
136 // If we are in the errored state, returns the exception, otherwise kj::none.
137 kj::Maybe<const kj::Exception&> isErrored() KJ_LIFETIMEBOUND;
138 
139 // Returns true if we are in the closed state.
140 bool isClosed();
141 
142 // Returns true if close() has been called but we are not yet closed.
143 bool isClosing();
144 
145 // If we are not in the closed or errored state, returns the desired
146 // size based on the configured high water mark and the number of
147 // bytes currently in flight. The desired size is the number of bytes
148 // that can be written before we exceed the high water mark. If the
149 // return value is <= 0 then backpressure is being signaled. If we are
150 // in the closed or errored states, returns kj::none.
151 kj::Maybe<ssize_t> getDesiredSize();
152 
153 // Writes a chunk to the underlying sink via the queued write mechanism.
154 // The implementation ensures that only one write is in flight with the
155 // underlying sink at a time, while additional writes are queued up
156 // behind it. It is not necessary to await the returned promise before
157 // calling write() again, though doing so is not an error. If the write
158 // fails, the returned promise will reject with the failure reason.
159 // Also if the write fails, the adapter will be transitioned to the
160 // errored state and all subsequent queued writes will fail. Once
161 // close() has been called, no additional writes will be accepted
162 // and the returned promise will reject with an error. If the adapter
163 // is already in the closed or errored state, the returned promise will
164 // be rejected.
165 //
166 // Values written may be ArrayBuffer, ArrayBufferView, SharedArrayBuffer,
167 // or string. Other types will cause the returned promise to reject.
168 //
169 // Backpressure is signaled when the number of bytes in flight (i.e.
170 // the total number of bytes passed to write() calls that have not yet
171 // completed) exceeds the configured high water mark. When backpressure
172 // is signaled, additional writes are still accepted and queued up, but
173 // the caller really should wait for the ready promise to resolve before
174 // continuing to write more. This works exactly like a WritableStream's
175 // backpressure mechanism. Callers keep writing until backpressure is
176 // signaled, then wait for the ready promise to resolve before continuing,
177 // etc.
178 jsg::Promise<void> write(jsg::Lock& js, const jsg::JsValue& value);
179 
180 // Inserts a flush signal into the write queue. The returned promise
181 // resolves once all prior writes have completed. This can be used
182 // as a synchronization point to ensure that all writes up to this
183 // point have been fully processed. If the adapter is in the closed
184 // or errored state, the returned promise will reject. If the stream
185 // errors while waiting for prior writes to complete, the returned
186 // promise will be rejected.
187 jsg::Promise<void> flush(jsg::Lock& js);
188 
189 // Transitions the adapter into the closing state. Once the write queue
190 // is empty, we will close the sink and transition to the closed state.
191 // If the adapter is already in the closing state, a new promise is
192 // returned that will resolve when the adapter is fully closed. If the
193 // adapter is already closed, a resolved promise is returned. If the
194 // adapter is in the errored state, a rejected promise is returned.
195 // All pending writes in the queue will be processed before closing
196 // the sink and transitioning to the closed state. If any pending
197 // writes fail, the adapter will transition to the errored state, and
198 // all subsequent pending writes will be rejected along with the close
199 // promise.
200 jsg::Promise<void> end(jsg::Lock& js);
201 
202 // Transitions the adapter to the errored state, even if we are already closed.
203 // All pending or in-flight writes, and a pending close, will all be rejected
204 // with the given exception. If we are already in the errored state, this
205 // is a no-op. This change is immediate. Once in the errored state, no
206 // further writes or closes are allowed.
207 void abort(kj::Exception&& exception);
208 
209 // Transitions the adapter to the errored state, even if we are already closed.
210 // All pending or in-flight writes, and a pending close, will all be rejected
211 // with the given exception. If we are already in the errored state, this
212 // is a no-op. This change is immediate. Once in the errored state, no
213 // further writes or closes are allowed. This variant is for use when
214 // the exception is coming from JavaScript. It will be converted into a
215 // tunneled kj::Exception.
216 void abort(jsg::Lock& js, const jsg::JsValue& reason);
217 
218 // Returns a promise that resolves when backpressure is released.
219 // Note that the identity of the returned promise will change as the
220 // backpressure state changes. Whenever backpressure is signaled, a new
221 // pending promise will be created, whenever backpressure is released
222 // again that promise will be resolved. As such, this promise should
223 // not be cached or stored. Instead, before every write() call, the
224 // caller should wait on the current getReady() promise.
225 jsg::Promise<void> getReady(jsg::Lock& js);
226 
227 // Returns a memoized identity for the ready promise. This can be used
228 // to return a stable reference to the ready promise out to JavaScript
229 // that will not change identity between calls unless the backpressure
230 // state changes. Like the getReady() promise, this should not be cached
231 // or stored, but it is safe to return this from a getter multiple times
232 // to JavaScript as it will ensure that the same JS promise object is
233 // always returned until the backpressure state changes. This variation
234 // is not suitable for use within C++ code that needs to await on the
235 // ready promise because the internal jsg::Promise<void> object will
236 // no longer exist once the reference is passed out to JavaScript.
237 jsg::MemoizedIdentity<jsg::Promise<void>>& getReadyStable();
238 
239 // Returns the options used to configure this adapter if the adapter
240 // is not closed or errored.
241 kj::Maybe<const Options&> getOptions();
242 
243 void visitForGc(jsg::GcVisitor& visitor);
244 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
245 
246 private:
247 // Represents the active state of the adapter. Importantly, this state
248 // holds both the underlying WritableStreamSink and the write queue.
249 // It must be held within an IoOwn.
250 struct Active;
251 
252 struct Closed final {
253 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
254 };
255 
256 struct Open {
257 static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj;
258 IoOwn<Active> active;
259 };
260 
261 // State machine for tracking writable sink adapter lifecycle:
262 // Open -> Closed (normal close via end())
263 // Open -> kj::Exception (error via abort() or write failure)
264 // Closed is terminal, kj::Exception is implicitly terminal via ErrorState.
265 using State = StateMachine<TerminalStates<Closed>,
266 ErrorState<kj::Exception>,
267 ActiveState<Open>,
268 Open,
269 Closed,
270 kj::Exception>;
271 State state;
272 
273 // Used for backpressure signaling. When backpressure is indicated, the
274 // readyResolver, ready, and readyWatcher will be replaced with a new set.
275 // When backpressure is relieved, the readyResolver will be resolved.
276 // The adapter will start out in a ready state.
277 struct BackpressureState final {
278 // Note that if the BackpressureState is dropped while in a waiting state,
279 // the ready promise will be left unresolved. This is OK.
280 kj::Maybe<jsg::Promise<void>::Resolver> readyResolver;
281 jsg::Promise<void> ready;
282 jsg::MemoizedIdentity<jsg::Promise<void>> readyWatcher;
283 
284 // Aborts backpressure signaling, likely because the adapter is being errored.
285 // Causes the ready promise to be rejected with the given reason.
286 void abort(jsg::Lock& js, const jsg::JsValue& reason);
287 
288 // Releases backpressure, resolving the ready promise.
289 void release(jsg::Lock& js);
290 
291 // Indicates that backpressure has been signaled and we are waiting
292 // for it to be released or aborted.
293 bool isWaiting() const;
294 
295 // Returns a promise that resolves when backpressure is released.
296 // Note that every call to this returns a new jsg::Promise<void>
297 // instance. Callers that need a stable identity should use
298 // getReadyStable() instead (generally this is only the case when
299 // returning the promise to JavaScript via a getter).
300 jsg::Promise<void> getReady(jsg::Lock& js);
301 
302 // Returns a memoized identity for the ready promise. This can be used
303 // to return a stable reference to the ready promise out to JavaScript
304 // that will not change identity between calls unless the backpressure
305 // state changes.
306 jsg::MemoizedIdentity<jsg::Promise<void>>& getReadyStable();
307 BackpressureState(jsg::Promise<void>::Resolver&& resolver,
308 jsg::Promise<void>&& promise,
309 jsg::MemoizedIdentity<jsg::Promise<void>>&& watcher);
310 };
311 BackpressureState backpressureState;
312 kj::Rc<WeakRef<WritableStreamSinkJsAdapter>> selfRef;
313 
314 // Replaces the backpressure state with a new one, indicating that backpressure
315 // is being applied. If we are already in a backpressure state, this is a no-op.
316 // This will cause the ready promise (and its stable identity) to change.
317 void maybeSignalBackpressure(jsg::Lock& js);
318 
319 // Conditionally releases backpressure if the desired size is now > 0.
320 void maybeReleaseBackpressure(jsg::Lock& js);
321 
322 // Creates a new BackpressureState in the waiting state.
323 static BackpressureState newBackpressureState(jsg::Lock& js);
324};
325 
326// ================================================================================
327 
328// Adapts a WritableStream to a KJ-frendly interface.
329// The adapter fully wraps the WritableStream instance,
330// using a WritableStreamDefaultWriter to push data to it.
331// Then the adapter is destroyed or aborted, the writer is
332// aborted and both the writer and the stream references
333// are dropped. Critically, the stream is not usable after
334// ownership is transferred to this adapter. Initializing the adapter
335// will fail if the stream is already locked.
336//
337// If the adapter is dropped, or aborted while there are pending writes,
338// the pending writes will be rejected with the same exception as the abort.
339//
340// While WritableStream itself allows multiple writes to be in flight
341// at the same time, the WritableStreamSink interface does not, so
342// the adapter will ensure that only one write is in flight at a time.
343//
344// While the caller is expected to follow the WritableStreamSink contract
345// and keep the adapter alive until the write promises resolve, there
346// are some protections in place to avoid use-after-free if the caller
347// drops the adapter. There's nothing we can do if the caller drops the
348// buffer, however, so that is still a hard requirement.
349// TODO(safety): This can be made safer by having write take a kj::Array
350// as input instead of a kj::ArrayPtr but that's a larger refactor.
351//
352// ┌───────────────────────────────────────────┐
353// │ WritableStreamSink │
354// │ │
355// │ • write(buffer) │
356// │ • write(pieces[]) │
357// │ • end() │
358// │ • abort(reason) │
359// └───────────────────────────────────────────┘
360// │
361// ▼
362// ┌───────────────────────────────────────────┐
363// │ WritableStreamSinkKjAdapter │
364// │ │
365// │ ┌─────────────────────────────────────┐ │
366// │ │ KJ Native API │ │
367// │ │ │ │
368// │ │ • write(ArrayPtr<byte>) │ │
369// │ │ • write(ArrayPtr<ArrayPtr<byte>>) │ │
370// │ │ • end() → Promise<void> │ │
371// │ │ • abort(exception) │ │
372// │ └─────────────────────────────────────┘ │
373// │ │ │
374// │ ▼ │
375// │ ┌─────────────────────────────────────┐ │
376// │ │ State Management │ │
377// │ │ │ │
378// │ │ Active ──► Closed │ │
379// │ │ │ │ │ │
380// │ │ │ ▼ │ │
381// │ │ └─────► Errored │ │
382// │ └─────────────────────────────────────┘ │
383// │ │ │
384// │ ▼ │
385// │ ┌─────────────────────────────────────┐ │
386// │ │ JavaScript Integration │ │
387// │ │ │ │
388// │ │ WritableStreamDefaultWriter │ │
389// │ │ WeakRef for safe references │ │
390// │ │ IoContext-aware JS operations │ │
391// │ │ Promise handling & async writes │ │
392// │ └─────────────────────────────────────┘ │
393// └───────────────────────────────────────────┘
394// │
395// ▼
396// ┌───────────────────────────────────────────┐
397// │ JavaScript WritableStream │
398// │ │
399// │ • getWriter() │
400// │ • write(chunk) → Promise<void> │
401// │ • close() → Promise<void> │
402// │ • abort(reason) → Promise<void> │
403// │ • locked, state properties │
404// └───────────────────────────────────────────┘
405//
406class WritableStreamSinkKjAdapter final: public WritableSink {
407 public:
408 WritableStreamSinkKjAdapter(jsg::Lock& js, IoContext& ioContext, jsg::Ref<WritableStream> stream);
409 ~WritableStreamSinkKjAdapter() noexcept(false);
410 
411 // Attempts to write the given buffer to the underlying stream.
412 // The returned promise resolves once the write has completed.
413 // If the stream is closed, the returned promise rejects with
414 // an exception. If the stream errors, the returned promise
415 // rejects with the same exception. If the write fails, the
416 // returned promise rejects with the failure reason.
417 //
418 // Per the contract of write, it is the caller's responsibility
419 // to ensure that the adapter and buffer remain alive until
420 // the returned promise resolves.
421 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override;
422 
423 // Attempts to write the given pieces to the underlying stream.
424 // The returned promise resolves once the full write has completed.
425 // If the stream is closed, the returned promise rejects with
426 // an exception. If the stream errors, the returned promise
427 // rejects with the same exception. If the write fails, the
428 // returned promise rejects with the failure reason.
429 // Per the contract of write, it is the caller's responsibility
430 // to ensure that the adapter and buffers remain alive until
431 // the returned promise resolves.
432 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override;
433 
434 // Closes the underlying stream. The returned promise resolves
435 // once the stream is fully closed. If the stream is already
436 // closed, the returned promise resolves immediately. If the
437 // stream errors, the returned promise rejects with the same
438 // exception. If the close fails, the returned promise rejects
439 // with the failure reason.
440 kj::Promise<void> end() override;
441 
442 // Immediately interrupts existing pending writes and errors the stream.
443 // All pending or in-flight writes will be rejected with the given
444 // exception. If we are already in the errored state, this is a no-op
445 // and the exception is ignored. This change is immediate. Once in
446 // the errored state, no further writes or closes are allowed.
447 void abort(kj::Exception reason) override;
448 
449 // A WritableStreamSinkKjAdapter always (currently) uses identity encoding.
450 rpc::StreamEncoding disownEncodingResponsibility() override {
451 return rpc::StreamEncoding::IDENTITY;
452 }
453 
454 rpc::StreamEncoding getEncoding() override {
455 return rpc::StreamEncoding::IDENTITY;
456 }
457 
458 private:
459 struct Active;
460 KJ_DECLARE_NON_POLYMORPHIC(Active);
461 
462 struct KjClosed {
463 static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj;
464 };
465 
466 struct KjOpen {
467 static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj;
468 kj::Own<Active> active;
469 };
470 
471 // State machine for tracking writable sink adapter lifecycle:
472 // KjOpen -> KjClosed (normal close via end())
473 // KjOpen -> kj::Exception (error via abort() or write failure)
474 // KjClosed is terminal, kj::Exception is implicitly terminal via ErrorState.
475 using KjState = StateMachine<TerminalStates<KjClosed>,
476 ErrorState<kj::Exception>,
477 ActiveState<KjOpen>,
478 KjOpen,
479 KjClosed,
480 kj::Exception>;
481 KjState state;
482 kj::Rc<WeakRef<WritableStreamSinkKjAdapter>> selfRef;
483};
484 
485} // namespace workerd::api::streams