File
Blob: src/workerd/api/streams/readable.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 <kj/function.h> |
| 10 | #include <workerd/util/state-machine.h> |
| 11 | |
| 12 | namespace workerd::api { |
| 13 | |
| 14 | class ReadableStreamDefaultReader; |
| 15 | class ReadableStreamBYOBReader; |
| 16 | |
| 17 | class ReaderImpl final { |
| 18 | public: |
| 19 | ReaderImpl(ReadableStreamController::Reader& reader); |
| 20 | |
| 21 | ~ReaderImpl() noexcept(false); |
| 22 | |
| 23 | void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise); |
| 24 | |
| 25 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason); |
| 26 | |
| 27 | void detach(); |
| 28 | |
| 29 | jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed(); |
| 30 | |
| 31 | void lockToStream(jsg::Lock& js, ReadableStream& stream); |
| 32 | |
| 33 | jsg::Promise<ReadResult> read(jsg::Lock& js, |
| 34 | kj::Maybe<ReadableStreamController::ByobOptions> byobOptions); |
| 35 | |
| 36 | void releaseLock(jsg::Lock& js); |
| 37 | |
| 38 | void visitForGc(jsg::GcVisitor& visitor); |
| 39 | |
| 40 | kj::StringPtr jsgGetMemoryName() const; |
| 41 | size_t jsgGetMemorySelfSize() const; |
| 42 | void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 43 | |
| 44 | private: |
| 45 | struct Initial { |
| 46 | static constexpr kj::StringPtr NAME KJ_UNUSED = "initial"_kj; |
| 47 | }; |
| 48 | // While a Reader is attached to a ReadableStream, it holds a strong reference to the |
| 49 | // ReadableStream to prevent it from being GC'ed so long as the Reader is available. |
| 50 | // Once the reader is closed, released, or GC'ed the reference to the ReadableStream |
| 51 | // is cleared and the ReadableStream can be GC'ed if there are no other references to |
| 52 | // it being held anywhere. If the reader is still attached to the ReadableStream when |
| 53 | // it is destroyed, the ReadableStream's reference to the reader is cleared but the |
| 54 | // ReadableStream remains in the "reader locked" state, per the spec. |
| 55 | struct Attached { |
| 56 | static constexpr kj::StringPtr NAME KJ_UNUSED = "attached"_kj; |
| 57 | jsg::Ref<ReadableStream> stream; |
| 58 | }; |
| 59 | // Released: The user explicitly called releaseLock() to detach the reader from the stream. |
| 60 | // The stream remains usable and can be locked by a new reader. |
| 61 | struct Released { |
| 62 | static constexpr kj::StringPtr NAME KJ_UNUSED = "released"_kj; |
| 63 | }; |
| 64 | // Closed: The underlying stream ended (closed or errored) while the reader was attached. |
| 65 | // The stream is no longer usable. |
| 66 | struct Closed { |
| 67 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 68 | }; |
| 69 | |
| 70 | // State machine for ReaderImpl: |
| 71 | // Initial -> Attached (attach() called) |
| 72 | // Attached -> Closed (detach() called when stream closes) |
| 73 | // Attached -> Released (releaseLock() called) |
| 74 | // Closed and Released are terminal states. |
| 75 | // Initial is not terminal but most methods assert if called in this state. |
| 76 | using ReaderState = StateMachine<TerminalStates<Closed, Released>, |
| 77 | ActiveState<Attached>, |
| 78 | Initial, |
| 79 | Attached, |
| 80 | Closed, |
| 81 | Released>; |
| 82 | |
| 83 | kj::Maybe<IoContext&> ioContext; |
| 84 | ReadableStreamController::Reader& reader; |
| 85 | |
| 86 | ReaderState state; |
| 87 | |
| 88 | inline void assertAttachedOrTerminal() const { |
| 89 | KJ_ASSERT(!state.is<Initial>(), "this reader was never attached"); |
| 90 | } |
| 91 | kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<void>>> closedPromise; |
| 92 | |
| 93 | friend class ReadableStreamDefaultReader; |
| 94 | friend class ReadableStreamBYOBReader; |
| 95 | }; |
| 96 | |
| 97 | class ReadableStreamDefaultReader : public jsg::Object, |
| 98 | public ReadableStreamController::Reader { |
| 99 | public: |
| 100 | explicit ReadableStreamDefaultReader(); |
| 101 | |
| 102 | // JavaScript API |
| 103 | |
| 104 | static jsg::Ref<ReadableStreamDefaultReader> constructor( |
| 105 | jsg::Lock& js, jsg::Ref<ReadableStream> stream); |
| 106 | |
| 107 | jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed(); |
| 108 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason); |
| 109 | jsg::Promise<ReadResult> read(jsg::Lock& js); |
| 110 | void releaseLock(jsg::Lock& js); |
| 111 | |
| 112 | JSG_RESOURCE_TYPE(ReadableStreamDefaultReader, CompatibilityFlags::Reader flags) { |
| 113 | if (flags.getJsgPropertyOnPrototypeTemplate()) { |
| 114 | JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed); |
| 115 | } else { |
| 116 | JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed); |
| 117 | } |
| 118 | JSG_METHOD(cancel); |
| 119 | JSG_METHOD(read); |
| 120 | JSG_METHOD(releaseLock); |
| 121 | |
| 122 | JSG_TS_OVERRIDE(<R = any> { |
| 123 | read(): Promise<ReadableStreamReadResult<R>>; |
| 124 | }); |
| 125 | } |
| 126 | |
| 127 | // Internal API |
| 128 | |
| 129 | void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) override; |
| 130 | |
| 131 | void detach() override; |
| 132 | |
| 133 | void lockToStream(jsg::Lock& js, ReadableStream& stream); |
| 134 | |
| 135 | inline bool isByteOriented() const override { return false; } |
| 136 | |
| 137 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 138 | tracker.trackField("impl", impl); |
| 139 | } |
| 140 | |
| 141 | private: |
| 142 | ReaderImpl impl; |
| 143 | |
| 144 | void visitForGc(jsg::GcVisitor& visitor); |
| 145 | }; |
| 146 | |
| 147 | class ReadableStreamBYOBReader: public jsg::Object, |
| 148 | public ReadableStreamController::Reader { |
| 149 | public: |
| 150 | explicit ReadableStreamBYOBReader(); |
| 151 | |
| 152 | // JavaScript API |
| 153 | |
| 154 | static jsg::Ref<ReadableStreamBYOBReader> constructor( |
| 155 | jsg::Lock& js, |
| 156 | jsg::Ref<ReadableStream> stream); |
| 157 | |
| 158 | jsg::MemoizedIdentity<jsg::Promise<void>>& getClosed(); |
| 159 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason); |
| 160 | |
| 161 | struct ReadableStreamBYOBReaderReadOptions { |
| 162 | jsg::Optional<int> min; |
| 163 | JSG_STRUCT(min); |
| 164 | }; |
| 165 | |
| 166 | jsg::Promise<ReadResult> read(jsg::Lock& js, v8::Local<v8::ArrayBufferView> byobBuffer, |
| 167 | jsg::Optional<ReadableStreamBYOBReaderReadOptions> options = kj::none); |
| 168 | |
| 169 | // Non-standard extension so that reads can specify a minimum number of elements to read. It's a |
| 170 | // struct so that we could eventually add things like timeouts if we need to. Since there's no |
| 171 | // existing spec that's a leading contender, this is behind a different method name to avoid |
| 172 | // conflicts with any changes to `read`. Fewer than `minElements` may be returned if EOF is hit |
| 173 | // or the underlying stream is closed/errors out. In all cases the read result is either |
| 174 | // {value: theChunk, done: false} or {value: undefined, done: true} as with read. |
| 175 | // TODO(soon): Like fetch() and Cache.match(), readAtLeast() returns a promise for a V8 object. |
| 176 | jsg::Promise<ReadResult> readAtLeast(jsg::Lock& js, |
| 177 | int minElements, |
| 178 | v8::Local<v8::ArrayBufferView> byobBuffer); |
| 179 | |
| 180 | void releaseLock(jsg::Lock& js); |
| 181 | |
| 182 | JSG_RESOURCE_TYPE(ReadableStreamBYOBReader, CompatibilityFlags::Reader flags) { |
| 183 | if (flags.getJsgPropertyOnPrototypeTemplate()) { |
| 184 | JSG_READONLY_PROTOTYPE_PROPERTY(closed, getClosed); |
| 185 | } else { |
| 186 | JSG_READONLY_INSTANCE_PROPERTY(closed, getClosed); |
| 187 | } |
| 188 | JSG_METHOD(cancel); |
| 189 | JSG_METHOD(read); |
| 190 | JSG_METHOD(releaseLock); |
| 191 | |
| 192 | // Non-standard extension that should only apply to BYOB byte streams. |
| 193 | JSG_METHOD(readAtLeast); |
| 194 | |
| 195 | JSG_TS_OVERRIDE(ReadableStreamBYOBReader { |
| 196 | read<T extends ArrayBufferView>(view: T): Promise<ReadableStreamReadResult<T>>; |
| 197 | readAtLeast<T extends ArrayBufferView>(minElements: number, view: T): Promise<ReadableStreamReadResult<T>>; |
| 198 | }); |
| 199 | } |
| 200 | |
| 201 | // Internal API |
| 202 | |
| 203 | void attach( |
| 204 | ReadableStreamController& controller, |
| 205 | jsg::Promise<void> closedPromise) override; |
| 206 | |
| 207 | void detach() override; |
| 208 | |
| 209 | void lockToStream(jsg::Lock& js, ReadableStream& stream); |
| 210 | |
| 211 | inline bool isByteOriented() const override { return true; } |
| 212 | |
| 213 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 214 | tracker.trackField("impl", impl); |
| 215 | } |
| 216 | |
| 217 | private: |
| 218 | ReaderImpl impl; |
| 219 | |
| 220 | void visitForGc(jsg::GcVisitor& visitor); |
| 221 | }; |
| 222 | |
| 223 | // DrainingReader is a C++ only reader (not exposed to JavaScript) that performs |
| 224 | // draining reads. It locks the stream like standard readers but uses drainingRead() |
| 225 | // instead of regular read() to drain all synchronously available data at once. |
| 226 | // This is intended for optimized pipe operations. |
| 227 | class DrainingReader: public ReadableStreamController::Reader { |
| 228 | public: |
| 229 | explicit DrainingReader(); |
| 230 | |
| 231 | // Factory method to create and lock to a stream. Returns nullptr if stream is locked. |
| 232 | static kj::Maybe<kj::Own<DrainingReader>> create(jsg::Lock& js, ReadableStream& stream); |
| 233 | |
| 234 | virtual ~DrainingReader() noexcept(false); |
| 235 | |
| 236 | // Performs a draining read, returning all synchronously available data as bytes. |
| 237 | // The maxRead parameter is a soft limit - see ReadableStreamController::drainingRead. |
| 238 | jsg::Promise<DrainingReadResult> read(jsg::Lock& js, size_t maxRead = kj::maxValue); |
| 239 | |
| 240 | // Cancels the stream. |
| 241 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> maybeReason); |
| 242 | |
| 243 | // Releases the lock on the stream. |
| 244 | void releaseLock(jsg::Lock& js); |
| 245 | |
| 246 | // Returns whether this reader is still attached to a stream. |
| 247 | bool isAttached() const; |
| 248 | |
| 249 | // ReadableStreamController::Reader interface |
| 250 | void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) override; |
| 251 | void detach() override; |
| 252 | bool isByteOriented() const override { return false; } |
| 253 | |
| 254 | void visitForGc(jsg::GcVisitor& visitor); |
| 255 | |
| 256 | private: |
| 257 | struct Initial {}; |
| 258 | using Attached = jsg::Ref<ReadableStream>; |
| 259 | struct Released {}; |
| 260 | |
| 261 | kj::Maybe<IoContext&> ioContext; |
| 262 | kj::OneOf<Initial, Attached, StreamStates::Closed, Released> state = Initial(); |
| 263 | kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<void>>> closedPromise; |
| 264 | }; |
| 265 | |
| 266 | class ReadableStream: public jsg::Object { |
| 267 | private: |
| 268 | |
| 269 | struct AsyncIteratorState { |
| 270 | kj::Maybe<IoContext&> ioContext; |
| 271 | jsg::Ref<ReadableStreamDefaultReader> reader; |
| 272 | bool preventCancel; |
| 273 | }; |
| 274 | |
| 275 | static jsg::Promise<kj::Maybe<jsg::Value>> nextFunction( |
| 276 | jsg::Lock& js, |
| 277 | AsyncIteratorState& state); |
| 278 | |
| 279 | static jsg::Promise<void> returnFunction( |
| 280 | jsg::Lock& js, |
| 281 | AsyncIteratorState& state, |
| 282 | jsg::Optional<jsg::Value>& value); |
| 283 | |
| 284 | public: |
| 285 | explicit ReadableStream(IoContext& ioContext, |
| 286 | kj::Own<ReadableStreamSource> source); |
| 287 | |
| 288 | explicit ReadableStream(kj::Own<ReadableStreamController> controller); |
| 289 | |
| 290 | ReadableStreamController& getController(); |
| 291 | |
| 292 | jsg::Ref<ReadableStream> addRef(); |
| 293 | |
| 294 | bool isDisturbed(); |
| 295 | |
| 296 | // --------------------------------------------------------------------------- |
| 297 | // JS interface |
| 298 | |
| 299 | // Creates a new JS-backed ReadableStream using the provided source and strategy. |
| 300 | // We use v8::Local<v8::Object>'s here instead of jsg structs because we need |
| 301 | // to preserve the object references within the implementation. |
| 302 | static jsg::Ref<ReadableStream> constructor( |
| 303 | jsg::Lock& js, |
| 304 | jsg::Optional<UnderlyingSource> underlyingSource, |
| 305 | jsg::Optional<StreamQueuingStrategy> queuingStrategy); |
| 306 | |
| 307 | static jsg::Ref<ReadableStream> from(jsg::Lock& js, jsg::AsyncGenerator<jsg::Value> generator); |
| 308 | |
| 309 | bool isLocked(); |
| 310 | |
| 311 | // Closes the stream. All present and future read requests are fulfilled with successful empty |
| 312 | // results. `reason` will be passed to the underlying source's cancel algorithm -- if this |
| 313 | // readable stream is one side of a transform stream, then its cancel algorithm causes the |
| 314 | // transform's writable side to become errored with `reason`. |
| 315 | jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason); |
| 316 | |
| 317 | using Reader = kj::OneOf<jsg::Ref<ReadableStreamDefaultReader>, |
| 318 | jsg::Ref<ReadableStreamBYOBReader>>; |
| 319 | |
| 320 | struct GetReaderOptions { |
| 321 | jsg::Optional<kj::String> mode; // can be "byob" or undefined |
| 322 | |
| 323 | JSG_STRUCT(mode); |
| 324 | |
| 325 | JSG_STRUCT_TS_OVERRIDE({ mode: "byob" }); |
| 326 | // Intentionally required, so we can use `GetReaderOptions` directly in the |
| 327 | // `ReadableStream#getReader()` overload. |
| 328 | }; |
| 329 | |
| 330 | Reader getReader(jsg::Lock& js, jsg::Optional<GetReaderOptions> options); |
| 331 | |
| 332 | // Options specifically for the values() function. |
| 333 | struct ValuesOptions { |
| 334 | jsg::Optional<bool> preventCancel = false; |
| 335 | JSG_STRUCT(preventCancel); |
| 336 | }; |
| 337 | |
| 338 | JSG_ASYNC_ITERATOR_WITH_OPTIONS(ReadableStreamAsyncIterator, |
| 339 | values, |
| 340 | jsg::Value, |
| 341 | AsyncIteratorState, |
| 342 | nextFunction, |
| 343 | returnFunction, |
| 344 | ValuesOptions); |
| 345 | struct Transform { |
| 346 | jsg::Ref<ReadableStream> readable; |
| 347 | jsg::Ref<WritableStream> writable; |
| 348 | |
| 349 | JSG_STRUCT(readable, writable); |
| 350 | JSG_STRUCT_TS_OVERRIDE(ReadableWritablePair<R = any, W = any> { |
| 351 | readable: ReadableStream<R>; |
| 352 | writable: WritableStream<W>; |
| 353 | }); |
| 354 | }; |
| 355 | |
| 356 | jsg::Ref<ReadableStream> pipeThrough( |
| 357 | jsg::Lock& js, |
| 358 | Transform transform, |
| 359 | jsg::Optional<PipeToOptions> options); |
| 360 | |
| 361 | jsg::Promise<void> pipeTo( |
| 362 | jsg::Lock& js, |
| 363 | jsg::Ref<WritableStream> destination, |
| 364 | jsg::Optional<PipeToOptions> options); |
| 365 | |
| 366 | // Locks the stream and returns a pair of two new ReadableStreams, each of which read the same |
| 367 | // data as this ReadableStream would. |
| 368 | kj::Array<jsg::Ref<ReadableStream>> tee(jsg::Lock& js); |
| 369 | |
| 370 | jsg::JsString inspectState(jsg::Lock& js); |
| 371 | bool inspectSupportsBYOB(); |
| 372 | jsg::Optional<uint64_t> inspectLength(); |
| 373 | |
| 374 | JSG_RESOURCE_TYPE(ReadableStream, CompatibilityFlags::Reader flags) { |
| 375 | if (flags.getJsgPropertyOnPrototypeTemplate()) { |
| 376 | JSG_READONLY_PROTOTYPE_PROPERTY(locked, isLocked); |
| 377 | } else { |
| 378 | JSG_READONLY_INSTANCE_PROPERTY(locked, isLocked); |
| 379 | } |
| 380 | JSG_METHOD(cancel); |
| 381 | JSG_METHOD(getReader); |
| 382 | JSG_METHOD(pipeThrough); |
| 383 | JSG_METHOD(pipeTo); |
| 384 | JSG_METHOD(tee); |
| 385 | JSG_METHOD(values); |
| 386 | JSG_STATIC_METHOD(from); |
| 387 | |
| 388 | JSG_INSPECT_PROPERTY(state, inspectState); |
| 389 | JSG_INSPECT_PROPERTY(supportsBYOB, inspectSupportsBYOB); |
| 390 | JSG_INSPECT_PROPERTY(length, inspectLength); |
| 391 | |
| 392 | JSG_ASYNC_ITERABLE(values); |
| 393 | |
| 394 | if (flags.getJsgPropertyOnPrototypeTemplate()) { |
| 395 | JSG_TS_DEFINE(interface ReadableStream<R = any> { |
| 396 | get locked(): boolean; |
| 397 | |
| 398 | cancel(reason?: any): Promise<void>; |
| 399 | |
| 400 | getReader(): ReadableStreamDefaultReader<R>; |
| 401 | getReader(options: ReadableStreamGetReaderOptions): ReadableStreamBYOBReader; |
| 402 | |
| 403 | pipeThrough<T>(transform: ReadableWritablePair<T, R>, options?: StreamPipeOptions): ReadableStream<T>; |
| 404 | pipeTo(destination: WritableStream<R>, options?: StreamPipeOptions): Promise<void>; |
| 405 | |
| 406 | tee(): [ReadableStream<R>, ReadableStream<R>]; |
| 407 | |
| 408 | values(options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>; |
| 409 | [Symbol.asyncIterator](options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>; |
| 410 | }); |
| 411 | } else { |
| 412 | JSG_TS_DEFINE(interface ReadableStream<R = any> { |
| 413 | readonly locked: boolean; |
| 414 | |
| 415 | cancel(reason?: any): Promise<void>; |
| 416 | |
| 417 | getReader(): ReadableStreamDefaultReader<R>; |
| 418 | getReader(options: ReadableStreamGetReaderOptions): ReadableStreamBYOBReader; |
| 419 | |
| 420 | pipeThrough<T>(transform: ReadableWritablePair<T, R>, options?: StreamPipeOptions): ReadableStream<T>; |
| 421 | pipeTo(destination: WritableStream<R>, options?: StreamPipeOptions): Promise<void>; |
| 422 | |
| 423 | tee(): [ReadableStream<R>, ReadableStream<R>]; |
| 424 | |
| 425 | values(options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>; |
| 426 | [Symbol.asyncIterator](options?: ReadableStreamValuesOptions): AsyncIterableIterator<R>; |
| 427 | }); |
| 428 | } |
| 429 | // Replace ReadableStream class with an interface and const, so we can have |
| 430 | // two constructors with differing type parameters for byte-oriented and |
| 431 | // value-oriented streams. |
| 432 | JSG_TS_OVERRIDE(const ReadableStream: { |
| 433 | prototype: ReadableStream; |
| 434 | new (underlyingSource: UnderlyingByteSource, strategy?: QueuingStrategy<Uint8Array>): ReadableStream<Uint8Array>; |
| 435 | new <R = any>(underlyingSource?: UnderlyingSource<R>, strategy?: QueuingStrategy<R>): ReadableStream<R>; |
| 436 | }); |
| 437 | } |
| 438 | |
| 439 | // Detaches this ReadableStream from its underlying controller state, returning a |
| 440 | // new ReadableStream instance that takes over the underlying state. This is used to |
| 441 | // support the "create a proxy" of a ReadableStream algorithm in the streams spec |
| 442 | // (see https://streams.spec.whatwg.org/#readablestream-create-a-proxy). In that |
| 443 | // algorithm, it says to create a proxy of a stream by creating a new TransformStream |
| 444 | // and piping the original through it. The readable side of the created transform |
| 445 | // becomes the proxy. That is quite inefficient so instead, we create a new |
| 446 | // ReadableStream that will take over ownership of the internal state of this one, |
| 447 | // leaving this ReadableStream locked and disturbed so that it is no longer usable. |
| 448 | // The name "detach" here is used in the sense of "detaching the internal state". |
| 449 | jsg::Ref<ReadableStream> detach(jsg::Lock& js, bool ignoreDisturbed=false); |
| 450 | |
| 451 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding); |
| 452 | |
| 453 | // A potentially optimized version of pipe that sends this stream's data to the given |
| 454 | // sink. The entire stream is consumed. The ReadableStream will be left locked and |
| 455 | // disturbed and the DeferredProxy returned will take over ownership of the internal |
| 456 | // state of the readable. |
| 457 | kj::Promise<DeferredProxy<void>> pumpTo(jsg::Lock& js, |
| 458 | kj::Own<WritableStreamSink> sink, |
| 459 | bool end); |
| 460 | |
| 461 | // Initializes signalling mechanism for EOF detection. Returns a promise that will resolve when |
| 462 | // EOF is reached. |
| 463 | // |
| 464 | // This method should only be called once. |
| 465 | jsg::Promise<void> onEof(jsg::Lock& js); |
| 466 | |
| 467 | // Used by ReadableStreamInternalController to signal EOF being reached. Can be called even if |
| 468 | // `onEof` wasn't called. |
| 469 | void signalEof(jsg::Lock& js); |
| 470 | |
| 471 | void serialize(jsg::Lock& js, jsg::Serializer& serializer); |
| 472 | static jsg::Ref<ReadableStream> deserialize( |
| 473 | jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer); |
| 474 | |
| 475 | JSG_SERIALIZABLE(rpc::SerializationTag::READABLE_STREAM); |
| 476 | |
| 477 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 478 | |
| 479 | private: |
| 480 | kj::Maybe<IoContext&> ioContext; |
| 481 | kj::Own<ReadableStreamController> controller; |
| 482 | |
| 483 | // Used to signal when this ReadableStream reads EOF. This signal is required for TCP sockets. |
| 484 | kj::Maybe<jsg::PromiseResolverPair<void>> eofResolverPair; |
| 485 | |
| 486 | void visitForGc(jsg::GcVisitor& visitor); |
| 487 | }; |
| 488 | |
| 489 | struct QueuingStrategyInit { |
| 490 | double highWaterMark; |
| 491 | JSG_STRUCT(highWaterMark); |
| 492 | }; |
| 493 | |
| 494 | using QueuingStrategySizeFunction = |
| 495 | jsg::Optional<uint32_t>(jsg::Optional<v8::Local<v8::Value>>); |
| 496 | |
| 497 | // Utility class defined by the streams spec that uses byteLength to calculate |
| 498 | // backpressure changes. |
| 499 | class ByteLengthQueuingStrategy: public jsg::Object { |
| 500 | public: |
| 501 | ByteLengthQueuingStrategy(QueuingStrategyInit init) : init(init) {} |
| 502 | |
| 503 | static jsg::Ref<ByteLengthQueuingStrategy> constructor(jsg::Lock& js, QueuingStrategyInit init) { |
| 504 | return js.alloc<ByteLengthQueuingStrategy>(init); |
| 505 | } |
| 506 | |
| 507 | double getHighWaterMark() const { return init.highWaterMark; } |
| 508 | |
| 509 | jsg::Function<QueuingStrategySizeFunction> getSize() const { return &size; } |
| 510 | |
| 511 | JSG_RESOURCE_TYPE(ByteLengthQueuingStrategy) { |
| 512 | JSG_READONLY_PROTOTYPE_PROPERTY(highWaterMark, getHighWaterMark); |
| 513 | JSG_READONLY_PROTOTYPE_PROPERTY(size, getSize); |
| 514 | |
| 515 | // QueuingStrategy requires the result of the size function to be defined |
| 516 | JSG_TS_OVERRIDE(implements QueuingStrategy<ArrayBufferView> { |
| 517 | get size(): (chunk?: any) => number; |
| 518 | }); |
| 519 | } |
| 520 | |
| 521 | private: |
| 522 | static jsg::Optional<uint32_t> size(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>>); |
| 523 | |
| 524 | QueuingStrategyInit init; |
| 525 | }; |
| 526 | |
| 527 | // Utility class defined by the streams spec that uses a fixed value of 1 to calculate |
| 528 | // backpressure change |
| 529 | class CountQueuingStrategy: public jsg::Object { |
| 530 | public: |
| 531 | CountQueuingStrategy(QueuingStrategyInit init) : init(init) {} |
| 532 | |
| 533 | static jsg::Ref<CountQueuingStrategy> constructor(jsg::Lock& js, QueuingStrategyInit init) { |
| 534 | return js.alloc<CountQueuingStrategy>(init); |
| 535 | } |
| 536 | |
| 537 | double getHighWaterMark() const { return init.highWaterMark; } |
| 538 | |
| 539 | jsg::Function<QueuingStrategySizeFunction> getSize() const { return &size; } |
| 540 | |
| 541 | JSG_RESOURCE_TYPE(CountQueuingStrategy) { |
| 542 | JSG_READONLY_PROTOTYPE_PROPERTY(highWaterMark, getHighWaterMark); |
| 543 | JSG_READONLY_PROTOTYPE_PROPERTY(size, getSize); |
| 544 | |
| 545 | // QueuingStrategy requires the result of the size function to be defined |
| 546 | JSG_TS_OVERRIDE(implements QueuingStrategy { |
| 547 | get size(): (chunk?: any) => number; |
| 548 | }); |
| 549 | } |
| 550 | |
| 551 | private: |
| 552 | static jsg::Optional<uint32_t> size(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>>) { |
| 553 | return 1; |
| 554 | } |
| 555 | |
| 556 | QueuingStrategyInit init; |
| 557 | }; |
| 558 | |
| 559 | } // namespace workerd::api |