File
Blob: src/node/internal/streams_readable.js
| 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 | // Copyright Joyent, Inc. and other Node contributors. |
| 6 | // |
| 7 | // Permission is hereby granted, free of charge, to any person obtaining a |
| 8 | // copy of this software and associated documentation files (the |
| 9 | // "Software"), to deal in the Software without restriction, including |
| 10 | // without limitation the rights to use, copy, modify, merge, publish, |
| 11 | // distribute, sublicense, and/or sell copies of the Software, and to permit |
| 12 | // persons to whom the Software is furnished to do so, subject to the |
| 13 | // following conditions: |
| 14 | // |
| 15 | // The above copyright notice and this permission notice shall be included |
| 16 | // in all copies or substantial portions of the Software. |
| 17 | // |
| 18 | // THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS |
| 19 | // OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF |
| 20 | // MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN |
| 21 | // NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, |
| 22 | // DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR |
| 23 | // OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE |
| 24 | // USE OR OTHER DEALINGS IN THE SOFTWARE. |
| 25 | |
| 26 | import { |
| 27 | kState, |
| 28 | // bitfields |
| 29 | kObjectMode, |
| 30 | kErrorEmitted, |
| 31 | kAutoDestroy, |
| 32 | kEmitClose, |
| 33 | kDestroyed, |
| 34 | kClosed, |
| 35 | kCloseEmitted, |
| 36 | kErrored, |
| 37 | kConstructed, |
| 38 | kOnConstructed, |
| 39 | isDestroyed, |
| 40 | isReadable, |
| 41 | isReadableStream, |
| 42 | handleKnownInternalErrors, |
| 43 | } from 'node-internal:streams_util'; |
| 44 | import { nextTick } from 'node-internal:internal_process'; |
| 45 | import { |
| 46 | destroy, |
| 47 | undestroy, |
| 48 | errorOrDestroy, |
| 49 | destroyer, |
| 50 | construct, |
| 51 | } from 'node-internal:streams_destroy'; |
| 52 | import { eos, finished, nop } from 'node-internal:streams_end_of_stream'; |
| 53 | import { |
| 54 | getHighWaterMark, |
| 55 | getDefaultHighWaterMark, |
| 56 | } from 'node-internal:streams_state'; |
| 57 | import { addAbortSignal } from 'node-internal:streams_add_abort_signal'; |
| 58 | import { EventEmitter } from 'node-internal:events'; |
| 59 | import { Stream } from 'node-internal:streams_legacy'; |
| 60 | import { Buffer } from 'node-internal:internal_buffer'; |
| 61 | |
| 62 | import { |
| 63 | AbortError, |
| 64 | aggregateTwoErrors, |
| 65 | ERR_INVALID_ARG_TYPE, |
| 66 | ERR_INVALID_ARG_VALUE, |
| 67 | ERR_METHOD_NOT_IMPLEMENTED, |
| 68 | ERR_MISSING_ARGS, |
| 69 | ERR_OUT_OF_RANGE, |
| 70 | ERR_STREAM_PUSH_AFTER_EOF, |
| 71 | ERR_STREAM_UNSHIFT_AFTER_END_EVENT, |
| 72 | ERR_STREAM_NULL_VALUES, |
| 73 | ERR_UNKNOWN_ENCODING, |
| 74 | } from 'node-internal:internal_errors'; |
| 75 | |
| 76 | import { |
| 77 | validateObject, |
| 78 | validateAbortSignal, |
| 79 | validateBoolean, |
| 80 | validateInteger, |
| 81 | } from 'node-internal:validators'; |
| 82 | |
| 83 | import { StringDecoder } from 'node-internal:internal_stringdecoder'; |
| 84 | |
| 85 | const streamsNodejsV24Compat = |
| 86 | Cloudflare.compatibilityFlags.enable_streams_nodejs_v24_compat; |
| 87 | |
| 88 | const kErroredValue = Symbol('kErroredValue'); |
| 89 | const kDefaultEncodingValue = Symbol('kDefaultEncodingValue'); |
| 90 | const kDecoderValue = Symbol('kDecoderValue'); |
| 91 | const kEncodingValue = Symbol('kEncodingValue'); |
| 92 | |
| 93 | // Bitfield flag constants for ReadableState. Each constant uses left-shift (<<) to set a specific |
| 94 | // bit position, allowing multiple boolean flags to be stored efficiently in a single integer (kState). |
| 95 | // For example, `1 << 9` creates a value with only bit 9 set (value: 512). |
| 96 | const kEnded = 1 << 9; |
| 97 | const kEndEmitted = 1 << 10; |
| 98 | const kReading = 1 << 11; |
| 99 | const kSync = 1 << 12; |
| 100 | const kNeedReadable = 1 << 13; |
| 101 | const kEmittedReadable = 1 << 14; |
| 102 | const kReadableListening = 1 << 15; |
| 103 | const kResumeScheduled = 1 << 16; |
| 104 | const kMultiAwaitDrain = 1 << 17; |
| 105 | const kReadingMore = 1 << 18; |
| 106 | const kDataEmitted = 1 << 19; |
| 107 | const kDefaultUTF8Encoding = 1 << 20; |
| 108 | const kDecoder = 1 << 21; |
| 109 | const kEncoding = 1 << 22; |
| 110 | const kHasFlowing = 1 << 23; |
| 111 | const kFlowing = 1 << 24; |
| 112 | const kHasPaused = 1 << 25; |
| 113 | const kPaused = 1 << 26; |
| 114 | const kDataListening = 1 << 27; |
| 115 | |
| 116 | // ====================================================================================== |
| 117 | // ReadableState |
| 118 | |
| 119 | // TODO(benjamingr) it is likely slower to do it this way than with free functions |
| 120 | function makeBitMapDescriptor(bit) { |
| 121 | return { |
| 122 | enumerable: false, |
| 123 | get() { |
| 124 | return (this[kState] & bit) !== 0; |
| 125 | }, |
| 126 | set(value) { |
| 127 | if (value) this[kState] |= bit; |
| 128 | else this[kState] &= ~bit; |
| 129 | }, |
| 130 | }; |
| 131 | } |
| 132 | Object.defineProperties(ReadableState.prototype, { |
| 133 | objectMode: makeBitMapDescriptor(kObjectMode), |
| 134 | ended: makeBitMapDescriptor(kEnded), |
| 135 | endEmitted: makeBitMapDescriptor(kEndEmitted), |
| 136 | reading: makeBitMapDescriptor(kReading), |
| 137 | // Stream is still being constructed and cannot be |
| 138 | // destroyed until construction finished or failed. |
| 139 | // Async construction is opt in, therefore we start as |
| 140 | // constructed. |
| 141 | constructed: makeBitMapDescriptor(kConstructed), |
| 142 | // A flag to be able to tell if the event 'readable'/'data' is emitted |
| 143 | // immediately, or on a later tick. We set this to true at first, because |
| 144 | // any actions that shouldn't happen until "later" should generally also |
| 145 | // not happen before the first read call. |
| 146 | sync: makeBitMapDescriptor(kSync), |
| 147 | // Whenever we return null, then we set a flag to say |
| 148 | // that we're awaiting a 'readable' event emission. |
| 149 | needReadable: makeBitMapDescriptor(kNeedReadable), |
| 150 | emittedReadable: makeBitMapDescriptor(kEmittedReadable), |
| 151 | readableListening: makeBitMapDescriptor(kReadableListening), |
| 152 | resumeScheduled: makeBitMapDescriptor(kResumeScheduled), |
| 153 | // True if the error was already emitted and should not be thrown again. |
| 154 | errorEmitted: makeBitMapDescriptor(kErrorEmitted), |
| 155 | emitClose: makeBitMapDescriptor(kEmitClose), |
| 156 | autoDestroy: makeBitMapDescriptor(kAutoDestroy), |
| 157 | // Has it been destroyed. |
| 158 | destroyed: makeBitMapDescriptor(kDestroyed), |
| 159 | // Indicates whether the stream has finished destroying. |
| 160 | closed: makeBitMapDescriptor(kClosed), |
| 161 | // True if close has been emitted or would have been emitted |
| 162 | // depending on emitClose. |
| 163 | closeEmitted: makeBitMapDescriptor(kCloseEmitted), |
| 164 | multiAwaitDrain: makeBitMapDescriptor(kMultiAwaitDrain), |
| 165 | // If true, a maybeReadMore has been scheduled. |
| 166 | readingMore: makeBitMapDescriptor(kReadingMore), |
| 167 | dataEmitted: makeBitMapDescriptor(kDataEmitted), |
| 168 | |
| 169 | // Indicates whether the stream has errored. When true no further |
| 170 | // _read calls, 'data' or 'readable' events should occur. This is needed |
| 171 | // since when autoDestroy is disabled we need a way to tell whether the |
| 172 | // stream has failed. |
| 173 | errored: { |
| 174 | __proto__: null, |
| 175 | enumerable: false, |
| 176 | get() { |
| 177 | return (this[kState] & kErrored) !== 0 ? this[kErroredValue] : null; |
| 178 | }, |
| 179 | set(value) { |
| 180 | if (value) { |
| 181 | this[kErroredValue] = value; |
| 182 | this[kState] |= kErrored; |
| 183 | } else { |
| 184 | this[kState] &= ~kErrored; |
| 185 | } |
| 186 | }, |
| 187 | }, |
| 188 | |
| 189 | defaultEncoding: { |
| 190 | __proto__: null, |
| 191 | enumerable: false, |
| 192 | get() { |
| 193 | return (this[kState] & kDefaultUTF8Encoding) !== 0 |
| 194 | ? 'utf8' |
| 195 | : this[kDefaultEncodingValue]; |
| 196 | }, |
| 197 | set(value) { |
| 198 | if (value === 'utf8' || value === 'utf-8') { |
| 199 | this[kState] |= kDefaultUTF8Encoding; |
| 200 | } else { |
| 201 | this[kState] &= ~kDefaultUTF8Encoding; |
| 202 | this[kDefaultEncodingValue] = value; |
| 203 | } |
| 204 | }, |
| 205 | }, |
| 206 | |
| 207 | decoder: { |
| 208 | __proto__: null, |
| 209 | enumerable: false, |
| 210 | get() { |
| 211 | return (this[kState] & kDecoder) !== 0 ? this[kDecoderValue] : null; |
| 212 | }, |
| 213 | set(value) { |
| 214 | if (value) { |
| 215 | this[kDecoderValue] = value; |
| 216 | this[kState] |= kDecoder; |
| 217 | } else { |
| 218 | this[kState] &= ~kDecoder; |
| 219 | } |
| 220 | }, |
| 221 | }, |
| 222 | |
| 223 | encoding: { |
| 224 | __proto__: null, |
| 225 | enumerable: false, |
| 226 | get() { |
| 227 | return (this[kState] & kEncoding) !== 0 ? this[kEncodingValue] : null; |
| 228 | }, |
| 229 | set(value) { |
| 230 | if (value) { |
| 231 | this[kEncodingValue] = value; |
| 232 | this[kState] |= kEncoding; |
| 233 | } else { |
| 234 | this[kState] &= ~kEncoding; |
| 235 | } |
| 236 | }, |
| 237 | }, |
| 238 | |
| 239 | flowing: { |
| 240 | __proto__: null, |
| 241 | enumerable: false, |
| 242 | get() { |
| 243 | return (this[kState] & kHasFlowing) !== 0 |
| 244 | ? (this[kState] & kFlowing) !== 0 |
| 245 | : null; |
| 246 | }, |
| 247 | set(value) { |
| 248 | if (value == null) { |
| 249 | this[kState] &= ~(kHasFlowing | kFlowing); |
| 250 | } else if (value) { |
| 251 | this[kState] |= kHasFlowing | kFlowing; |
| 252 | } else { |
| 253 | this[kState] |= kHasFlowing; |
| 254 | this[kState] &= ~kFlowing; |
| 255 | } |
| 256 | }, |
| 257 | }, |
| 258 | }); |
| 259 | |
| 260 | export function ReadableState(options, _stream, isDuplex) { |
| 261 | // Bit map field to store ReadableState more efficiently with 1 bit per field |
| 262 | // instead of a V8 slot per field. |
| 263 | this[kState] = kEmitClose | kAutoDestroy | kConstructed | kSync; |
| 264 | |
| 265 | // Object stream flag. Used to make read(n) ignore n and to |
| 266 | // make all the buffer merging and length checks go away. |
| 267 | if (options?.objectMode) this[kState] |= kObjectMode; |
| 268 | |
| 269 | if (isDuplex && options?.readableObjectMode) this[kState] |= kObjectMode; |
| 270 | |
| 271 | // The point at which it stops calling _read() to fill the buffer |
| 272 | // Note: 0 is a valid value, means "don't call _read preemptively ever" |
| 273 | this.highWaterMark = options |
| 274 | ? getHighWaterMark(this, options, 'readableHighWaterMark', isDuplex) |
| 275 | : getDefaultHighWaterMark(false); |
| 276 | |
| 277 | this.buffer = []; |
| 278 | this.bufferIndex = 0; |
| 279 | this.length = 0; |
| 280 | this.pipes = []; |
| 281 | |
| 282 | // Should close be emitted on destroy. Defaults to true. |
| 283 | if (options && options.emitClose === false) this[kState] &= ~kEmitClose; |
| 284 | |
| 285 | // Should .destroy() be called after 'end' (and potentially 'finish'). |
| 286 | if (options && options.autoDestroy === false) this[kState] &= ~kAutoDestroy; |
| 287 | |
| 288 | // Crypto is kind of old and crusty. Historically, its default string |
| 289 | // encoding is 'binary' so we have to make this configurable. |
| 290 | // Everything else in the universe uses 'utf8', though. |
| 291 | const defaultEncoding = options?.defaultEncoding; |
| 292 | if ( |
| 293 | defaultEncoding == null || |
| 294 | defaultEncoding === 'utf8' || |
| 295 | defaultEncoding === 'utf-8' |
| 296 | ) { |
| 297 | this[kState] |= kDefaultUTF8Encoding; |
| 298 | } else if (Buffer.isEncoding(defaultEncoding)) { |
| 299 | this.defaultEncoding = defaultEncoding; |
| 300 | } else if (streamsNodejsV24Compat) { |
| 301 | // This is a semver-major change. Ref: https://github.com/nodejs/node/pull/46430 |
| 302 | throw new ERR_UNKNOWN_ENCODING(defaultEncoding); |
| 303 | } else { |
| 304 | this.defaultEncoding = defaultEncoding; |
| 305 | } |
| 306 | |
| 307 | // Ref the piped dest which we need a drain event on it |
| 308 | // type: null | Writable | Set<Writable>. |
| 309 | this.awaitDrainWriters = null; |
| 310 | |
| 311 | if (options?.encoding) { |
| 312 | this.decoder = new StringDecoder(options.encoding); |
| 313 | this.encoding = options.encoding; |
| 314 | } |
| 315 | } |
| 316 | |
| 317 | ReadableState.prototype[kOnConstructed] = function onConstructed(stream) { |
| 318 | if ((this[kState] & kNeedReadable) !== 0) { |
| 319 | maybeReadMore(stream, this); |
| 320 | } |
| 321 | }; |
| 322 | |
| 323 | // ====================================================================================== |
| 324 | // Readable |
| 325 | |
| 326 | Readable.ReadableState = ReadableState; |
| 327 | |
| 328 | Object.setPrototypeOf(Readable.prototype, Stream.prototype); |
| 329 | Object.setPrototypeOf(Readable, Stream); |
| 330 | |
| 331 | export function Readable(options) { |
| 332 | if (!(this instanceof Readable)) return new Readable(options); |
| 333 | |
| 334 | this._events ??= { |
| 335 | close: undefined, |
| 336 | error: undefined, |
| 337 | data: undefined, |
| 338 | end: undefined, |
| 339 | readable: undefined, |
| 340 | // Skip uncommon events... |
| 341 | // pause: undefined, |
| 342 | // resume: undefined, |
| 343 | // pipe: undefined, |
| 344 | // unpipe: undefined, |
| 345 | // [destroyImpl.kConstruct]: undefined, |
| 346 | // [destroyImpl.kDestroy]: undefined, |
| 347 | }; |
| 348 | |
| 349 | this._readableState = new ReadableState(options, this, false); |
| 350 | |
| 351 | if (options) { |
| 352 | if (typeof options.read === 'function') this._read = options.read; |
| 353 | |
| 354 | if (typeof options.destroy === 'function') this._destroy = options.destroy; |
| 355 | |
| 356 | if (typeof options.construct === 'function') |
| 357 | this._construct = options.construct; |
| 358 | |
| 359 | if (options.signal) addAbortSignal(options.signal, this); |
| 360 | } |
| 361 | |
| 362 | Stream.call(this, options); |
| 363 | |
| 364 | if (this._construct != null) { |
| 365 | construct(this, () => { |
| 366 | this._readableState[kOnConstructed](this); |
| 367 | }); |
| 368 | } |
| 369 | } |
| 370 | Readable.prototype.destroy = destroy; |
| 371 | Readable.prototype._undestroy = undestroy; |
| 372 | Readable.prototype._destroy = function (err, cb) { |
| 373 | if (cb) cb(err); |
| 374 | }; |
| 375 | |
| 376 | Readable.prototype[EventEmitter.captureRejectionSymbol] = function (err) { |
| 377 | this.destroy(err); |
| 378 | }; |
| 379 | |
| 380 | Readable.prototype[Symbol.asyncDispose] = async function () { |
| 381 | let error; |
| 382 | if (!this.destroyed) { |
| 383 | error = this.readableEnded ? null : new AbortError(); |
| 384 | this.destroy(error); |
| 385 | } |
| 386 | await new Promise((resolve, reject) => |
| 387 | eos(this, (err) => (err && err !== error ? reject(err) : resolve(null))) |
| 388 | ); |
| 389 | }; |
| 390 | |
| 391 | // Manually shove something into the read() buffer. |
| 392 | // This returns true if the highWaterMark has not been hit yet, |
| 393 | // similar to how Writable.write() returns true if you should |
| 394 | // write() some more. |
| 395 | Readable.prototype.push = function (chunk, encoding) { |
| 396 | const state = this._readableState; |
| 397 | return (state[kState] & kObjectMode) === 0 |
| 398 | ? readableAddChunkPushByteMode(this, state, chunk, encoding) |
| 399 | : readableAddChunkPushObjectMode(this, state, chunk, encoding); |
| 400 | }; |
| 401 | |
| 402 | // Unshift should *always* be something directly out of read(). |
| 403 | Readable.prototype.unshift = function (chunk, encoding) { |
| 404 | const state = this._readableState; |
| 405 | return (state[kState] & kObjectMode) === 0 |
| 406 | ? readableAddChunkUnshiftByteMode(this, state, chunk, encoding) |
| 407 | : readableAddChunkUnshiftObjectMode(this, state, chunk); |
| 408 | }; |
| 409 | |
| 410 | function readableAddChunkUnshiftByteMode(stream, state, chunk, encoding) { |
| 411 | if (chunk === null) { |
| 412 | state[kState] &= ~kReading; |
| 413 | onEofChunk(stream, state); |
| 414 | |
| 415 | return false; |
| 416 | } |
| 417 | |
| 418 | if (typeof chunk === 'string') { |
| 419 | encoding ||= state.defaultEncoding; |
| 420 | if (state.encoding !== encoding) { |
| 421 | if (state.encoding) { |
| 422 | // When unshifting, if state.encoding is set, we have to save |
| 423 | // the string in the BufferList with the state encoding. |
| 424 | chunk = Buffer.from(chunk, encoding).toString(state.encoding); |
| 425 | } else { |
| 426 | chunk = Buffer.from(chunk, encoding); |
| 427 | } |
| 428 | } |
| 429 | } else if (Stream._isArrayBufferView(chunk)) { |
| 430 | chunk = Stream._uint8ArrayToBuffer(chunk); |
| 431 | } else if (chunk !== undefined && !(chunk instanceof Buffer)) { |
| 432 | errorOrDestroy( |
| 433 | stream, |
| 434 | new ERR_INVALID_ARG_TYPE( |
| 435 | 'chunk', |
| 436 | ['string', 'Buffer', 'TypedArray', 'DataView'], |
| 437 | chunk |
| 438 | ) |
| 439 | ); |
| 440 | return false; |
| 441 | } |
| 442 | |
| 443 | if (!(chunk && chunk.length > 0)) { |
| 444 | return canPushMore(state); |
| 445 | } |
| 446 | |
| 447 | return readableAddChunkUnshiftValue(stream, state, chunk); |
| 448 | } |
| 449 | |
| 450 | function readableAddChunkUnshiftObjectMode(stream, state, chunk) { |
| 451 | if (chunk === null) { |
| 452 | state[kState] &= ~kReading; |
| 453 | onEofChunk(stream, state); |
| 454 | |
| 455 | return false; |
| 456 | } |
| 457 | |
| 458 | return readableAddChunkUnshiftValue(stream, state, chunk); |
| 459 | } |
| 460 | |
| 461 | function readableAddChunkUnshiftValue(stream, state, chunk) { |
| 462 | if ((state[kState] & kEndEmitted) !== 0) |
| 463 | errorOrDestroy(stream, new ERR_STREAM_UNSHIFT_AFTER_END_EVENT()); |
| 464 | else if ((state[kState] & (kDestroyed | kErrored)) !== 0) return false; |
| 465 | else addChunk(stream, state, chunk, true); |
| 466 | |
| 467 | return canPushMore(state); |
| 468 | } |
| 469 | |
| 470 | function readableAddChunkPushByteMode(stream, state, chunk, encoding) { |
| 471 | if (chunk === null) { |
| 472 | state[kState] &= ~kReading; |
| 473 | onEofChunk(stream, state); |
| 474 | return false; |
| 475 | } |
| 476 | |
| 477 | if (typeof chunk === 'string') { |
| 478 | encoding ||= state.defaultEncoding; |
| 479 | if (state.encoding !== encoding) { |
| 480 | chunk = Buffer.from(chunk, encoding); |
| 481 | encoding = ''; |
| 482 | } |
| 483 | } else if (chunk instanceof Buffer) { |
| 484 | encoding = ''; |
| 485 | } else if (Stream._isArrayBufferView(chunk)) { |
| 486 | chunk = Stream._uint8ArrayToBuffer(chunk); |
| 487 | encoding = ''; |
| 488 | } else if (chunk !== undefined) { |
| 489 | errorOrDestroy( |
| 490 | stream, |
| 491 | new ERR_INVALID_ARG_TYPE( |
| 492 | 'chunk', |
| 493 | ['string', 'Buffer', 'TypedArray', 'DataView'], |
| 494 | chunk |
| 495 | ) |
| 496 | ); |
| 497 | return false; |
| 498 | } |
| 499 | |
| 500 | if (!chunk || chunk.length <= 0) { |
| 501 | state[kState] &= ~kReading; |
| 502 | maybeReadMore(stream, state); |
| 503 | |
| 504 | return canPushMore(state); |
| 505 | } |
| 506 | |
| 507 | if ((state[kState] & kEnded) !== 0) { |
| 508 | errorOrDestroy(stream, new ERR_STREAM_PUSH_AFTER_EOF()); |
| 509 | return false; |
| 510 | } |
| 511 | |
| 512 | if ((state[kState] & (kDestroyed | kErrored)) !== 0) { |
| 513 | return false; |
| 514 | } |
| 515 | |
| 516 | state[kState] &= ~kReading; |
| 517 | if ((state[kState] & kDecoder) !== 0 && !encoding) { |
| 518 | chunk = state[kDecoderValue].write(chunk); |
| 519 | if (chunk.length === 0) { |
| 520 | maybeReadMore(stream, state); |
| 521 | return canPushMore(state); |
| 522 | } |
| 523 | } |
| 524 | |
| 525 | addChunk(stream, state, chunk, false); |
| 526 | return canPushMore(state); |
| 527 | } |
| 528 | |
| 529 | function readableAddChunkPushObjectMode(stream, state, chunk, encoding) { |
| 530 | if (chunk === null) { |
| 531 | state[kState] &= ~kReading; |
| 532 | onEofChunk(stream, state); |
| 533 | return false; |
| 534 | } |
| 535 | |
| 536 | if ((state[kState] & kEnded) !== 0) { |
| 537 | errorOrDestroy(stream, new ERR_STREAM_PUSH_AFTER_EOF()); |
| 538 | return false; |
| 539 | } |
| 540 | |
| 541 | if ((state[kState] & (kDestroyed | kErrored)) !== 0) { |
| 542 | return false; |
| 543 | } |
| 544 | |
| 545 | state[kState] &= ~kReading; |
| 546 | |
| 547 | if ((state[kState] & kDecoder) !== 0 && !encoding) { |
| 548 | chunk = state[kDecoderValue].write(chunk); |
| 549 | } |
| 550 | |
| 551 | addChunk(stream, state, chunk, false); |
| 552 | return canPushMore(state); |
| 553 | } |
| 554 | |
| 555 | function canPushMore(state) { |
| 556 | // We can push more data if we are below the highWaterMark. |
| 557 | // Also, if we have no data yet, we can stand some more bytes. |
| 558 | // This is to work around cases where hwm=0, such as the repl. |
| 559 | return ( |
| 560 | (state[kState] & kEnded) === 0 && |
| 561 | (state.length < state.highWaterMark || state.length === 0) |
| 562 | ); |
| 563 | } |
| 564 | |
| 565 | function addChunk(stream, state, chunk, addToFront) { |
| 566 | if ( |
| 567 | (state[kState] & (kFlowing | kSync | kDataListening)) === |
| 568 | (kFlowing | kDataListening) && |
| 569 | state.length === 0 |
| 570 | ) { |
| 571 | // Use the guard to avoid creating `Set()` repeatedly |
| 572 | // when we have multiple pipes. |
| 573 | if ((state[kState] & kMultiAwaitDrain) !== 0) { |
| 574 | state.awaitDrainWriters.clear(); |
| 575 | } else { |
| 576 | state.awaitDrainWriters = null; |
| 577 | } |
| 578 | |
| 579 | state[kState] |= kDataEmitted; |
| 580 | stream.emit('data', chunk); |
| 581 | } else { |
| 582 | // Update the buffer info. |
| 583 | state.length += (state[kState] & kObjectMode) !== 0 ? 1 : chunk.length; |
| 584 | if (addToFront) { |
| 585 | if (state.bufferIndex > 0) { |
| 586 | state.buffer[--state.bufferIndex] = chunk; |
| 587 | } else { |
| 588 | state.buffer.unshift(chunk); // Slow path |
| 589 | } |
| 590 | } else { |
| 591 | state.buffer.push(chunk); |
| 592 | } |
| 593 | |
| 594 | if ((state[kState] & kNeedReadable) !== 0) emitReadable(stream); |
| 595 | } |
| 596 | maybeReadMore(stream, state); |
| 597 | } |
| 598 | |
| 599 | Readable.prototype.isPaused = function () { |
| 600 | const state = this._readableState; |
| 601 | return ( |
| 602 | (state[kState] & kPaused) !== 0 || |
| 603 | (state[kState] & (kHasFlowing | kFlowing)) === kHasFlowing |
| 604 | ); |
| 605 | }; |
| 606 | |
| 607 | // Backwards compatibility. |
| 608 | Readable.prototype.setEncoding = function (enc) { |
| 609 | const state = this._readableState; |
| 610 | |
| 611 | const decoder = new StringDecoder(enc); |
| 612 | state.decoder = decoder; |
| 613 | // If setEncoding(null), decoder.encoding equals utf8. |
| 614 | state.encoding = state.decoder.encoding; |
| 615 | |
| 616 | // Iterate over current buffer to convert already stored Buffers: |
| 617 | let content = ''; |
| 618 | for (const data of state.buffer.slice(state.bufferIndex)) { |
| 619 | content += decoder.write(data); |
| 620 | } |
| 621 | state.buffer.length = 0; |
| 622 | state.bufferIndex = 0; |
| 623 | |
| 624 | if (content !== '') state.buffer.push(content); |
| 625 | state.length = content.length; |
| 626 | return this; |
| 627 | }; |
| 628 | |
| 629 | // Don't raise the hwm > 1GB. |
| 630 | const MAX_HWM = 0x40000000; |
| 631 | function computeNewHighWaterMark(n) { |
| 632 | if (n > MAX_HWM) { |
| 633 | throw new ERR_OUT_OF_RANGE('size', '<= 1GiB', n); |
| 634 | } else { |
| 635 | // Get the next highest power of 2 to prevent increasing hwm excessively in |
| 636 | // tiny amounts. |
| 637 | n--; |
| 638 | n |= n >>> 1; |
| 639 | n |= n >>> 2; |
| 640 | n |= n >>> 4; |
| 641 | n |= n >>> 8; |
| 642 | n |= n >>> 16; |
| 643 | n++; |
| 644 | } |
| 645 | return n; |
| 646 | } |
| 647 | |
| 648 | // This function is designed to be inlinable, so please take care when making |
| 649 | // changes to the function body. |
| 650 | function howMuchToRead(n, state) { |
| 651 | if (n <= 0 || (state.length === 0 && (state[kState] & kEnded) !== 0)) |
| 652 | return 0; |
| 653 | if ((state[kState] & kObjectMode) !== 0) return 1; |
| 654 | if (Number.isNaN(n)) { |
| 655 | // Only flow one buffer at a time. |
| 656 | if ((state[kState] & kFlowing) !== 0 && state.length) |
| 657 | return state.buffer[state.bufferIndex].length; |
| 658 | return state.length; |
| 659 | } |
| 660 | if (n <= state.length) return n; |
| 661 | return (state[kState] & kEnded) !== 0 ? state.length : 0; |
| 662 | } |
| 663 | |
| 664 | // You can override either this method, or the async _read(n) below. |
| 665 | Readable.prototype.read = function (n) { |
| 666 | // Same as parseInt(undefined, 10), however V8 7.3 performance regressed |
| 667 | // in this scenario, so we are doing it manually. |
| 668 | if (n === undefined) { |
| 669 | n = NaN; |
| 670 | } else if (!Number.isInteger(n)) { |
| 671 | n = Number.parseInt(n, 10); |
| 672 | } |
| 673 | const state = this._readableState; |
| 674 | const nOrig = n; |
| 675 | |
| 676 | // If we're asking for more than the current hwm, then raise the hwm. |
| 677 | if (n > state.highWaterMark) state.highWaterMark = computeNewHighWaterMark(n); |
| 678 | |
| 679 | if (n !== 0) state[kState] &= ~kEmittedReadable; |
| 680 | |
| 681 | // If we're doing read(0) to trigger a readable event, but we |
| 682 | // already have a bunch of data in the buffer, then just trigger |
| 683 | // the 'readable' event and move on. |
| 684 | if ( |
| 685 | n === 0 && |
| 686 | (state[kState] & kNeedReadable) !== 0 && |
| 687 | ((state.highWaterMark !== 0 |
| 688 | ? state.length >= state.highWaterMark |
| 689 | : state.length > 0) || |
| 690 | (state[kState] & kEnded) !== 0) |
| 691 | ) { |
| 692 | if (state.length === 0 && (state[kState] & kEnded) !== 0) endReadable(this); |
| 693 | else emitReadable(this); |
| 694 | return null; |
| 695 | } |
| 696 | |
| 697 | n = howMuchToRead(n, state); |
| 698 | |
| 699 | // If we've ended, and we're now clear, then finish it up. |
| 700 | if (n === 0 && (state[kState] & kEnded) !== 0) { |
| 701 | if (state.length === 0) endReadable(this); |
| 702 | return null; |
| 703 | } |
| 704 | |
| 705 | // All the actual chunk generation logic needs to be |
| 706 | // *below* the call to _read. The reason is that in certain |
| 707 | // synthetic stream cases, such as passthrough streams, _read |
| 708 | // may be a completely synchronous operation which may change |
| 709 | // the state of the read buffer, providing enough data when |
| 710 | // before there was *not* enough. |
| 711 | // |
| 712 | // So, the steps are: |
| 713 | // 1. Figure out what the state of things will be after we do |
| 714 | // a read from the buffer. |
| 715 | // |
| 716 | // 2. If that resulting state will trigger a _read, then call _read. |
| 717 | // Note that this may be asynchronous, or synchronous. Yes, it is |
| 718 | // deeply ugly to write APIs this way, but that still doesn't mean |
| 719 | // that the Readable class should behave improperly, as streams are |
| 720 | // designed to be sync/async agnostic. |
| 721 | // Take note if the _read call is sync or async (ie, if the read call |
| 722 | // has returned yet), so that we know whether or not it's safe to emit |
| 723 | // 'readable' etc. |
| 724 | // |
| 725 | // 3. Actually pull the requested chunks out of the buffer and return. |
| 726 | |
| 727 | // if we need a readable event, then we need to do some reading. |
| 728 | let doRead = (state[kState] & kNeedReadable) !== 0; |
| 729 | |
| 730 | // If we currently have less than the highWaterMark, then also read some. |
| 731 | if (state.length === 0 || state.length - n < state.highWaterMark) { |
| 732 | doRead = true; |
| 733 | } |
| 734 | |
| 735 | // However, if we've ended, then there's no point, if we're already |
| 736 | // reading, then it's unnecessary, if we're constructing we have to wait, |
| 737 | // and if we're destroyed or errored, then it's not allowed, |
| 738 | if ( |
| 739 | (state[kState] & |
| 740 | (kReading | kEnded | kDestroyed | kErrored | kConstructed)) !== |
| 741 | kConstructed |
| 742 | ) { |
| 743 | doRead = false; |
| 744 | } else if (doRead) { |
| 745 | state[kState] |= kReading | kSync; |
| 746 | // If the length is currently zero, then we *need* a readable event. |
| 747 | if (state.length === 0) state[kState] |= kNeedReadable; |
| 748 | |
| 749 | // Call internal read method |
| 750 | try { |
| 751 | this._read(state.highWaterMark); |
| 752 | } catch (err) { |
| 753 | errorOrDestroy(this, err); |
| 754 | } |
| 755 | state[kState] &= ~kSync; |
| 756 | |
| 757 | // If _read pushed data synchronously, then `reading` will be false, |
| 758 | // and we need to re-evaluate how much data we can return to the user. |
| 759 | if ((state[kState] & kReading) === 0) n = howMuchToRead(nOrig, state); |
| 760 | } |
| 761 | |
| 762 | let ret; |
| 763 | if (n > 0) ret = fromList(n, state); |
| 764 | else ret = null; |
| 765 | |
| 766 | if (ret === null) { |
| 767 | state[kState] |= state.length <= state.highWaterMark ? kNeedReadable : 0; |
| 768 | n = 0; |
| 769 | } else { |
| 770 | state.length -= n; |
| 771 | if ((state[kState] & kMultiAwaitDrain) !== 0) { |
| 772 | state.awaitDrainWriters.clear(); |
| 773 | } else { |
| 774 | state.awaitDrainWriters = null; |
| 775 | } |
| 776 | } |
| 777 | |
| 778 | if (state.length === 0) { |
| 779 | // If we have nothing in the buffer, then we want to know |
| 780 | // as soon as we *do* get something into the buffer. |
| 781 | if ((state[kState] & kEnded) === 0) state[kState] |= kNeedReadable; |
| 782 | |
| 783 | // If we tried to read() past the EOF, then emit end on the next tick. |
| 784 | if (nOrig !== n && (state[kState] & kEnded) !== 0) endReadable(this); |
| 785 | } |
| 786 | |
| 787 | if (ret !== null && (state[kState] & (kErrorEmitted | kCloseEmitted)) === 0) { |
| 788 | state[kState] |= kDataEmitted; |
| 789 | this.emit('data', ret); |
| 790 | } |
| 791 | |
| 792 | return ret; |
| 793 | }; |
| 794 | |
| 795 | function onEofChunk(stream, state) { |
| 796 | if ((state[kState] & kEnded) !== 0) return; |
| 797 | const decoder = |
| 798 | (state[kState] & kDecoder) !== 0 ? state[kDecoderValue] : null; |
| 799 | if (decoder) { |
| 800 | const chunk = decoder.end(); |
| 801 | if (chunk?.length) { |
| 802 | state.buffer.push(chunk); |
| 803 | state.length += (state[kState] & kObjectMode) !== 0 ? 1 : chunk.length; |
| 804 | } |
| 805 | } |
| 806 | state[kState] |= kEnded; |
| 807 | |
| 808 | if ((state[kState] & kSync) !== 0) { |
| 809 | // If we are sync, wait until next tick to emit the data. |
| 810 | // Otherwise we risk emitting data in the flow() |
| 811 | // the readable code triggers during a read() call. |
| 812 | emitReadable(stream); |
| 813 | } else { |
| 814 | // Emit 'readable' now to make sure it gets picked up. |
| 815 | state[kState] &= ~kNeedReadable; |
| 816 | state[kState] |= kEmittedReadable; |
| 817 | // We have to emit readable now that we are EOF. Modules |
| 818 | // in the ecosystem (e.g. dicer) rely on this event being sync. |
| 819 | emitReadable_(stream); |
| 820 | } |
| 821 | } |
| 822 | |
| 823 | // Don't emit readable right away in sync mode, because this can trigger |
| 824 | // another read() call => stack overflow. This way, it might trigger |
| 825 | // a nextTick recursion warning, but that's not so bad. |
| 826 | function emitReadable(stream) { |
| 827 | const state = stream._readableState; |
| 828 | state[kState] &= ~kNeedReadable; |
| 829 | if ((state[kState] & kEmittedReadable) === 0) { |
| 830 | state[kState] |= kEmittedReadable; |
| 831 | nextTick(emitReadable_, stream); |
| 832 | } |
| 833 | } |
| 834 | |
| 835 | function emitReadable_(stream) { |
| 836 | const state = stream._readableState; |
| 837 | if ( |
| 838 | (state[kState] & (kDestroyed | kErrored)) === 0 && |
| 839 | (state.length || (state[kState] & kEnded) !== 0) |
| 840 | ) { |
| 841 | stream.emit('readable'); |
| 842 | state[kState] &= ~kEmittedReadable; |
| 843 | } |
| 844 | |
| 845 | // The stream needs another readable event if: |
| 846 | // 1. It is not flowing, as the flow mechanism will take |
| 847 | // care of it. |
| 848 | // 2. It is not ended. |
| 849 | // 3. It is below the highWaterMark, so we can schedule |
| 850 | // another readable later. |
| 851 | state[kState] |= |
| 852 | (state[kState] & (kFlowing | kEnded)) === 0 && |
| 853 | state.length <= state.highWaterMark |
| 854 | ? kNeedReadable |
| 855 | : 0; |
| 856 | flow(stream); |
| 857 | } |
| 858 | |
| 859 | // At this point, the user has presumably seen the 'readable' event, |
| 860 | // and called read() to consume some data. that may have triggered |
| 861 | // in turn another _read(n) call, in which case reading = true if |
| 862 | // it's in progress. |
| 863 | // However, if we're not ended, or reading, and the length < hwm, |
| 864 | // then go ahead and try to read some more preemptively. |
| 865 | function maybeReadMore(stream, state) { |
| 866 | if ((state[kState] & (kReadingMore | kConstructed)) === kConstructed) { |
| 867 | state[kState] |= kReadingMore; |
| 868 | nextTick(maybeReadMore_, stream, state); |
| 869 | } |
| 870 | } |
| 871 | |
| 872 | function maybeReadMore_(stream, state) { |
| 873 | // Attempt to read more data if we should. |
| 874 | // |
| 875 | // The conditions for reading more data are (one of): |
| 876 | // - Not enough data buffered (state.length < state.highWaterMark). The loop |
| 877 | // is responsible for filling the buffer with enough data if such data |
| 878 | // is available. If highWaterMark is 0 and we are not in the flowing mode |
| 879 | // we should _not_ attempt to buffer any extra data. We'll get more data |
| 880 | // when the stream consumer calls read() instead. |
| 881 | // - No data in the buffer, and the stream is in flowing mode. In this mode |
| 882 | // the loop below is responsible for ensuring read() is called. Failing to |
| 883 | // call read here would abort the flow and there's no other mechanism for |
| 884 | // continuing the flow if the stream consumer has just subscribed to the |
| 885 | // 'data' event. |
| 886 | // |
| 887 | // In addition to the above conditions to keep reading data, the following |
| 888 | // conditions prevent the data from being read: |
| 889 | // - The stream has ended (state.ended). |
| 890 | // - There is already a pending 'read' operation (state.reading). This is a |
| 891 | // case where the stream has called the implementation defined _read() |
| 892 | // method, but they are processing the call asynchronously and have _not_ |
| 893 | // called push() with new data. In this case we skip performing more |
| 894 | // read()s. The execution ends in this method again after the _read() ends |
| 895 | // up calling push() with more data. |
| 896 | while ( |
| 897 | (state[kState] & (kReading | kEnded)) === 0 && |
| 898 | (state.length < state.highWaterMark || |
| 899 | ((state[kState] & kFlowing) !== 0 && state.length === 0)) |
| 900 | ) { |
| 901 | const len = state.length; |
| 902 | stream.read(0); |
| 903 | if (len === state.length) |
| 904 | // Didn't get any data, stop spinning. |
| 905 | break; |
| 906 | } |
| 907 | state[kState] &= ~kReadingMore; |
| 908 | } |
| 909 | |
| 910 | // Abstract method. to be overridden in specific implementation classes. |
| 911 | // call cb(er, data) where data is <= n in length. |
| 912 | // for virtual (non-string, non-buffer) streams, "length" is somewhat |
| 913 | // arbitrary, and perhaps not very meaningful. |
| 914 | Readable.prototype._read = function (_size) { |
| 915 | throw new ERR_METHOD_NOT_IMPLEMENTED('_read()'); |
| 916 | }; |
| 917 | |
| 918 | Readable.prototype.pipe = function (dest, pipeOpts) { |
| 919 | const src = this; // eslint-disable-line @typescript-eslint/no-this-alias |
| 920 | const state = this._readableState; |
| 921 | |
| 922 | if (state.pipes.length === 1) { |
| 923 | if ((state[kState] & kMultiAwaitDrain) === 0) { |
| 924 | state[kState] |= kMultiAwaitDrain; |
| 925 | state.awaitDrainWriters = new Set( |
| 926 | state.awaitDrainWriters ? [state.awaitDrainWriters] : [] |
| 927 | ); |
| 928 | } |
| 929 | } |
| 930 | |
| 931 | state.pipes.push(dest); |
| 932 | |
| 933 | const doEnd = !pipeOpts || pipeOpts.end !== false; |
| 934 | |
| 935 | const endFn = doEnd ? onend : unpipe; |
| 936 | if ((state[kState] & kEndEmitted) !== 0) nextTick(endFn); |
| 937 | else src.once('end', endFn); |
| 938 | |
| 939 | dest.on('unpipe', onunpipe); |
| 940 | function onunpipe(readable, unpipeInfo) { |
| 941 | if (readable === src) { |
| 942 | if (unpipeInfo && unpipeInfo.hasUnpiped === false) { |
| 943 | unpipeInfo.hasUnpiped = true; |
| 944 | cleanup(); |
| 945 | } |
| 946 | } |
| 947 | } |
| 948 | |
| 949 | function onend() { |
| 950 | dest.end(); |
| 951 | } |
| 952 | |
| 953 | let ondrain; |
| 954 | |
| 955 | let cleanedUp = false; |
| 956 | function cleanup() { |
| 957 | // Cleanup event handlers once the pipe is broken. |
| 958 | dest.removeListener('close', onclose); |
| 959 | dest.removeListener('finish', onfinish); |
| 960 | if (ondrain) { |
| 961 | dest.removeListener('drain', ondrain); |
| 962 | } |
| 963 | dest.removeListener('error', onerror); |
| 964 | dest.removeListener('unpipe', onunpipe); |
| 965 | src.removeListener('end', onend); |
| 966 | src.removeListener('end', unpipe); |
| 967 | src.removeListener('data', ondata); |
| 968 | |
| 969 | cleanedUp = true; |
| 970 | |
| 971 | // If the reader is waiting for a drain event from this |
| 972 | // specific writer, then it would cause it to never start |
| 973 | // flowing again. |
| 974 | // So, if this is awaiting a drain, then we just call it now. |
| 975 | // If we don't know, then assume that we are waiting for one. |
| 976 | if ( |
| 977 | ondrain && |
| 978 | state.awaitDrainWriters && |
| 979 | (!dest._writableState || dest._writableState.needDrain) |
| 980 | ) |
| 981 | ondrain(); |
| 982 | } |
| 983 | |
| 984 | function pause() { |
| 985 | // If the user unpiped during `dest.write()`, it is possible |
| 986 | // to get stuck in a permanently paused state if that write |
| 987 | // also returned false. |
| 988 | // => Check whether `dest` is still a piping destination. |
| 989 | if (!cleanedUp) { |
| 990 | if (state.pipes.length === 1 && state.pipes[0] === dest) { |
| 991 | state.awaitDrainWriters = dest; |
| 992 | state[kState] &= ~kMultiAwaitDrain; |
| 993 | } else if (state.pipes.length > 1 && state.pipes.includes(dest)) { |
| 994 | state.awaitDrainWriters.add(dest); |
| 995 | } |
| 996 | src.pause(); |
| 997 | } |
| 998 | if (!ondrain) { |
| 999 | // When the dest drains, it reduces the awaitDrain counter |
| 1000 | // on the source. This would be more elegant with a .once() |
| 1001 | // handler in flow(), but adding and removing repeatedly is |
| 1002 | // too slow. |
| 1003 | ondrain = pipeOnDrain(src, dest); |
| 1004 | dest.on('drain', ondrain); |
| 1005 | } |
| 1006 | } |
| 1007 | |
| 1008 | src.on('data', ondata); |
| 1009 | function ondata(chunk) { |
| 1010 | // This is a semver-major change. Ref: https://github.com/nodejs/node/pull/55270 |
| 1011 | if (streamsNodejsV24Compat) { |
| 1012 | try { |
| 1013 | const ret = dest.write(chunk); |
| 1014 | if (ret === false) { |
| 1015 | pause(); |
| 1016 | } |
| 1017 | } catch (error) { |
| 1018 | dest.destroy(error); |
| 1019 | } |
| 1020 | } else { |
| 1021 | const ret = dest.write(chunk); |
| 1022 | if (ret === false) { |
| 1023 | pause(); |
| 1024 | } |
| 1025 | } |
| 1026 | } |
| 1027 | |
| 1028 | // If the dest has an error, then stop piping into it. |
| 1029 | // However, don't suppress the throwing behavior for this. |
| 1030 | function onerror(er) { |
| 1031 | unpipe(); |
| 1032 | dest.removeListener('error', onerror); |
| 1033 | if (dest.listenerCount('error') === 0) { |
| 1034 | const s = dest._writableState || dest._readableState; |
| 1035 | if (s && !s.errorEmitted) { |
| 1036 | // User incorrectly emitted 'error' directly on the stream. |
| 1037 | errorOrDestroy(dest, er); |
| 1038 | } else { |
| 1039 | dest.emit('error', er); |
| 1040 | } |
| 1041 | } |
| 1042 | } |
| 1043 | |
| 1044 | // Make sure our error handler is attached before userland ones. |
| 1045 | dest.prependListener('error', onerror); |
| 1046 | |
| 1047 | // Both close and finish should trigger unpipe, but only once. |
| 1048 | function onclose() { |
| 1049 | dest.removeListener('finish', onfinish); |
| 1050 | unpipe(); |
| 1051 | } |
| 1052 | dest.once('close', onclose); |
| 1053 | function onfinish() { |
| 1054 | dest.removeListener('close', onclose); |
| 1055 | unpipe(); |
| 1056 | } |
| 1057 | dest.once('finish', onfinish); |
| 1058 | |
| 1059 | function unpipe() { |
| 1060 | src.unpipe(dest); |
| 1061 | } |
| 1062 | |
| 1063 | // Tell the dest that it's being piped to. |
| 1064 | dest.emit('pipe', src); |
| 1065 | |
| 1066 | // Start the flow if it hasn't been started already. |
| 1067 | |
| 1068 | if (dest.writableNeedDrain === true) { |
| 1069 | pause(); |
| 1070 | } else if ((state[kState] & kFlowing) === 0) { |
| 1071 | src.resume(); |
| 1072 | } |
| 1073 | |
| 1074 | return dest; |
| 1075 | }; |
| 1076 | |
| 1077 | function pipeOnDrain(src, dest) { |
| 1078 | return function pipeOnDrainFunctionResult() { |
| 1079 | const state = src._readableState; |
| 1080 | |
| 1081 | // `ondrain` will call directly, |
| 1082 | // `this` maybe not a reference to dest, |
| 1083 | // so we use the real dest here. |
| 1084 | if (state.awaitDrainWriters === dest) { |
| 1085 | state.awaitDrainWriters = null; |
| 1086 | } else if ((state[kState] & kMultiAwaitDrain) !== 0) { |
| 1087 | state.awaitDrainWriters.delete(dest); |
| 1088 | } |
| 1089 | |
| 1090 | if ( |
| 1091 | (!state.awaitDrainWriters || state.awaitDrainWriters.size === 0) && |
| 1092 | (state[kState] & kDataListening) !== 0 |
| 1093 | ) { |
| 1094 | src.resume(); |
| 1095 | } |
| 1096 | }; |
| 1097 | } |
| 1098 | |
| 1099 | Readable.prototype.unpipe = function (dest) { |
| 1100 | const state = this._readableState; |
| 1101 | const unpipeInfo = { hasUnpiped: false }; |
| 1102 | |
| 1103 | // If we're not piping anywhere, then do nothing. |
| 1104 | if (state.pipes.length === 0) return this; |
| 1105 | |
| 1106 | if (!dest) { |
| 1107 | // remove all. |
| 1108 | const dests = state.pipes; |
| 1109 | state.pipes = []; |
| 1110 | this.pause(); |
| 1111 | |
| 1112 | for (let i = 0; i < dests.length; i++) |
| 1113 | dests[i].emit('unpipe', this, { hasUnpiped: false }); |
| 1114 | return this; |
| 1115 | } |
| 1116 | |
| 1117 | // Try to find the right one. |
| 1118 | const index = state.pipes.indexOf(dest); |
| 1119 | if (index === -1) return this; |
| 1120 | |
| 1121 | state.pipes.splice(index, 1); |
| 1122 | if (state.pipes.length === 0) this.pause(); |
| 1123 | dest.emit('unpipe', this, unpipeInfo); |
| 1124 | return this; |
| 1125 | }; |
| 1126 | |
| 1127 | // Set up data events if they are asked for |
| 1128 | // Ensure readable listeners eventually get something. |
| 1129 | Readable.prototype.on = function (ev, fn) { |
| 1130 | const res = Stream.prototype.on.call(this, ev, fn); |
| 1131 | const state = this._readableState; |
| 1132 | |
| 1133 | if (ev === 'data') { |
| 1134 | state[kState] |= kDataListening; |
| 1135 | |
| 1136 | // Update readableListening so that resume() may be a no-op |
| 1137 | // a few lines down. This is needed to support once('readable'). |
| 1138 | state[kState] |= |
| 1139 | this.listenerCount('readable') > 0 ? kReadableListening : 0; |
| 1140 | |
| 1141 | // Try start flowing on next tick if stream isn't explicitly paused. |
| 1142 | if ((state[kState] & (kHasFlowing | kFlowing)) !== kHasFlowing) { |
| 1143 | this.resume(); |
| 1144 | } |
| 1145 | } else if (ev === 'readable') { |
| 1146 | if ((state[kState] & (kEndEmitted | kReadableListening)) === 0) { |
| 1147 | state[kState] |= kReadableListening | kNeedReadable | kHasFlowing; |
| 1148 | state[kState] &= ~(kFlowing | kEmittedReadable); |
| 1149 | if (state.length) { |
| 1150 | emitReadable(this); |
| 1151 | } else if ((state[kState] & kReading) === 0) { |
| 1152 | nextTick(nReadingNextTick, this); |
| 1153 | } |
| 1154 | } |
| 1155 | } |
| 1156 | |
| 1157 | return res; |
| 1158 | }; |
| 1159 | Readable.prototype.addListener = Readable.prototype.on; |
| 1160 | |
| 1161 | Readable.prototype.removeListener = function (ev, fn) { |
| 1162 | const state = this._readableState; |
| 1163 | |
| 1164 | const res = Stream.prototype.removeListener.call(this, ev, fn); |
| 1165 | |
| 1166 | if (ev === 'readable') { |
| 1167 | // We need to check if there is someone still listening to |
| 1168 | // readable and reset the state. However this needs to happen |
| 1169 | // after readable has been emitted but before I/O (nextTick) to |
| 1170 | // support once('readable', fn) cycles. This means that calling |
| 1171 | // resume within the same tick will have no |
| 1172 | // effect. |
| 1173 | nextTick(updateReadableListening, this); |
| 1174 | } else if (ev === 'data' && this.listenerCount('data') === 0) { |
| 1175 | state[kState] &= ~kDataListening; |
| 1176 | } |
| 1177 | |
| 1178 | return res; |
| 1179 | }; |
| 1180 | Readable.prototype.off = Readable.prototype.removeListener; |
| 1181 | |
| 1182 | Readable.prototype.removeAllListeners = function (ev) { |
| 1183 | const res = Stream.prototype.removeAllListeners.apply(this, arguments); |
| 1184 | |
| 1185 | if (ev === 'readable' || ev === undefined) { |
| 1186 | // We need to check if there is someone still listening to |
| 1187 | // readable and reset the state. However this needs to happen |
| 1188 | // after readable has been emitted but before I/O (nextTick) to |
| 1189 | // support once('readable', fn) cycles. This means that calling |
| 1190 | // resume within the same tick will have no |
| 1191 | // effect. |
| 1192 | nextTick(updateReadableListening, this); |
| 1193 | } |
| 1194 | |
| 1195 | return res; |
| 1196 | }; |
| 1197 | |
| 1198 | function updateReadableListening(self) { |
| 1199 | const state = self._readableState; |
| 1200 | |
| 1201 | if (self.listenerCount('readable') > 0) { |
| 1202 | state[kState] |= kReadableListening; |
| 1203 | } else { |
| 1204 | state[kState] &= ~kReadableListening; |
| 1205 | } |
| 1206 | |
| 1207 | if ( |
| 1208 | (state[kState] & (kHasPaused | kPaused | kResumeScheduled)) === |
| 1209 | (kHasPaused | kResumeScheduled) |
| 1210 | ) { |
| 1211 | // Flowing needs to be set to true now, otherwise |
| 1212 | // the upcoming resume will not flow. |
| 1213 | state[kState] |= kHasFlowing | kFlowing; |
| 1214 | |
| 1215 | // Crude way to check if we should resume. |
| 1216 | } else if ((state[kState] & kDataListening) !== 0) { |
| 1217 | self.resume(); |
| 1218 | } else if ((state[kState] & kReadableListening) === 0) { |
| 1219 | state[kState] &= ~(kHasFlowing | kFlowing); |
| 1220 | } |
| 1221 | } |
| 1222 | |
| 1223 | function nReadingNextTick(self) { |
| 1224 | self.read(0); |
| 1225 | } |
| 1226 | |
| 1227 | // pause() and resume() are remnants of the legacy readable stream API |
| 1228 | // If the user uses them, then switch into old mode. |
| 1229 | Readable.prototype.resume = function () { |
| 1230 | const state = this._readableState; |
| 1231 | if ((state[kState] & kFlowing) === 0) { |
| 1232 | // We flow only if there is no one listening |
| 1233 | // for readable, but we still have to call |
| 1234 | // resume(). |
| 1235 | state[kState] |= kHasFlowing; |
| 1236 | if ((state[kState] & kReadableListening) === 0) { |
| 1237 | state[kState] |= kFlowing; |
| 1238 | } else { |
| 1239 | state[kState] &= ~kFlowing; |
| 1240 | } |
| 1241 | resume(this, state); |
| 1242 | } |
| 1243 | state[kState] |= kHasPaused; |
| 1244 | state[kState] &= ~kPaused; |
| 1245 | return this; |
| 1246 | }; |
| 1247 | |
| 1248 | function resume(stream, state) { |
| 1249 | if ((state[kState] & kResumeScheduled) === 0) { |
| 1250 | state[kState] |= kResumeScheduled; |
| 1251 | nextTick(resume_, stream, state); |
| 1252 | } |
| 1253 | } |
| 1254 | |
| 1255 | function resume_(stream, state) { |
| 1256 | if ((state[kState] & kReading) === 0) { |
| 1257 | stream.read(0); |
| 1258 | } |
| 1259 | |
| 1260 | state[kState] &= ~kResumeScheduled; |
| 1261 | stream.emit('resume'); |
| 1262 | flow(stream); |
| 1263 | if ((state[kState] & (kFlowing | kReading)) === kFlowing) stream.read(0); |
| 1264 | } |
| 1265 | |
| 1266 | Readable.prototype.pause = function () { |
| 1267 | const state = this._readableState; |
| 1268 | if ((state[kState] & (kHasFlowing | kFlowing)) !== kHasFlowing) { |
| 1269 | state[kState] |= kHasFlowing; |
| 1270 | state[kState] &= ~kFlowing; |
| 1271 | this.emit('pause'); |
| 1272 | } |
| 1273 | state[kState] |= kHasPaused | kPaused; |
| 1274 | return this; |
| 1275 | }; |
| 1276 | |
| 1277 | function flow(stream) { |
| 1278 | const state = stream._readableState; |
| 1279 | while ((state[kState] & kFlowing) !== 0 && stream.read() !== null); |
| 1280 | } |
| 1281 | |
| 1282 | // Wrap an old-style stream as the async data source. |
| 1283 | // This is *not* part of the readable stream interface. |
| 1284 | // It is an ugly unfortunate mess of history. |
| 1285 | Readable.prototype.wrap = function (stream) { |
| 1286 | let paused = false; |
| 1287 | |
| 1288 | // TODO (ronag): Should this.destroy(err) emit |
| 1289 | // 'error' on the wrapped stream? Would require |
| 1290 | // a static factory method, e.g. Readable.wrap(stream). |
| 1291 | |
| 1292 | stream.on('data', (chunk) => { |
| 1293 | if (!this.push(chunk) && stream.pause) { |
| 1294 | paused = true; |
| 1295 | stream.pause(); |
| 1296 | } |
| 1297 | }); |
| 1298 | |
| 1299 | stream.on('end', () => { |
| 1300 | this.push(null); |
| 1301 | }); |
| 1302 | |
| 1303 | stream.on('error', (err) => { |
| 1304 | errorOrDestroy(this, err); |
| 1305 | }); |
| 1306 | |
| 1307 | stream.on('close', () => { |
| 1308 | this.destroy(); |
| 1309 | }); |
| 1310 | |
| 1311 | stream.on('destroy', () => { |
| 1312 | this.destroy(); |
| 1313 | }); |
| 1314 | |
| 1315 | this._read = () => { |
| 1316 | if (paused && stream.resume) { |
| 1317 | paused = false; |
| 1318 | stream.resume(); |
| 1319 | } |
| 1320 | }; |
| 1321 | |
| 1322 | // Proxy all the other methods. Important when wrapping filters and duplexes. |
| 1323 | const streamKeys = Object.keys(stream); |
| 1324 | for (let j = 1; j < streamKeys.length; j++) { |
| 1325 | const i = streamKeys[j]; |
| 1326 | if (this[i] === undefined && typeof stream[i] === 'function') { |
| 1327 | this[i] = stream[i].bind(stream); |
| 1328 | } |
| 1329 | } |
| 1330 | |
| 1331 | return this; |
| 1332 | }; |
| 1333 | |
| 1334 | Readable.prototype[Symbol.asyncIterator] = function () { |
| 1335 | return streamToAsyncIterator(this); |
| 1336 | }; |
| 1337 | |
| 1338 | Readable.prototype.iterator = function (options) { |
| 1339 | if (options !== undefined) { |
| 1340 | validateObject(options, 'options'); |
| 1341 | } |
| 1342 | return streamToAsyncIterator(this, options); |
| 1343 | }; |
| 1344 | |
| 1345 | function streamToAsyncIterator(stream, options) { |
| 1346 | if (typeof stream.read !== 'function') { |
| 1347 | stream = Readable.wrap(stream, { objectMode: true }); |
| 1348 | } |
| 1349 | |
| 1350 | const iter = createAsyncIterator(stream, options); |
| 1351 | iter.stream = stream; |
| 1352 | return iter; |
| 1353 | } |
| 1354 | |
| 1355 | async function* createAsyncIterator(stream, options) { |
| 1356 | let callback = nop; |
| 1357 | |
| 1358 | function next(resolve) { |
| 1359 | if (this === stream) { |
| 1360 | callback(); |
| 1361 | callback = nop; |
| 1362 | } else { |
| 1363 | callback = resolve; |
| 1364 | } |
| 1365 | } |
| 1366 | |
| 1367 | stream.on('readable', next); |
| 1368 | |
| 1369 | let error; |
| 1370 | const cleanup = eos(stream, { writable: false }, (err) => { |
| 1371 | error = err ? aggregateTwoErrors(error, err) : null; |
| 1372 | callback(); |
| 1373 | callback = nop; |
| 1374 | }); |
| 1375 | |
| 1376 | try { |
| 1377 | while (true) { |
| 1378 | const chunk = stream.destroyed ? null : stream.read(); |
| 1379 | if (chunk !== null) { |
| 1380 | yield chunk; |
| 1381 | } else if (error) { |
| 1382 | throw error; |
| 1383 | } else if (error === null) { |
| 1384 | return; |
| 1385 | } else { |
| 1386 | await new Promise(next); |
| 1387 | } |
| 1388 | } |
| 1389 | } catch (err) { |
| 1390 | error = aggregateTwoErrors(error, err); |
| 1391 | throw error; |
| 1392 | } finally { |
| 1393 | if ( |
| 1394 | (error || options?.destroyOnReturn !== false) && |
| 1395 | (error === undefined || stream._readableState.autoDestroy) |
| 1396 | ) { |
| 1397 | destroyer(stream, null); |
| 1398 | } else { |
| 1399 | stream.off('readable', next); |
| 1400 | cleanup(); |
| 1401 | } |
| 1402 | } |
| 1403 | } |
| 1404 | |
| 1405 | // Making it explicit these properties are not enumerable |
| 1406 | // because otherwise some prototype manipulation in |
| 1407 | // userland will fail. |
| 1408 | Object.defineProperties(Readable.prototype, { |
| 1409 | readable: { |
| 1410 | __proto__: null, |
| 1411 | get() { |
| 1412 | const r = this._readableState; |
| 1413 | // r.readable === false means that this is part of a Duplex stream |
| 1414 | // where the readable side was disabled upon construction. |
| 1415 | // Compat. The user might manually disable readable side through |
| 1416 | // deprecated setter. |
| 1417 | return ( |
| 1418 | !!r && |
| 1419 | r.readable !== false && |
| 1420 | !r.destroyed && |
| 1421 | !r.errorEmitted && |
| 1422 | !r.endEmitted |
| 1423 | ); |
| 1424 | }, |
| 1425 | set(val) { |
| 1426 | // Backwards compat. |
| 1427 | if (this._readableState) { |
| 1428 | this._readableState.readable = !!val; |
| 1429 | } |
| 1430 | }, |
| 1431 | }, |
| 1432 | |
| 1433 | readableDidRead: { |
| 1434 | __proto__: null, |
| 1435 | enumerable: false, |
| 1436 | get: function () { |
| 1437 | return this._readableState.dataEmitted; |
| 1438 | }, |
| 1439 | }, |
| 1440 | |
| 1441 | readableAborted: { |
| 1442 | __proto__: null, |
| 1443 | enumerable: false, |
| 1444 | get: function () { |
| 1445 | return !!( |
| 1446 | this._readableState.readable !== false && |
| 1447 | (this._readableState.destroyed || this._readableState.errored) && |
| 1448 | !this._readableState.endEmitted |
| 1449 | ); |
| 1450 | }, |
| 1451 | }, |
| 1452 | |
| 1453 | readableHighWaterMark: { |
| 1454 | __proto__: null, |
| 1455 | enumerable: false, |
| 1456 | get: function () { |
| 1457 | return this._readableState.highWaterMark; |
| 1458 | }, |
| 1459 | }, |
| 1460 | |
| 1461 | readableBuffer: { |
| 1462 | __proto__: null, |
| 1463 | enumerable: false, |
| 1464 | get: function () { |
| 1465 | return this._readableState?.buffer; |
| 1466 | }, |
| 1467 | }, |
| 1468 | |
| 1469 | readableFlowing: { |
| 1470 | __proto__: null, |
| 1471 | enumerable: false, |
| 1472 | get: function () { |
| 1473 | return this._readableState.flowing; |
| 1474 | }, |
| 1475 | set: function (state) { |
| 1476 | if (this._readableState) { |
| 1477 | this._readableState.flowing = state; |
| 1478 | } |
| 1479 | }, |
| 1480 | }, |
| 1481 | |
| 1482 | readableLength: { |
| 1483 | __proto__: null, |
| 1484 | enumerable: false, |
| 1485 | get() { |
| 1486 | return this._readableState.length; |
| 1487 | }, |
| 1488 | }, |
| 1489 | |
| 1490 | readableObjectMode: { |
| 1491 | __proto__: null, |
| 1492 | enumerable: false, |
| 1493 | get() { |
| 1494 | return this._readableState ? this._readableState.objectMode : false; |
| 1495 | }, |
| 1496 | }, |
| 1497 | |
| 1498 | readableEncoding: { |
| 1499 | __proto__: null, |
| 1500 | enumerable: false, |
| 1501 | get() { |
| 1502 | return this._readableState ? this._readableState.encoding : null; |
| 1503 | }, |
| 1504 | }, |
| 1505 | |
| 1506 | errored: { |
| 1507 | __proto__: null, |
| 1508 | enumerable: false, |
| 1509 | get() { |
| 1510 | return this._readableState ? this._readableState.errored : null; |
| 1511 | }, |
| 1512 | }, |
| 1513 | |
| 1514 | closed: { |
| 1515 | __proto__: null, |
| 1516 | get() { |
| 1517 | return this._readableState ? this._readableState.closed : false; |
| 1518 | }, |
| 1519 | }, |
| 1520 | |
| 1521 | destroyed: { |
| 1522 | __proto__: null, |
| 1523 | enumerable: false, |
| 1524 | get() { |
| 1525 | return this._readableState ? this._readableState.destroyed : false; |
| 1526 | }, |
| 1527 | set(value) { |
| 1528 | // We ignore the value if the stream |
| 1529 | // has not been initialized yet. |
| 1530 | if (!this._readableState) { |
| 1531 | return; |
| 1532 | } |
| 1533 | |
| 1534 | // Backward compatibility, the user is explicitly |
| 1535 | // managing destroyed. |
| 1536 | this._readableState.destroyed = value; |
| 1537 | }, |
| 1538 | }, |
| 1539 | |
| 1540 | readableEnded: { |
| 1541 | __proto__: null, |
| 1542 | enumerable: false, |
| 1543 | get() { |
| 1544 | return this._readableState ? this._readableState.endEmitted : false; |
| 1545 | }, |
| 1546 | }, |
| 1547 | }); |
| 1548 | |
| 1549 | Object.defineProperties(ReadableState.prototype, { |
| 1550 | // Legacy getter for `pipesCount`. |
| 1551 | pipesCount: { |
| 1552 | __proto__: null, |
| 1553 | get() { |
| 1554 | return this.pipes.length; |
| 1555 | }, |
| 1556 | }, |
| 1557 | |
| 1558 | // Legacy property for `paused`. |
| 1559 | paused: { |
| 1560 | __proto__: null, |
| 1561 | get() { |
| 1562 | return (this[kState] & kPaused) !== 0; |
| 1563 | }, |
| 1564 | set(value) { |
| 1565 | this[kState] |= kHasPaused; |
| 1566 | if (value) { |
| 1567 | this[kState] |= kPaused; |
| 1568 | } else { |
| 1569 | this[kState] &= ~kPaused; |
| 1570 | } |
| 1571 | }, |
| 1572 | }, |
| 1573 | }); |
| 1574 | |
| 1575 | // Exposed for testing purposes only. |
| 1576 | Readable._fromList = fromList; |
| 1577 | |
| 1578 | // Pluck off n bytes from an array of buffers. |
| 1579 | // Length is the combined lengths of all the buffers in the list. |
| 1580 | // This function is designed to be inlinable, so please take care when making |
| 1581 | // changes to the function body. |
| 1582 | function fromList(n, state) { |
| 1583 | // nothing buffered. |
| 1584 | if (state.length === 0) return null; |
| 1585 | |
| 1586 | let idx = state.bufferIndex; |
| 1587 | let ret; |
| 1588 | |
| 1589 | const buf = state.buffer; |
| 1590 | const len = buf.length; |
| 1591 | |
| 1592 | if ((state[kState] & kObjectMode) !== 0) { |
| 1593 | ret = buf[idx]; |
| 1594 | buf[idx++] = null; |
| 1595 | } else if (!n || n >= state.length) { |
| 1596 | // Read it all, truncate the list. |
| 1597 | if ((state[kState] & kDecoder) !== 0) { |
| 1598 | ret = ''; |
| 1599 | while (idx < len) { |
| 1600 | ret += buf[idx]; |
| 1601 | buf[idx++] = null; |
| 1602 | } |
| 1603 | } else if (len - idx === 0) { |
| 1604 | ret = Buffer.alloc(0); |
| 1605 | } else if (len - idx === 1) { |
| 1606 | ret = buf[idx]; |
| 1607 | buf[idx++] = null; |
| 1608 | } else { |
| 1609 | ret = Buffer.allocUnsafe(state.length); |
| 1610 | |
| 1611 | let i = 0; |
| 1612 | while (idx < len) { |
| 1613 | ret.set(buf[idx], i); |
| 1614 | i += buf[idx].length; |
| 1615 | buf[idx++] = null; |
| 1616 | } |
| 1617 | } |
| 1618 | } else if (n < buf[idx].length) { |
| 1619 | // `slice` is the same for buffers and strings. |
| 1620 | ret = buf[idx].slice(0, n); |
| 1621 | buf[idx] = buf[idx].slice(n); |
| 1622 | } else if (n === buf[idx].length) { |
| 1623 | // First chunk is a perfect match. |
| 1624 | ret = buf[idx]; |
| 1625 | buf[idx++] = null; |
| 1626 | } else if ((state[kState] & kDecoder) !== 0) { |
| 1627 | ret = ''; |
| 1628 | while (idx < len) { |
| 1629 | const str = buf[idx]; |
| 1630 | if (n > str.length) { |
| 1631 | ret += str; |
| 1632 | n -= str.length; |
| 1633 | buf[idx++] = null; |
| 1634 | } else { |
| 1635 | if (n === buf.length) { |
| 1636 | ret += str; |
| 1637 | buf[idx++] = null; |
| 1638 | } else { |
| 1639 | ret += str.slice(0, n); |
| 1640 | buf[idx] = str.slice(n); |
| 1641 | } |
| 1642 | break; |
| 1643 | } |
| 1644 | } |
| 1645 | } else { |
| 1646 | ret = Buffer.allocUnsafe(n); |
| 1647 | |
| 1648 | const retLen = n; |
| 1649 | while (idx < len) { |
| 1650 | const data = buf[idx]; |
| 1651 | if (n > data.length) { |
| 1652 | ret.set(data, retLen - n); |
| 1653 | n -= data.length; |
| 1654 | buf[idx++] = null; |
| 1655 | } else { |
| 1656 | if (n === data.length) { |
| 1657 | ret.set(data, retLen - n); |
| 1658 | buf[idx++] = null; |
| 1659 | } else { |
| 1660 | ret.set(Buffer.from(data.buffer, data.byteOffset, n), retLen - n); |
| 1661 | buf[idx] = Buffer.from( |
| 1662 | data.buffer, |
| 1663 | data.byteOffset + n, |
| 1664 | data.length - n |
| 1665 | ); |
| 1666 | } |
| 1667 | break; |
| 1668 | } |
| 1669 | } |
| 1670 | } |
| 1671 | |
| 1672 | if (idx === len) { |
| 1673 | state.buffer.length = 0; |
| 1674 | state.bufferIndex = 0; |
| 1675 | } else if (idx > 1024) { |
| 1676 | state.buffer.splice(0, idx); |
| 1677 | state.bufferIndex = 0; |
| 1678 | } else { |
| 1679 | state.bufferIndex = idx; |
| 1680 | } |
| 1681 | |
| 1682 | return ret; |
| 1683 | } |
| 1684 | |
| 1685 | function endReadable(stream) { |
| 1686 | const state = stream._readableState; |
| 1687 | |
| 1688 | if ((state[kState] & kEndEmitted) === 0) { |
| 1689 | state[kState] |= kEnded; |
| 1690 | nextTick(endReadableNT, state, stream); |
| 1691 | } |
| 1692 | } |
| 1693 | |
| 1694 | function endReadableNT(state, stream) { |
| 1695 | // Check that we didn't get one last unshift. |
| 1696 | if ( |
| 1697 | (state[kState] & (kErrored | kCloseEmitted | kEndEmitted)) === 0 && |
| 1698 | state.length === 0 |
| 1699 | ) { |
| 1700 | state[kState] |= kEndEmitted; |
| 1701 | stream.emit('end'); |
| 1702 | |
| 1703 | if (stream.writable && stream.allowHalfOpen === false) { |
| 1704 | nextTick(endWritableNT, stream); |
| 1705 | } else if (state.autoDestroy) { |
| 1706 | // In case of duplex streams we need a way to detect |
| 1707 | // if the writable side is ready for autoDestroy as well. |
| 1708 | const wState = stream._writableState; |
| 1709 | const autoDestroy = |
| 1710 | !wState || |
| 1711 | (wState.autoDestroy && |
| 1712 | // We don't expect the writable to ever 'finish' |
| 1713 | // if writable is explicitly set to false. |
| 1714 | (wState.finished || wState.writable === false)); |
| 1715 | |
| 1716 | if (autoDestroy) { |
| 1717 | stream.destroy(); |
| 1718 | } |
| 1719 | } |
| 1720 | } |
| 1721 | } |
| 1722 | |
| 1723 | function endWritableNT(stream) { |
| 1724 | const writable = |
| 1725 | stream.writable && !stream.writableEnded && !stream.destroyed; |
| 1726 | if (writable) { |
| 1727 | stream.end(); |
| 1728 | } |
| 1729 | } |
| 1730 | |
| 1731 | export function fromWeb(readableStream, options) { |
| 1732 | return newStreamReadableFromReadableStream(readableStream, options); |
| 1733 | } |
| 1734 | |
| 1735 | export function toWeb(streamReadable, options) { |
| 1736 | return newReadableStreamFromStreamReadable(streamReadable, options); |
| 1737 | } |
| 1738 | |
| 1739 | export function wrap(src, options) { |
| 1740 | let _ref, _src$readableObjectMo; |
| 1741 | return new Readable({ |
| 1742 | objectMode: |
| 1743 | (_ref = |
| 1744 | (_src$readableObjectMo = src.readableObjectMode) !== null && |
| 1745 | _src$readableObjectMo !== undefined |
| 1746 | ? _src$readableObjectMo |
| 1747 | : src.objectMode) !== null && _ref !== undefined |
| 1748 | ? _ref |
| 1749 | : true, |
| 1750 | ...options, |
| 1751 | destroy(err, callback) { |
| 1752 | destroyer(src, err); |
| 1753 | callback(err); |
| 1754 | }, |
| 1755 | }).wrap(src); |
| 1756 | } |
| 1757 | |
| 1758 | Readable.toWeb = toWeb; |
| 1759 | Readable.fromWeb = fromWeb; |
| 1760 | Readable.wrap = wrap; |
| 1761 | |
| 1762 | // ====================================================================================== |
| 1763 | // |
| 1764 | |
| 1765 | Readable.from = function (iterable, opts) { |
| 1766 | return from(Readable, iterable, opts); |
| 1767 | }; |
| 1768 | |
| 1769 | export function from(Readable, iterable, opts) { |
| 1770 | let iterator; |
| 1771 | if (typeof iterable === 'string' || iterable instanceof Buffer) { |
| 1772 | return new Readable({ |
| 1773 | objectMode: true, |
| 1774 | ...opts, |
| 1775 | read() { |
| 1776 | this.push(iterable); |
| 1777 | this.push(null); |
| 1778 | }, |
| 1779 | }); |
| 1780 | } |
| 1781 | let isAsync; |
| 1782 | if (iterable && iterable[Symbol.asyncIterator]) { |
| 1783 | isAsync = true; |
| 1784 | iterator = iterable[Symbol.asyncIterator](); |
| 1785 | } else if (iterable && iterable[Symbol.iterator]) { |
| 1786 | isAsync = false; |
| 1787 | iterator = iterable[Symbol.iterator](); |
| 1788 | } else { |
| 1789 | throw new ERR_INVALID_ARG_TYPE('iterable', ['Iterable'], iterable); |
| 1790 | } |
| 1791 | const readable = new Readable({ |
| 1792 | objectMode: true, |
| 1793 | highWaterMark: 1, |
| 1794 | // TODO(ronag): What options should be allowed? |
| 1795 | ...opts, |
| 1796 | }); |
| 1797 | |
| 1798 | // Flag to protect against _read |
| 1799 | // being called before last iteration completion. |
| 1800 | let reading = false; |
| 1801 | readable._read = function () { |
| 1802 | if (!reading) { |
| 1803 | reading = true; |
| 1804 | next(); |
| 1805 | } |
| 1806 | }; |
| 1807 | readable._destroy = function (error, cb) { |
| 1808 | close(error).then( |
| 1809 | () => nextTick(cb, error), |
| 1810 | (err) => nextTick(cb, err || error) |
| 1811 | ); |
| 1812 | }; |
| 1813 | async function close(error) { |
| 1814 | const hadError = error !== undefined && error !== null; |
| 1815 | const hasThrow = typeof iterator.throw === 'function'; |
| 1816 | if (hadError && hasThrow) { |
| 1817 | const { value, done } = await iterator.throw(error); |
| 1818 | await value; |
| 1819 | if (done) { |
| 1820 | return; |
| 1821 | } |
| 1822 | } |
| 1823 | if (typeof iterator.return === 'function') { |
| 1824 | const { value } = await iterator.return(); |
| 1825 | await value; |
| 1826 | } |
| 1827 | } |
| 1828 | async function next() { |
| 1829 | for (;;) { |
| 1830 | try { |
| 1831 | const { value, done } = isAsync |
| 1832 | ? await iterator.next() |
| 1833 | : iterator.next(); |
| 1834 | if (done) { |
| 1835 | readable.push(null); |
| 1836 | } else { |
| 1837 | const res = |
| 1838 | value && typeof value.then === 'function' ? await value : value; |
| 1839 | if (res === null) { |
| 1840 | reading = false; |
| 1841 | throw new ERR_STREAM_NULL_VALUES(); |
| 1842 | } else if (readable.push(res)) { |
| 1843 | continue; |
| 1844 | } else { |
| 1845 | reading = false; |
| 1846 | } |
| 1847 | } |
| 1848 | } catch (err) { |
| 1849 | readable.destroy(err); |
| 1850 | } |
| 1851 | break; |
| 1852 | } |
| 1853 | } |
| 1854 | return readable; |
| 1855 | } |
| 1856 | |
| 1857 | // ====================================================================================== |
| 1858 | // Operators |
| 1859 | |
| 1860 | const kWeakHandler = Symbol('kWeak'); |
| 1861 | const kEmpty = Symbol('kEmpty'); |
| 1862 | const kEof = Symbol('kEof'); |
| 1863 | |
| 1864 | function map(fn, options) { |
| 1865 | if (typeof fn !== 'function') { |
| 1866 | throw new ERR_INVALID_ARG_TYPE('fn', ['Function', 'AsyncFunction'], fn); |
| 1867 | } |
| 1868 | if (options != null) { |
| 1869 | validateObject(options, 'options', options); |
| 1870 | } |
| 1871 | if (options?.signal != null) { |
| 1872 | validateAbortSignal(options.signal, 'options.signal'); |
| 1873 | } |
| 1874 | let concurrency = 1; |
| 1875 | if (options?.concurrency != null) { |
| 1876 | concurrency = Math.floor(options.concurrency); |
| 1877 | } |
| 1878 | validateInteger(concurrency, 'concurrency', 1); |
| 1879 | return async function* map() { |
| 1880 | let _options$signal, _options$signal2; |
| 1881 | const ac = new globalThis.AbortController(); |
| 1882 | const stream = this; // eslint-disable-line @typescript-eslint/no-this-alias |
| 1883 | const queue = []; |
| 1884 | const signal = ac.signal; |
| 1885 | const signalOpt = { |
| 1886 | signal, |
| 1887 | }; |
| 1888 | const abort = () => ac.abort(); |
| 1889 | if ( |
| 1890 | options != null && |
| 1891 | (_options$signal = options.signal) !== null && |
| 1892 | _options$signal !== undefined && |
| 1893 | _options$signal.aborted |
| 1894 | ) { |
| 1895 | abort(); |
| 1896 | } |
| 1897 | // eslint-disable-next-line @typescript-eslint/no-unused-expressions |
| 1898 | options == null |
| 1899 | ? undefined |
| 1900 | : (_options$signal2 = options.signal) === null || |
| 1901 | _options$signal2 === undefined |
| 1902 | ? undefined |
| 1903 | : _options$signal2.addEventListener('abort', abort); |
| 1904 | let next; |
| 1905 | let resume; |
| 1906 | let done = false; |
| 1907 | function onDone() { |
| 1908 | done = true; |
| 1909 | } |
| 1910 | async function pump() { |
| 1911 | try { |
| 1912 | for await (let val of stream) { |
| 1913 | let _val; |
| 1914 | if (done) { |
| 1915 | return; |
| 1916 | } |
| 1917 | if (signal.aborted) { |
| 1918 | throw new AbortError(); |
| 1919 | } |
| 1920 | try { |
| 1921 | val = fn(val, signalOpt); |
| 1922 | } catch (err) { |
| 1923 | val = Promise.reject(err); |
| 1924 | } |
| 1925 | if (val === kEmpty) { |
| 1926 | continue; |
| 1927 | } |
| 1928 | if ( |
| 1929 | typeof ((_val = val) === null || _val === undefined |
| 1930 | ? undefined |
| 1931 | : _val.catch) === 'function' |
| 1932 | ) { |
| 1933 | val.catch(onDone); |
| 1934 | } |
| 1935 | queue.push(val); |
| 1936 | if (next) { |
| 1937 | next(); |
| 1938 | next = null; |
| 1939 | } |
| 1940 | if (!done && queue.length && queue.length >= concurrency) { |
| 1941 | await new Promise((resolve) => { |
| 1942 | resume = resolve; |
| 1943 | }); |
| 1944 | } |
| 1945 | } |
| 1946 | queue.push(kEof); |
| 1947 | } catch (err) { |
| 1948 | const val = Promise.reject(err); |
| 1949 | val.then(undefined, onDone); |
| 1950 | queue.push(val); |
| 1951 | } finally { |
| 1952 | let _options$signal3; |
| 1953 | done = true; |
| 1954 | if (next) { |
| 1955 | next(); |
| 1956 | next = null; |
| 1957 | } |
| 1958 | // eslint-disable-next-line @typescript-eslint/no-unused-expressions |
| 1959 | options == null |
| 1960 | ? undefined |
| 1961 | : (_options$signal3 = options.signal) === null || |
| 1962 | _options$signal3 === undefined |
| 1963 | ? undefined |
| 1964 | : _options$signal3.removeEventListener('abort', abort); |
| 1965 | } |
| 1966 | } |
| 1967 | pump(); |
| 1968 | try { |
| 1969 | while (true) { |
| 1970 | while (queue.length > 0) { |
| 1971 | const val = await queue[0]; |
| 1972 | if (val === kEof) { |
| 1973 | return; |
| 1974 | } |
| 1975 | if (signal.aborted) { |
| 1976 | throw new AbortError(); |
| 1977 | } |
| 1978 | if (val !== kEmpty) { |
| 1979 | yield val; |
| 1980 | } |
| 1981 | queue.shift(); |
| 1982 | if (resume) { |
| 1983 | resume(); |
| 1984 | resume = null; |
| 1985 | } |
| 1986 | } |
| 1987 | await new Promise((resolve) => { |
| 1988 | next = resolve; |
| 1989 | }); |
| 1990 | } |
| 1991 | } finally { |
| 1992 | ac.abort(); |
| 1993 | done = true; |
| 1994 | if (resume) { |
| 1995 | resume(); |
| 1996 | resume = null; |
| 1997 | } |
| 1998 | } |
| 1999 | }.call(this); |
| 2000 | } |
| 2001 | |
| 2002 | function asIndexedPairs(options) { |
| 2003 | if (options != null) { |
| 2004 | validateObject(options, 'options', options); |
| 2005 | } |
| 2006 | if ((options == null ? undefined : options.signal) != null) { |
| 2007 | validateAbortSignal(options.signal, 'options.signal'); |
| 2008 | } |
| 2009 | return async function* asIndexedPairs() { |
| 2010 | let index = 0; |
| 2011 | for await (const val of this) { |
| 2012 | let _options$signal4; |
| 2013 | if ( |
| 2014 | options !== null && |
| 2015 | options !== undefined && |
| 2016 | (_options$signal4 = options.signal) !== null && |
| 2017 | _options$signal4 !== undefined && |
| 2018 | _options$signal4.aborted |
| 2019 | ) { |
| 2020 | throw new AbortError('Aborted', { |
| 2021 | cause: options.signal?.reason, |
| 2022 | }); |
| 2023 | } |
| 2024 | yield [index++, val]; |
| 2025 | } |
| 2026 | }.call(this); |
| 2027 | } |
| 2028 | |
| 2029 | async function some(fn, options) { |
| 2030 | if (typeof fn !== 'function') { |
| 2031 | throw new ERR_INVALID_ARG_TYPE('fn', ['Function', 'AsyncFunction'], fn); |
| 2032 | } |
| 2033 | for await (const _ of filter.call(this, fn, options)) { |
| 2034 | return true; |
| 2035 | } |
| 2036 | return false; |
| 2037 | } |
| 2038 | |
| 2039 | async function every(fn, options) { |
| 2040 | if (typeof fn !== 'function') { |
| 2041 | throw new ERR_INVALID_ARG_TYPE('fn', ['Function', 'AsyncFunction'], fn); |
| 2042 | } |
| 2043 | // https://en.wikipedia.org/wiki/De_Morgan%27s_laws |
| 2044 | return !(await some.call( |
| 2045 | this, |
| 2046 | async (...args) => { |
| 2047 | return !(await fn(...args)); |
| 2048 | }, |
| 2049 | options |
| 2050 | )); |
| 2051 | } |
| 2052 | |
| 2053 | async function find(fn, options) { |
| 2054 | for await (const result of filter.call(this, fn, options)) { |
| 2055 | return result; |
| 2056 | } |
| 2057 | return undefined; |
| 2058 | } |
| 2059 | |
| 2060 | async function forEach(fn, options) { |
| 2061 | if (typeof fn !== 'function') { |
| 2062 | throw new ERR_INVALID_ARG_TYPE('fn', ['Function', 'AsyncFunction'], fn); |
| 2063 | } |
| 2064 | async function forEachFn(value, options) { |
| 2065 | await fn(value, options); |
| 2066 | return kEmpty; |
| 2067 | } |
| 2068 | for await (const _ of map.call(this, forEachFn, options)); |
| 2069 | } |
| 2070 | |
| 2071 | function filter(fn, options) { |
| 2072 | if (typeof fn !== 'function') { |
| 2073 | throw new ERR_INVALID_ARG_TYPE('fn', ['Function', 'AsyncFunction'], fn); |
| 2074 | } |
| 2075 | async function filterFn(value, options) { |
| 2076 | if (await fn(value, options)) { |
| 2077 | return value; |
| 2078 | } |
| 2079 | return kEmpty; |
| 2080 | } |
| 2081 | return map.call(this, filterFn, options); |
| 2082 | } |
| 2083 | |
| 2084 | // Specific to provide better error to reduce since the argument is only |
| 2085 | // missing if the stream has no items in it - but the code is still appropriate |
| 2086 | class ReduceAwareErrMissingArgs extends ERR_MISSING_ARGS { |
| 2087 | constructor() { |
| 2088 | super('reduce'); |
| 2089 | this.message = 'Reduce of an empty stream requires an initial value'; |
| 2090 | } |
| 2091 | } |
| 2092 | |
| 2093 | async function reduce(reducer, initialValue, options) { |
| 2094 | let _options$signal5; |
| 2095 | if (typeof reducer !== 'function') { |
| 2096 | throw new ERR_INVALID_ARG_TYPE( |
| 2097 | 'reducer', |
| 2098 | ['Function', 'AsyncFunction'], |
| 2099 | reducer |
| 2100 | ); |
| 2101 | } |
| 2102 | if (options != null) { |
| 2103 | validateObject(options, 'options', options); |
| 2104 | } |
| 2105 | if (options?.signal != null) { |
| 2106 | validateAbortSignal(options?.signal, 'options.signal'); |
| 2107 | } |
| 2108 | let hasInitialValue = arguments.length > 1; |
| 2109 | if ( |
| 2110 | options !== null && |
| 2111 | options !== undefined && |
| 2112 | (_options$signal5 = options.signal) !== null && |
| 2113 | _options$signal5 !== undefined && |
| 2114 | _options$signal5.aborted |
| 2115 | ) { |
| 2116 | const err = new AbortError(undefined, { |
| 2117 | cause: options.signal?.reason, |
| 2118 | }); |
| 2119 | this.once('error', () => {}); // The error is already propagated |
| 2120 | await finished(this.destroy(err)); |
| 2121 | throw err; |
| 2122 | } |
| 2123 | const ac = new globalThis.AbortController(); |
| 2124 | const signal = ac.signal; |
| 2125 | if (options?.signal) { |
| 2126 | const opts = { |
| 2127 | once: true, |
| 2128 | [kWeakHandler]: this, |
| 2129 | }; |
| 2130 | options.signal.addEventListener('abort', () => ac.abort(), opts); |
| 2131 | } |
| 2132 | let gotAnyItemFromStream = false; |
| 2133 | try { |
| 2134 | for await (const value of this) { |
| 2135 | let _options$signal6; |
| 2136 | gotAnyItemFromStream = true; |
| 2137 | if ( |
| 2138 | options !== null && |
| 2139 | options !== undefined && |
| 2140 | (_options$signal6 = options.signal) !== null && |
| 2141 | _options$signal6 !== undefined && |
| 2142 | _options$signal6.aborted |
| 2143 | ) { |
| 2144 | throw new AbortError(); |
| 2145 | } |
| 2146 | if (!hasInitialValue) { |
| 2147 | initialValue = value; |
| 2148 | hasInitialValue = true; |
| 2149 | } else { |
| 2150 | initialValue = await reducer(initialValue, value, { |
| 2151 | signal, |
| 2152 | }); |
| 2153 | } |
| 2154 | } |
| 2155 | if (!gotAnyItemFromStream && !hasInitialValue) { |
| 2156 | throw new ReduceAwareErrMissingArgs(); |
| 2157 | } |
| 2158 | } finally { |
| 2159 | ac.abort(); |
| 2160 | } |
| 2161 | return initialValue; |
| 2162 | } |
| 2163 | |
| 2164 | async function toArray(options) { |
| 2165 | if (options != null) { |
| 2166 | validateObject(options, 'options', options); |
| 2167 | } |
| 2168 | if (options?.signal != null) { |
| 2169 | validateAbortSignal(options?.signal, 'options.signal'); |
| 2170 | } |
| 2171 | const result = []; |
| 2172 | for await (const val of this) { |
| 2173 | let _options$signal7; |
| 2174 | if ( |
| 2175 | options !== null && |
| 2176 | options !== undefined && |
| 2177 | (_options$signal7 = options.signal) !== null && |
| 2178 | _options$signal7 !== undefined && |
| 2179 | _options$signal7.aborted |
| 2180 | ) { |
| 2181 | throw new AbortError(undefined, { |
| 2182 | cause: options.signal?.reason, |
| 2183 | }); |
| 2184 | } |
| 2185 | result.push(val); |
| 2186 | } |
| 2187 | return result; |
| 2188 | } |
| 2189 | |
| 2190 | function flatMap(fn, options) { |
| 2191 | const values = map.call(this, fn, options); |
| 2192 | return async function* flatMap() { |
| 2193 | for await (const val of values) { |
| 2194 | yield* val; |
| 2195 | } |
| 2196 | }.call(this); |
| 2197 | } |
| 2198 | |
| 2199 | function toIntegerOrInfinity(number) { |
| 2200 | // We coerce here to align with the spec |
| 2201 | // https://github.com/tc39/proposal-iterator-helpers/issues/169 |
| 2202 | number = Number(number); |
| 2203 | if (Number.isNaN(number)) { |
| 2204 | return 0; |
| 2205 | } |
| 2206 | if (number < 0) { |
| 2207 | throw new ERR_OUT_OF_RANGE('number', '>= 0', number); |
| 2208 | } |
| 2209 | return number; |
| 2210 | } |
| 2211 | |
| 2212 | function drop(number, options) { |
| 2213 | if (options != null) { |
| 2214 | validateObject(options, 'options', options); |
| 2215 | } |
| 2216 | if (options?.signal != null) { |
| 2217 | validateAbortSignal(options?.signal, 'options.signal'); |
| 2218 | } |
| 2219 | number = toIntegerOrInfinity(number); |
| 2220 | return async function* drop() { |
| 2221 | let _options$signal8; |
| 2222 | if ( |
| 2223 | options !== null && |
| 2224 | options !== undefined && |
| 2225 | (_options$signal8 = options.signal) !== null && |
| 2226 | _options$signal8 !== undefined && |
| 2227 | _options$signal8.aborted |
| 2228 | ) { |
| 2229 | throw new AbortError(); |
| 2230 | } |
| 2231 | for await (const val of this) { |
| 2232 | let _options$signal9; |
| 2233 | if ( |
| 2234 | options !== null && |
| 2235 | options !== undefined && |
| 2236 | (_options$signal9 = options.signal) !== null && |
| 2237 | _options$signal9 !== undefined && |
| 2238 | _options$signal9.aborted |
| 2239 | ) { |
| 2240 | throw new AbortError(); |
| 2241 | } |
| 2242 | if (number-- <= 0) { |
| 2243 | yield val; |
| 2244 | } |
| 2245 | } |
| 2246 | }.call(this); |
| 2247 | } |
| 2248 | |
| 2249 | function take(number, options) { |
| 2250 | if (options != null) { |
| 2251 | validateObject(options, 'options', options); |
| 2252 | } |
| 2253 | if (options?.signal != null) { |
| 2254 | validateAbortSignal(options?.signal, 'options.signal'); |
| 2255 | } |
| 2256 | number = toIntegerOrInfinity(number); |
| 2257 | return async function* take() { |
| 2258 | let _options$signal10; |
| 2259 | if ( |
| 2260 | options !== null && |
| 2261 | options !== undefined && |
| 2262 | (_options$signal10 = options.signal) !== null && |
| 2263 | _options$signal10 !== undefined && |
| 2264 | _options$signal10.aborted |
| 2265 | ) { |
| 2266 | throw new AbortError(); |
| 2267 | } |
| 2268 | for await (const val of this) { |
| 2269 | let _options$signal11; |
| 2270 | if ( |
| 2271 | options !== null && |
| 2272 | options !== undefined && |
| 2273 | (_options$signal11 = options.signal) !== null && |
| 2274 | _options$signal11 !== undefined && |
| 2275 | _options$signal11.aborted |
| 2276 | ) { |
| 2277 | throw new AbortError(); |
| 2278 | } |
| 2279 | if (number-- > 0) { |
| 2280 | yield val; |
| 2281 | } else { |
| 2282 | return; |
| 2283 | } |
| 2284 | } |
| 2285 | }.call(this); |
| 2286 | } |
| 2287 | |
| 2288 | Readable.prototype.map = function (fn, options) { |
| 2289 | return from(Readable, map.call(this, fn, options)); |
| 2290 | }; |
| 2291 | |
| 2292 | Readable.prototype.asIndexedPairs = function (options) { |
| 2293 | return from(Readable, asIndexedPairs.call(this, options)); |
| 2294 | }; |
| 2295 | |
| 2296 | Readable.prototype.drop = function (number, options) { |
| 2297 | return from(Readable, drop.call(this, number, options)); |
| 2298 | }; |
| 2299 | |
| 2300 | Readable.prototype.filter = function (fn, options) { |
| 2301 | return from(Readable, filter.call(this, fn, options)); |
| 2302 | }; |
| 2303 | |
| 2304 | Readable.prototype.flatMap = function (fn, options) { |
| 2305 | return from(Readable, flatMap.call(this, fn, options)); |
| 2306 | }; |
| 2307 | |
| 2308 | Readable.prototype.take = function (number, options) { |
| 2309 | return from(Readable, take.call(this, number, options)); |
| 2310 | }; |
| 2311 | |
| 2312 | Readable.prototype.every = every; |
| 2313 | Readable.prototype.forEach = forEach; |
| 2314 | Readable.prototype.reduce = reduce; |
| 2315 | Readable.prototype.toArray = toArray; |
| 2316 | Readable.prototype.some = some; |
| 2317 | Readable.prototype.find = find; |
| 2318 | |
| 2319 | /** |
| 2320 | * @typedef {import('./queuingstrategies').QueuingStrategy} QueuingStrategy |
| 2321 | * @param {Readable} streamReadable |
| 2322 | * @param {{ |
| 2323 | * strategy : QueuingStrategy |
| 2324 | * }} [options] |
| 2325 | * @returns {ReadableStream} |
| 2326 | */ |
| 2327 | export function newReadableStreamFromStreamReadable( |
| 2328 | streamReadable, |
| 2329 | options = {}, |
| 2330 | createTypeBytes = false |
| 2331 | ) { |
| 2332 | // Not using the internal/streams/utils isReadableNodeStream utility |
| 2333 | // here because it will return false if streamReadable is a Duplex |
| 2334 | // whose readable option is false. For a Duplex that is not readable, |
| 2335 | // we want it to pass this check but return a closed ReadableStream. |
| 2336 | if (typeof streamReadable?._readableState !== 'object') { |
| 2337 | throw new ERR_INVALID_ARG_TYPE( |
| 2338 | 'streamReadable', |
| 2339 | 'stream.Readable', |
| 2340 | streamReadable |
| 2341 | ); |
| 2342 | } |
| 2343 | |
| 2344 | if (isDestroyed(streamReadable) || !isReadable(streamReadable)) { |
| 2345 | const readable = new globalThis.ReadableStream(); |
| 2346 | readable.cancel(); |
| 2347 | return readable; |
| 2348 | } |
| 2349 | |
| 2350 | const objectMode = streamReadable.readableObjectMode; |
| 2351 | const highWaterMark = streamReadable.readableHighWaterMark; |
| 2352 | |
| 2353 | const evaluateStrategyOrFallback = (strategy) => { |
| 2354 | // If there is a strategy available, use it |
| 2355 | if (strategy) return strategy; |
| 2356 | |
| 2357 | if (objectMode) { |
| 2358 | // When running in objectMode explicitly but no strategy, we just fall |
| 2359 | // back to CountQueuingStrategy |
| 2360 | return new globalThis.CountQueuingStrategy({ highWaterMark }); |
| 2361 | } |
| 2362 | |
| 2363 | return new globalThis.ByteLengthQueuingStrategy({ highWaterMark }); |
| 2364 | }; |
| 2365 | |
| 2366 | const strategy = evaluateStrategyOrFallback(options?.strategy); |
| 2367 | |
| 2368 | let controller; |
| 2369 | let wasCanceled = false; |
| 2370 | |
| 2371 | function onData(chunk) { |
| 2372 | // Copy the Buffer to detach it from the pool. |
| 2373 | if (Buffer.isBuffer(chunk) && !objectMode) chunk = new Uint8Array(chunk); |
| 2374 | controller.enqueue(chunk); |
| 2375 | if (controller.desiredSize <= 0) streamReadable.pause(); |
| 2376 | } |
| 2377 | |
| 2378 | streamReadable.pause(); |
| 2379 | |
| 2380 | const cleanup = eos(streamReadable, (error) => { |
| 2381 | error = handleKnownInternalErrors(error); |
| 2382 | |
| 2383 | cleanup(); |
| 2384 | // This is a protection against non-standard, legacy streams |
| 2385 | // that happen to emit an error event again after finished is called. |
| 2386 | streamReadable.on('error', () => {}); |
| 2387 | if (error) return controller.error(error); |
| 2388 | // Was already canceled |
| 2389 | if (wasCanceled) { |
| 2390 | return; |
| 2391 | } |
| 2392 | controller.close(); |
| 2393 | }); |
| 2394 | |
| 2395 | streamReadable.on('data', onData); |
| 2396 | |
| 2397 | return new globalThis.ReadableStream( |
| 2398 | { |
| 2399 | start(c) { |
| 2400 | controller = c; |
| 2401 | }, |
| 2402 | |
| 2403 | pull() { |
| 2404 | streamReadable.resume(); |
| 2405 | }, |
| 2406 | |
| 2407 | cancel(reason) { |
| 2408 | wasCanceled = true; |
| 2409 | destroy(streamReadable, reason); |
| 2410 | }, |
| 2411 | type: createTypeBytes ? 'bytes' : undefined, |
| 2412 | }, |
| 2413 | strategy |
| 2414 | ); |
| 2415 | } |
| 2416 | |
| 2417 | /** |
| 2418 | * @param {ReadableStream} readableStream |
| 2419 | * @param {{ |
| 2420 | * highWaterMark? : number, |
| 2421 | * encoding? : string, |
| 2422 | * objectMode? : boolean, |
| 2423 | * signal? : AbortSignal, |
| 2424 | * }} [options] |
| 2425 | * @returns {Readable} |
| 2426 | */ |
| 2427 | export function newStreamReadableFromReadableStream( |
| 2428 | readableStream, |
| 2429 | options = {} |
| 2430 | ) { |
| 2431 | if (!isReadableStream(readableStream)) { |
| 2432 | throw new ERR_INVALID_ARG_TYPE( |
| 2433 | 'readableStream', |
| 2434 | 'ReadableStream', |
| 2435 | readableStream |
| 2436 | ); |
| 2437 | } |
| 2438 | |
| 2439 | validateObject(options, 'options'); |
| 2440 | const { highWaterMark, encoding, objectMode = false, signal } = options; |
| 2441 | |
| 2442 | if (encoding !== undefined && !Buffer.isEncoding(encoding)) |
| 2443 | throw new ERR_INVALID_ARG_VALUE('options.encoding', encoding); |
| 2444 | validateBoolean(objectMode, 'options.objectMode'); |
| 2445 | |
| 2446 | const reader = readableStream.getReader(); |
| 2447 | let closed = false; |
| 2448 | |
| 2449 | const readable = new Readable({ |
| 2450 | objectMode, |
| 2451 | highWaterMark, |
| 2452 | encoding, |
| 2453 | signal, |
| 2454 | |
| 2455 | read() { |
| 2456 | reader.read().then( |
| 2457 | (chunk) => { |
| 2458 | if (chunk.done) { |
| 2459 | // Value should always be undefined here. |
| 2460 | readable.push(null); |
| 2461 | } else { |
| 2462 | readable.push(chunk.value); |
| 2463 | } |
| 2464 | }, |
| 2465 | (error) => destroy.call(readable, error) |
| 2466 | ); |
| 2467 | }, |
| 2468 | |
| 2469 | destroy(error, callback) { |
| 2470 | function done() { |
| 2471 | try { |
| 2472 | callback(error); |
| 2473 | } catch (error) { |
| 2474 | // In a next tick because this is happening within |
| 2475 | // a promise context, and if there are any errors |
| 2476 | // thrown we don't want those to cause an unhandled |
| 2477 | // rejection. Let's just escape the promise and |
| 2478 | // handle it separately. |
| 2479 | nextTick(() => { |
| 2480 | throw error; |
| 2481 | }); |
| 2482 | } |
| 2483 | } |
| 2484 | |
| 2485 | if (!closed) { |
| 2486 | reader.cancel(error).then(done, done); |
| 2487 | return; |
| 2488 | } |
| 2489 | done(); |
| 2490 | }, |
| 2491 | }); |
| 2492 | |
| 2493 | reader.closed.then( |
| 2494 | () => { |
| 2495 | closed = true; |
| 2496 | }, |
| 2497 | (error) => { |
| 2498 | closed = true; |
| 2499 | destroy.call(readable, error); |
| 2500 | } |
| 2501 | ); |
| 2502 | |
| 2503 | return readable; |
| 2504 | } |