Skip to content
File

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

cpp501 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#include "writable.h"
9 
10#include <workerd/io/io-context.h>
11#include <workerd/io/observer.h>
12#include <workerd/util/ring-buffer.h>
13#include <workerd/util/state-machine.h>
14 
15#include <kj/refcount.h>
16 
17namespace workerd::api {
18 
19// =======================================================================================
20// The ReadableStreamInternalController and WritableStreamInternalController provide the
21// internal (original) implementation of the ReadableStream/WritableStream objects and are
22// each backed by the ReadableStreamSource and WritableStreamSink respectively. Every stream
23// implementation that originates from *within* the Workers runtime will use these.
24//
25// It is important to understand that the behavior of these are not entirely compliant with
26// the streams specification.
27 
28// The ReadableStreamInternalController is always in one of three states: Readable, Closed,
29// or Errored. When the state is Readable, the controller has an associated ReadableStreamSource.
30// When the state is Errored, the ReadableStreamSource has been released and the controller
31// stores a jsg::Value with whatever value was used to error. When Closed, the
32// ReadableStreamSource has been released.
33 
34// Likewise, the WritableStreamInternalController is always either Writable, Closed, or Errored.
35// When the state is Writable, the controller has an associated WritableStreamSink. In either of
36// the other two states, the sink has been released.
37 
38class WritableStreamInternalController;
39 
40class ReadableStreamInternalController: public ReadableStreamController {
41 public:
42 using Readable = IoOwn<ReadableStreamSource>;
43 
44 explicit ReadableStreamInternalController(StreamStates::Closed closed)
45 : state(State::create<StreamStates::Closed>()) {}
46 explicit ReadableStreamInternalController(StreamStates::Errored errored)
47 : state(State::create<StreamStates::Errored>(kj::mv(errored))) {}
48 explicit ReadableStreamInternalController(Readable readable)
49 : state(State::create<Readable>(kj::mv(readable))) {}
50 
51 KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamInternalController);
52 
53 ~ReadableStreamInternalController() noexcept(false) override;
54 
55 void setOwnerRef(ReadableStream& stream) override {
56 owner = stream;
57 }
58 
59 jsg::Ref<ReadableStream> addRef() override;
60 
61 bool isByteOriented() const override {
62 return true;
63 }
64 
65 kj::Maybe<jsg::Promise<ReadResult>> read(
66 jsg::Lock& js, kj::Maybe<ByobOptions> byobOptions) override;
67 
68 kj::Maybe<jsg::Promise<DrainingReadResult>> drainingRead(
69 jsg::Lock& js, size_t maxRead = kj::maxValue) override;
70 
71 jsg::Promise<void> pipeTo(
72 jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) override;
73 
74 jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) override;
75 
76 Tee tee(jsg::Lock& js) override;
77 
78 kj::Maybe<kj::Own<ReadableStreamSource>> removeSource(
79 jsg::Lock& js, bool ignoreDisturbed = false);
80 
81 bool isClosedOrErrored() const override {
82 return state.is<StreamStates::Closed>() || state.is<StreamStates::Errored>();
83 }
84 
85 bool isClosed() const override {
86 return state.is<StreamStates::Closed>();
87 }
88 
89 bool isDisturbed() override {
90 return disturbed;
91 }
92 
93 bool isLockedToReader() const override {
94 return !readState.is<Unlocked>();
95 }
96 
97 bool lockReader(jsg::Lock& js, Reader& reader) override;
98 
99 void releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) override;
100 // See the comment for releaseReader in common.h for details on the use of maybeJs
101 
102 kj::Maybe<PipeController&> tryPipeLock() override;
103 
104 void visitForGc(jsg::GcVisitor& visitor) override;
105 
106 jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit) override;
107 jsg::Promise<kj::String> readAllText(jsg::Lock& js, uint64_t limit) override;
108 
109 kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) override;
110 
111 kj::Promise<DeferredProxy<void>> pumpTo(
112 jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) override;
113 
114 StreamEncoding getPreferredEncoding() override;
115 
116 kj::Own<ReadableStreamController> detach(jsg::Lock& js, bool ignoreDisturbed) override;
117 
118 void setPendingClosure() override {
119 isPendingClosure = true;
120 }
121 
122 kj::StringPtr jsgGetMemoryName() const override;
123 size_t jsgGetMemorySelfSize() const override;
124 void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override;
125 
126 private:
127 void doCancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason);
128 void doClose(jsg::Lock& js);
129 void doError(jsg::Lock& js, v8::Local<v8::Value> reason);
130 
131 class PipeLocked: public PipeController {
132 public:
133 static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj;
134 PipeLocked(ReadableStreamInternalController& inner): inner(inner) {}
135 
136 bool isClosed() override;
137 
138 kj::Maybe<v8::Local<v8::Value>> tryGetErrored(jsg::Lock& js) override;
139 
140 void cancel(jsg::Lock& js, v8::Local<v8::Value> reason) override;
141 
142 void close(jsg::Lock& js) override;
143 
144 void error(jsg::Lock& js, v8::Local<v8::Value> reason) override;
145 
146 void release(jsg::Lock& js, kj::Maybe<v8::Local<v8::Value>> maybeError = kj::none) override;
147 
148 kj::Maybe<kj::Promise<void>> tryPumpTo(WritableStreamSink& sink, bool end) override;
149 
150 jsg::Promise<ReadResult> read(jsg::Lock& js) override;
151 
152 private:
153 ReadableStreamInternalController& inner;
154 };
155 
156 kj::Maybe<ReadableStream&> owner;
157 
158 // State machine for ReadableStreamInternalController:
159 // Closed is terminal, Errored is implicitly terminal via ErrorState.
160 // Readable is the active state (stream has data).
161 using State = StateMachine<TerminalStates<StreamStates::Closed>,
162 ErrorState<StreamStates::Errored>,
163 ActiveState<Readable>,
164 StreamStates::Closed,
165 StreamStates::Errored,
166 Readable>;
167 State state;
168 
169 // Lock state machine for ReadableStreamInternalController:
170 // All states can transition to any other state (no terminal states).
171 // Unlocked -> Locked (removeSink() or pumpTo() called)
172 // Unlocked -> ReaderLocked (lockReader() called)
173 // Unlocked -> PipeLocked (tryPipeLock() called)
174 // ReaderLocked -> Unlocked (releaseReader() called)
175 // PipeLocked -> Unlocked (release() or doClose/doError called)
176 // Locked -> (remains until stream is done)
177 using ReadLockState = StateMachine<Unlocked, Locked, PipeLocked, ReaderLocked>;
178 ReadLockState readState = ReadLockState::create<Unlocked>();
179 bool disturbed = false;
180 bool readPending = false;
181 
182 // Used by Sockets code to signal to the ReadableStream that it should error when read from
183 // because the socket is currently being closed.
184 bool isPendingClosure = false;
185 
186 friend class ReadableStream;
187 friend class WritableStreamInternalController;
188 friend class PipeLocked;
189};
190 
191class WritableStreamInternalController: public WritableStreamController {
192 public:
193 struct Writable {
194 kj::Own<WritableStreamSink> sink;
195 kj::Canceler canceler;
196 Writable(kj::Own<WritableStreamSink> sink): sink(kj::mv(sink)) {}
197 void abort(kj::Exception&& ex);
198 };
199 
200 explicit WritableStreamInternalController(StreamStates::Closed closed)
201 : state(State::create<StreamStates::Closed>()) {}
202 explicit WritableStreamInternalController(StreamStates::Errored errored)
203 : state(State::create<StreamStates::Errored>(kj::mv(errored))) {}
204 explicit WritableStreamInternalController(kj::Own<WritableStreamSink> writable,
205 kj::Maybe<kj::Own<ByteStreamObserver>> observer,
206 kj::Maybe<uint64_t> maybeHighWaterMark = kj::none,
207 kj::Maybe<jsg::Promise<void>> maybeClosureWaitable = kj::none)
208 : state(State::create<IoOwn<Writable>>(
209 IoContext::current().addObject(kj::heap<Writable>(kj::mv(writable))))),
210 observer(kj::mv(observer)),
211 maybeHighWaterMark(maybeHighWaterMark),
212 maybeClosureWaitable(kj::mv(maybeClosureWaitable)) {}
213 
214 WritableStreamInternalController(WritableStreamInternalController&& other) = default;
215 WritableStreamInternalController& operator=(WritableStreamInternalController&& other) = default;
216 
217 ~WritableStreamInternalController() noexcept(false) override;
218 
219 void setOwnerRef(WritableStream& stream) override {
220 owner = stream;
221 }
222 
223 jsg::Ref<WritableStream> addRef() override;
224 
225 jsg::Promise<void> write(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) override;
226 
227 jsg::Promise<void> close(jsg::Lock& js, bool markAsHandled = false) override;
228 
229 jsg::Promise<void> flush(jsg::Lock& js, bool markAsHandled = false) override;
230 
231 jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) override;
232 
233 kj::Maybe<jsg::Promise<void>> tryPipeFrom(
234 jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) override;
235 
236 kj::Maybe<kj::Own<WritableStreamSink>> removeSink(jsg::Lock& js) override;
237 void detach(jsg::Lock& js) override;
238 
239 kj::Maybe<int> getDesiredSize() override;
240 
241 bool isLockedToWriter() const override {
242 return !writeState.is<Unlocked>();
243 }
244 
245 bool lockWriter(jsg::Lock& js, Writer& writer) override;
246 
247 void releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) override;
248 // See the comment for releaseWriter in common.h for details on the use of maybeJs
249 
250 kj::Maybe<v8::Local<v8::Value>> isErroring(jsg::Lock& js) override {
251 // TODO(later): The internal controller has no concept of an "erroring"
252 // state, so for now we just return kj::none here.
253 return kj::none;
254 }
255 
256 void visitForGc(jsg::GcVisitor& visitor) override;
257 
258 void setHighWaterMark(uint64_t highWaterMark);
259 
260 bool isClosedOrClosing() override;
261 bool isPiping();
262 bool isErrored() override;
263 
264 inline bool isByteOriented() const override {
265 return true;
266 }
267 
268 void setPendingClosure() override {
269 isPendingClosure = true;
270 }
271 
272 kj::StringPtr jsgGetMemoryName() const override;
273 size_t jsgGetMemorySelfSize() const override;
274 void jsgGetMemoryInfo(jsg::MemoryTracker& info) const override;
275 
276 private:
277 struct AbortOptions {
278 bool reject = false;
279 bool handled = false;
280 };
281 
282 jsg::Promise<void> doAbort(jsg::Lock& js,
283 v8::Local<v8::Value> reason,
284 AbortOptions options = {.reject = false, .handled = false});
285 void doClose(jsg::Lock& js);
286 void doError(jsg::Lock& js, v8::Local<v8::Value> reason);
287 void ensureWriting(jsg::Lock& js);
288 jsg::Promise<void> writeLoop(jsg::Lock& js, IoContext& ioContext);
289 jsg::Promise<void> writeLoopAfterFrontOutputLock(jsg::Lock& js);
290 
291 void drain(jsg::Lock& js, v8::Local<v8::Value> reason);
292 void finishClose(jsg::Lock& js);
293 void finishError(jsg::Lock& js, v8::Local<v8::Value> reason);
294 jsg::Promise<void> closeImpl(jsg::Lock& js, bool markAsHandled);
295 
296 struct PipeLocked {
297 static constexpr kj::StringPtr NAME KJ_UNUSED = "pipe-locked"_kj;
298 ReadableStream& ref;
299 };
300 
301 kj::Maybe<WritableStream&> owner;
302 
303 // State machine for WritableStreamInternalController:
304 // Closed is terminal, Errored is implicitly terminal via ErrorState.
305 // IoOwn<Writable> is the active state (stream is writable).
306 using State = StateMachine<TerminalStates<StreamStates::Closed>,
307 ErrorState<StreamStates::Errored>,
308 ActiveState<IoOwn<Writable>>,
309 StreamStates::Closed,
310 StreamStates::Errored,
311 IoOwn<Writable>>;
312 State state;
313 
314 // Lock state machine for WritableStreamInternalController:
315 // All states can transition to any other state (no terminal states).
316 // Unlocked -> Locked (removeSink() or detach() called)
317 // Unlocked -> WriterLocked (lockWriter() called)
318 // Unlocked -> PipeLocked (tryPipeFrom() called)
319 // WriterLocked -> Unlocked (releaseWriter() called)
320 // WriterLocked -> Locked (doClose/doError called - stream closed but writer still attached)
321 // PipeLocked -> Unlocked (pipe completes)
322 using WriteLockState = StateMachine<Unlocked, Locked, PipeLocked, WriterLocked>;
323 WriteLockState writeState = WriteLockState::create<Unlocked>();
324 
325 kj::Maybe<kj::Own<ByteStreamObserver>> observer;
326 
327 kj::Maybe<kj::Own<PendingAbort>> maybePendingAbort;
328 
329 uint64_t currentWriteBufferSize = 0;
330 
331 // The highWaterMark is the total amount of data currently buffered in
332 // the controller waiting to be flushed out to the underlying WritableStreamSink.
333 // It is used to implement backpressure signaling using desiredSize and the ready
334 // promise on the writer.
335 kj::Maybe<uint64_t> maybeHighWaterMark;
336 
337 // Used by Sockets code to ensure the connection is established before the associated
338 // WritableStream is closed.
339 kj::Maybe<jsg::Promise<void>> maybeClosureWaitable;
340 bool waitingOnClosureWritableAlready = false;
341 
342 // Used by Sockets code to signal to the WritableStream that it should error when written to
343 // because the socket is currently being closed.
344 bool isPendingClosure = false;
345 
346 void adjustWriteBufferSize(jsg::Lock& js, int64_t amount);
347 void updateBackpressure(jsg::Lock& js, bool backpressure);
348 
349 struct Write {
350 kj::Maybe<jsg::Promise<void>::Resolver> promise;
351 size_t totalBytes;
352 kj::Array<kj::byte> ownBytes;
353 kj::ArrayPtr<const kj::byte> bytes;
354 
355 JSG_MEMORY_INFO(Write) {
356 tracker.trackField("resolver", promise);
357 if (ownBytes != nullptr) {
358 tracker.trackFieldWithSize("backing", totalBytes);
359 }
360 }
361 };
362 struct Close {
363 kj::Maybe<jsg::Promise<void>::Resolver> promise;
364 JSG_MEMORY_INFO(Close) {
365 tracker.trackField("promise", promise);
366 }
367 };
368 struct Flush {
369 kj::Maybe<jsg::Promise<void>::Resolver> promise;
370 JSG_MEMORY_INFO(Flush) {
371 tracker.trackField("promise", promise);
372 }
373 };
374 struct Pipe {
375 // PipeState is ref-counted so that it can be safely captured by lambdas in pipeLoop().
376 // When drain() destroys the Pipe, the state survives as long as pending callbacks need it.
377 // The `aborted` flag is set when the Pipe is destroyed.
378 struct State: public kj::Refcounted {
379 WritableStreamInternalController& parent;
380 ReadableStreamController::PipeController& source;
381 kj::Maybe<jsg::Promise<void>::Resolver> promise;
382 kj::Maybe<jsg::Ref<AbortSignal>> maybeSignal;
383 
384 bool preventAbort;
385 bool preventClose;
386 bool preventCancel;
387 
388 // True when the Pipe is being destroyed
389 bool aborted = false;
390 
391 State(WritableStreamInternalController& parent,
392 ReadableStreamController::PipeController& source,
393 kj::Maybe<jsg::Promise<void>::Resolver> promise,
394 bool preventAbort,
395 bool preventClose,
396 bool preventCancel,
397 kj::Maybe<jsg::Ref<AbortSignal>> maybeSignal)
398 : parent(parent),
399 source(source),
400 promise(kj::mv(promise)),
401 maybeSignal(kj::mv(maybeSignal)),
402 preventAbort(preventAbort),
403 preventClose(preventClose),
404 preventCancel(preventCancel) {}
405 
406 bool checkSignal(jsg::Lock& js);
407 jsg::Promise<void> pipeLoop(jsg::Lock& js);
408 jsg::Promise<void> write(v8::Local<v8::Value> value);
409 
410 JSG_MEMORY_INFO(State) {
411 tracker.trackField("resolver", promise);
412 tracker.trackField("signal", maybeSignal);
413 }
414 };
415 
416 kj::Own<State> state;
417 
418 Pipe(WritableStreamInternalController& parent,
419 ReadableStreamController::PipeController& source,
420 kj::Maybe<jsg::Promise<void>::Resolver> promise,
421 bool preventAbort,
422 bool preventClose,
423 bool preventCancel,
424 kj::Maybe<jsg::Ref<AbortSignal>> maybeSignal)
425 : state(kj::refcounted<State>(parent,
426 source,
427 kj::mv(promise),
428 preventAbort,
429 preventClose,
430 preventCancel,
431 kj::mv(maybeSignal))) {}
432 
433 ~Pipe() noexcept(false) {
434 state->aborted = true;
435 }
436 
437 WritableStreamInternalController& parent() {
438 return state->parent;
439 }
440 ReadableStreamController::PipeController& source() {
441 return state->source;
442 }
443 kj::Maybe<jsg::Promise<void>::Resolver>& promise() {
444 return state->promise;
445 }
446 bool preventAbort() const {
447 return state->preventAbort;
448 }
449 bool preventClose() const {
450 return state->preventClose;
451 }
452 bool preventCancel() const {
453 return state->preventCancel;
454 }
455 kj::Maybe<jsg::Ref<AbortSignal>>& maybeSignal() {
456 return state->maybeSignal;
457 }
458 
459 bool checkSignal(jsg::Lock& js) {
460 return state->checkSignal(js);
461 }
462 jsg::Promise<void> pipeLoop(jsg::Lock& js) {
463 return state->pipeLoop(js);
464 }
465 jsg::Promise<void> write(v8::Local<v8::Value> value) {
466 return state->write(value);
467 }
468 
469 JSG_MEMORY_INFO(Pipe) {
470 tracker.trackField("state", state);
471 }
472 };
473 struct WriteEvent {
474 kj::Maybe<IoOwn<kj::Promise<void>>> outputLock; // must wait for this before actually writing
475 kj::OneOf<kj::Own<Write>, kj::Own<Pipe>, kj::Own<Close>, kj::Own<Flush>> event;
476 
477 JSG_MEMORY_INFO(WriteEvent) {
478 if (outputLock != kj::none) {
479 tracker.trackFieldWithSize("outputLock", sizeof(IoOwn<kj::Promise<void>>));
480 }
481 KJ_SWITCH_ONEOF(event) {
482 KJ_CASE_ONEOF(w, kj::Own<Write>) {
483 tracker.trackField("inner", w);
484 }
485 KJ_CASE_ONEOF(p, kj::Own<Pipe>) {
486 tracker.trackField("inner", p);
487 }
488 KJ_CASE_ONEOF(c, kj::Own<Close>) {
489 tracker.trackField("inner", c);
490 }
491 KJ_CASE_ONEOF(f, kj::Own<Flush>) {
492 tracker.trackField("inner", f);
493 }
494 }
495 }
496 };
497 
498 RingBuffer<WriteEvent, 8> queue;
499};
500} // namespace workerd::api