Skip to content
File

Blob: src/node/internal/streams_readable.js

javascript2505 lines
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4//
5// 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 
26import {
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';
44import { nextTick } from 'node-internal:internal_process';
45import {
46 destroy,
47 undestroy,
48 errorOrDestroy,
49 destroyer,
50 construct,
51} from 'node-internal:streams_destroy';
52import { eos, finished, nop } from 'node-internal:streams_end_of_stream';
53import {
54 getHighWaterMark,
55 getDefaultHighWaterMark,
56} from 'node-internal:streams_state';
57import { addAbortSignal } from 'node-internal:streams_add_abort_signal';
58import { EventEmitter } from 'node-internal:events';
59import { Stream } from 'node-internal:streams_legacy';
60import { Buffer } from 'node-internal:internal_buffer';
61 
62import {
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 
76import {
77 validateObject,
78 validateAbortSignal,
79 validateBoolean,
80 validateInteger,
81} from 'node-internal:validators';
82 
83import { StringDecoder } from 'node-internal:internal_stringdecoder';
84 
85const streamsNodejsV24Compat =
86 Cloudflare.compatibilityFlags.enable_streams_nodejs_v24_compat;
87 
88const kErroredValue = Symbol('kErroredValue');
89const kDefaultEncodingValue = Symbol('kDefaultEncodingValue');
90const kDecoderValue = Symbol('kDecoderValue');
91const 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).
96const kEnded = 1 << 9;
97const kEndEmitted = 1 << 10;
98const kReading = 1 << 11;
99const kSync = 1 << 12;
100const kNeedReadable = 1 << 13;
101const kEmittedReadable = 1 << 14;
102const kReadableListening = 1 << 15;
103const kResumeScheduled = 1 << 16;
104const kMultiAwaitDrain = 1 << 17;
105const kReadingMore = 1 << 18;
106const kDataEmitted = 1 << 19;
107const kDefaultUTF8Encoding = 1 << 20;
108const kDecoder = 1 << 21;
109const kEncoding = 1 << 22;
110const kHasFlowing = 1 << 23;
111const kFlowing = 1 << 24;
112const kHasPaused = 1 << 25;
113const kPaused = 1 << 26;
114const kDataListening = 1 << 27;
115 
116// ======================================================================================
117// ReadableState
118 
119// TODO(benjamingr) it is likely slower to do it this way than with free functions
120function 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}
132Object.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 
260export 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 
317ReadableState.prototype[kOnConstructed] = function onConstructed(stream) {
318 if ((this[kState] & kNeedReadable) !== 0) {
319 maybeReadMore(stream, this);
320 }
321};
322 
323// ======================================================================================
324// Readable
325 
326Readable.ReadableState = ReadableState;
327 
328Object.setPrototypeOf(Readable.prototype, Stream.prototype);
329Object.setPrototypeOf(Readable, Stream);
330 
331export 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}
370Readable.prototype.destroy = destroy;
371Readable.prototype._undestroy = undestroy;
372Readable.prototype._destroy = function (err, cb) {
373 if (cb) cb(err);
374};
375 
376Readable.prototype[EventEmitter.captureRejectionSymbol] = function (err) {
377 this.destroy(err);
378};
379 
380Readable.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.
395Readable.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().
403Readable.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 
410function 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 
450function 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 
461function 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 
470function 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 
529function 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 
555function 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 
565function 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 
599Readable.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.
608Readable.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.
630const MAX_HWM = 0x40000000;
631function 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.
650function 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.
665Readable.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 
795function 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.
826function 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 
835function 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.
865function maybeReadMore(stream, state) {
866 if ((state[kState] & (kReadingMore | kConstructed)) === kConstructed) {
867 state[kState] |= kReadingMore;
868 nextTick(maybeReadMore_, stream, state);
869 }
870}
871 
872function 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.
914Readable.prototype._read = function (_size) {
915 throw new ERR_METHOD_NOT_IMPLEMENTED('_read()');
916};
917 
918Readable.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 
1077function 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 
1099Readable.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.
1129Readable.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};
1159Readable.prototype.addListener = Readable.prototype.on;
1160 
1161Readable.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};
1180Readable.prototype.off = Readable.prototype.removeListener;
1181 
1182Readable.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 
1198function 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 
1223function 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.
1229Readable.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 
1248function resume(stream, state) {
1249 if ((state[kState] & kResumeScheduled) === 0) {
1250 state[kState] |= kResumeScheduled;
1251 nextTick(resume_, stream, state);
1252 }
1253}
1254 
1255function 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 
1266Readable.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 
1277function 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.
1285Readable.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 
1334Readable.prototype[Symbol.asyncIterator] = function () {
1335 return streamToAsyncIterator(this);
1336};
1337 
1338Readable.prototype.iterator = function (options) {
1339 if (options !== undefined) {
1340 validateObject(options, 'options');
1341 }
1342 return streamToAsyncIterator(this, options);
1343};
1344 
1345function 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 
1355async 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.
1408Object.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 
1549Object.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.
1576Readable._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.
1582function 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 
1685function 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 
1694function 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 
1723function endWritableNT(stream) {
1724 const writable =
1725 stream.writable && !stream.writableEnded && !stream.destroyed;
1726 if (writable) {
1727 stream.end();
1728 }
1729}
1730 
1731export function fromWeb(readableStream, options) {
1732 return newStreamReadableFromReadableStream(readableStream, options);
1733}
1734 
1735export function toWeb(streamReadable, options) {
1736 return newReadableStreamFromStreamReadable(streamReadable, options);
1737}
1738 
1739export 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 
1758Readable.toWeb = toWeb;
1759Readable.fromWeb = fromWeb;
1760Readable.wrap = wrap;
1761 
1762// ======================================================================================
1763//
1764 
1765Readable.from = function (iterable, opts) {
1766 return from(Readable, iterable, opts);
1767};
1768 
1769export 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 
1860const kWeakHandler = Symbol('kWeak');
1861const kEmpty = Symbol('kEmpty');
1862const kEof = Symbol('kEof');
1863 
1864function 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 
2002function 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 
2029async 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 
2039async 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 
2053async function find(fn, options) {
2054 for await (const result of filter.call(this, fn, options)) {
2055 return result;
2056 }
2057 return undefined;
2058}
2059 
2060async 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 
2071function 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
2086class 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 
2093async 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 
2164async 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 
2190function 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 
2199function 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 
2212function 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 
2249function 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 
2288Readable.prototype.map = function (fn, options) {
2289 return from(Readable, map.call(this, fn, options));
2290};
2291 
2292Readable.prototype.asIndexedPairs = function (options) {
2293 return from(Readable, asIndexedPairs.call(this, options));
2294};
2295 
2296Readable.prototype.drop = function (number, options) {
2297 return from(Readable, drop.call(this, number, options));
2298};
2299 
2300Readable.prototype.filter = function (fn, options) {
2301 return from(Readable, filter.call(this, fn, options));
2302};
2303 
2304Readable.prototype.flatMap = function (fn, options) {
2305 return from(Readable, flatMap.call(this, fn, options));
2306};
2307 
2308Readable.prototype.take = function (number, options) {
2309 return from(Readable, take.call(this, number, options));
2310};
2311 
2312Readable.prototype.every = every;
2313Readable.prototype.forEach = forEach;
2314Readable.prototype.reduce = reduce;
2315Readable.prototype.toArray = toArray;
2316Readable.prototype.some = some;
2317Readable.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 */
2327export 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 */
2427export 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}