Skip to content
File

Blob: src/node/internal/streams_duplex.js

javascript930 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 
26/* TODO: the following is adopted code, enabling linting one day */
27/* eslint-disable */
28 
29import { Buffer } from 'node-internal:internal_buffer';
30import {
31 Readable,
32 newReadableStreamFromStreamReadable,
33} from 'node-internal:streams_readable';
34import {
35 Writable,
36 newWritableStreamFromStreamWritable,
37} from 'node-internal:streams_writable';
38import { ok as assert } from 'node-internal:internal_assert';
39import { Stream } from 'node-internal:streams_legacy';
40import { nextTick } from 'node-internal:internal_process';
41import { validateBoolean, validateObject } from 'node-internal:validators';
42import { normalizeEncoding } from 'node-internal:internal_utils';
43import { addAbortSignal } from 'node-internal:streams_add_abort_signal';
44 
45import {
46 isDestroyed,
47 isReadable,
48 isWritable,
49 isIterable,
50 isNodeStream,
51 isWritableEnded,
52 isReadableNodeStream,
53 isWritableNodeStream,
54 isDuplexNodeStream,
55 kOnConstructed,
56} from 'node-internal:streams_util';
57import {
58 construct as destroyConstruct,
59 destroyer,
60} from 'node-internal:streams_destroy';
61import { eos } from 'node-internal:streams_end_of_stream';
62 
63import {
64 AbortError,
65 ERR_INVALID_ARG_TYPE,
66 ERR_INVALID_ARG_VALUE,
67 ERR_INVALID_RETURN_VALUE,
68 ERR_STREAM_PREMATURE_CLOSE,
69} from 'node-internal:internal_errors';
70 
71/**
72 * @typedef {import('./readablestream').ReadableWritablePair
73 * } ReadableWritablePair
74 * @typedef {import('../../stream').Duplex} Duplex
75 */
76const encoder = new TextEncoder();
77 
78Object.setPrototypeOf(Duplex.prototype, Readable.prototype);
79Object.setPrototypeOf(Duplex, Readable);
80{
81 const keys = Object.keys(Writable.prototype);
82 // Allow the keys array to be GC'ed.
83 for (let i = 0; i < keys.length; i++) {
84 const method = keys[i];
85 Duplex.prototype[method] ||= Writable.prototype[method];
86 }
87}
88 
89// Use the `destroy` method of `Writable`.
90Duplex.prototype.destroy = Writable.prototype.destroy;
91 
92export function Duplex(options) {
93 if (!(this instanceof Duplex)) return new Duplex(options);
94 
95 this._events ??= {
96 close: undefined,
97 error: undefined,
98 prefinish: undefined,
99 finish: undefined,
100 drain: undefined,
101 data: undefined,
102 end: undefined,
103 readable: undefined,
104 // Skip uncommon events...
105 // pause: undefined,
106 // resume: undefined,
107 // pipe: undefined,
108 // unpipe: undefined,
109 // [destroyImpl.kConstruct]: undefined,
110 // [destroyImpl.kDestroy]: undefined,
111 };
112 
113 this._readableState = new Readable.ReadableState(options, this, true);
114 this._writableState = new Writable.WritableState(options, this, true);
115 
116 if (options) {
117 this.allowHalfOpen = options.allowHalfOpen !== false;
118 
119 if (options.readable === false) {
120 this._readableState.readable = false;
121 this._readableState.ended = true;
122 this._readableState.endEmitted = true;
123 }
124 
125 if (options.writable === false) {
126 this._writableState.writable = false;
127 this._writableState.ending = true;
128 this._writableState.ended = true;
129 this._writableState.finished = true;
130 }
131 
132 if (typeof options.read === 'function') this._read = options.read;
133 
134 if (typeof options.write === 'function') this._write = options.write;
135 
136 if (typeof options.writev === 'function') this._writev = options.writev;
137 
138 if (typeof options.destroy === 'function') this._destroy = options.destroy;
139 
140 if (typeof options.final === 'function') this._final = options.final;
141 
142 if (typeof options.construct === 'function')
143 this._construct = options.construct;
144 
145 if (options.signal) {
146 addAbortSignal(options.signal, this);
147 }
148 } else {
149 this.allowHalfOpen = true;
150 }
151 
152 Stream.call(this, options);
153 
154 if (this._construct != null) {
155 destroyConstruct(this, () => {
156 this._readableState[kOnConstructed](this);
157 this._writableState[kOnConstructed](this);
158 });
159 }
160}
161 
162// Use the `destroy` method of `Writable`.
163Duplex.prototype.destroy = Writable.prototype.destroy;
164 
165Object.defineProperties(Duplex.prototype, {
166 writable: {
167 __proto__: null,
168 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writable'),
169 },
170 writableHighWaterMark: {
171 __proto__: null,
172 ...Object.getOwnPropertyDescriptor(
173 Writable.prototype,
174 'writableHighWaterMark'
175 ),
176 },
177 writableObjectMode: {
178 __proto__: null,
179 ...Object.getOwnPropertyDescriptor(
180 Writable.prototype,
181 'writableObjectMode'
182 ),
183 },
184 writableBuffer: {
185 __proto__: null,
186 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writableBuffer'),
187 },
188 writableLength: {
189 __proto__: null,
190 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writableLength'),
191 },
192 writableFinished: {
193 __proto__: null,
194 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writableFinished'),
195 },
196 writableCorked: {
197 __proto__: null,
198 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writableCorked'),
199 },
200 writableEnded: {
201 __proto__: null,
202 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writableEnded'),
203 },
204 writableNeedDrain: {
205 __proto__: null,
206 ...Object.getOwnPropertyDescriptor(Writable.prototype, 'writableNeedDrain'),
207 },
208 
209 destroyed: {
210 __proto__: null,
211 get() {
212 if (
213 this._readableState === undefined ||
214 this._writableState === undefined
215 ) {
216 return false;
217 }
218 return this._readableState.destroyed && this._writableState.destroyed;
219 },
220 set(value) {
221 // Backward compatibility, the user is explicitly
222 // managing destroyed.
223 if (this._readableState && this._writableState) {
224 this._readableState.destroyed = value;
225 this._writableState.destroyed = value;
226 }
227 },
228 },
229});
230 
231export function fromWeb(pair, options) {
232 return newStreamDuplexFromReadableWritablePair(pair, options);
233}
234 
235export function toWeb(duplex) {
236 return newReadableWritablePairFromDuplex(duplex);
237}
238 
239export function toBYOBWeb(duplex) {
240 return newReadableWritablePairFromDuplex(duplex, true /* createTypeBytes */);
241}
242 
243export function from(body) {
244 return duplexify(body, 'body');
245}
246 
247Duplex.fromWeb = fromWeb;
248Duplex.toWeb = toWeb;
249Duplex.from = from;
250 
251// ======================================================================================
252 
253function isBlob(b) {
254 return b instanceof Blob;
255}
256 
257// This is needed for pre node 17.
258class Duplexify extends Duplex {
259 constructor(options) {
260 super(options);
261 // https://github.com/nodejs/node/pull/34385
262 
263 if (
264 (options === null || options === undefined
265 ? undefined
266 : options.readable) === false
267 ) {
268 this['_readableState'].readable = false;
269 this['_readableState'].ended = true;
270 this['_readableState'].endEmitted = true;
271 }
272 if (
273 (options === null || options === undefined
274 ? undefined
275 : options.writable) === false
276 ) {
277 this['_readableState'].writable = false;
278 this['_readableState'].ending = true;
279 this['_readableState'].ended = true;
280 this['_readableState'].finished = true;
281 }
282 }
283}
284 
285function duplexify(body, name) {
286 if (isDuplexNodeStream(body)) {
287 return body;
288 }
289 if (isReadableNodeStream(body)) {
290 return _duplexify({
291 readable: body,
292 });
293 }
294 if (isWritableNodeStream(body)) {
295 return _duplexify({
296 writable: body,
297 });
298 }
299 if (isNodeStream(body)) {
300 return _duplexify({
301 writable: false,
302 readable: false,
303 });
304 }
305 
306 if (body instanceof ReadableStream) {
307 return _duplexify({ readable: Readable.fromWeb(body) });
308 }
309 
310 if (body instanceof WritableStream) {
311 return _duplexify({ writable: Writable.fromWeb(body) });
312 }
313 
314 if (typeof body === 'function') {
315 const { value, write, final, destroy } = fromAsyncGen(body);
316 if (isIterable(value)) {
317 return Readable.from(Duplexify, value, {
318 // TODO (ronag): highWaterMark?
319 objectMode: true,
320 write,
321 final,
322 destroy,
323 });
324 }
325 const then = value.then;
326 if (typeof then === 'function') {
327 let d;
328 const promise = Reflect.apply(then, value, [
329 (val) => {
330 if (val != null) {
331 throw new ERR_INVALID_RETURN_VALUE('nully', 'body', val);
332 }
333 },
334 (err) => {
335 destroyer(d, err);
336 },
337 ]);
338 
339 return (d = new Duplexify({
340 // TODO (ronag): highWaterMark?
341 objectMode: true,
342 readable: false,
343 write,
344 final(cb) {
345 final(async () => {
346 try {
347 await promise;
348 nextTick(cb, null);
349 } catch (err) {
350 nextTick(cb, err);
351 }
352 });
353 },
354 destroy,
355 }));
356 }
357 throw new ERR_INVALID_RETURN_VALUE(
358 'Iterable, AsyncIterable or AsyncFunction',
359 name,
360 value
361 );
362 }
363 if (isBlob(body)) {
364 return duplexify(body.arrayBuffer(), name);
365 }
366 if (isIterable(body)) {
367 return Readable.from(Duplexify, body, {
368 // TODO (ronag): highWaterMark?
369 objectMode: true,
370 writable: false,
371 });
372 }
373 
374 if (
375 body?.readable instanceof ReadableStream &&
376 body?.writable instanceof WritableStream
377 ) {
378 return Duplexify.fromWeb(body);
379 }
380 
381 if (
382 typeof (body === null || body === undefined ? undefined : body.writable) ===
383 'object' ||
384 typeof (body === null || body === undefined ? undefined : body.readable) ===
385 'object'
386 ) {
387 const readable =
388 body !== null && body !== undefined && body.readable
389 ? isReadableNodeStream(
390 body === null || body === undefined ? undefined : body.readable
391 )
392 ? body === null || body === undefined
393 ? undefined
394 : body.readable
395 : duplexify(body.readable, name)
396 : undefined;
397 const writable =
398 body !== null && body !== undefined && body.writable
399 ? isWritableNodeStream(
400 body === null || body === undefined ? undefined : body.writable
401 )
402 ? body === null || body === undefined
403 ? undefined
404 : body.writable
405 : duplexify(body.writable, name)
406 : undefined;
407 return _duplexify({
408 readable,
409 writable,
410 });
411 }
412 const then = body?.then;
413 if (typeof then === 'function') {
414 let d;
415 Reflect.apply(then, body, [
416 (val) => {
417 if (val != null) {
418 d.push(val);
419 }
420 d.push(null);
421 },
422 (err) => {
423 destroyer(d, err);
424 },
425 ]);
426 
427 return (d = new Duplexify({
428 objectMode: true,
429 writable: false,
430 read() {},
431 }));
432 }
433 throw new ERR_INVALID_ARG_TYPE(
434 name,
435 [
436 'Blob',
437 'ReadableStream',
438 'WritableStream',
439 'Stream',
440 'Iterable',
441 'AsyncIterable',
442 'Function',
443 '{ readable, writable } pair',
444 'Promise',
445 ],
446 body
447 );
448}
449 
450function fromAsyncGen(fn) {
451 let { promise, resolve } = Promise.withResolvers();
452 const ac = new AbortController();
453 const signal = ac.signal;
454 const value = fn(
455 (async function* () {
456 while (true) {
457 const _promise = promise;
458 promise = null;
459 const { chunk, done, cb } = await _promise;
460 nextTick(cb);
461 if (done) return;
462 if (signal.aborted)
463 throw new AbortError(undefined, {
464 cause: signal.reason,
465 });
466 ({ promise, resolve } = Promise.withResolvers());
467 yield chunk;
468 }
469 })(),
470 {
471 signal,
472 }
473 );
474 return {
475 value,
476 write(chunk, _encoding, cb) {
477 const _resolve = resolve;
478 resolve = null;
479 _resolve({
480 chunk,
481 done: false,
482 cb,
483 });
484 },
485 final(cb) {
486 const _resolve = resolve;
487 resolve = null;
488 _resolve({
489 done: true,
490 cb,
491 });
492 },
493 destroy(err, cb) {
494 ac.abort();
495 cb(err);
496 },
497 };
498}
499 
500function _duplexify(pair) {
501 const r =
502 pair.readable && typeof pair.readable.read !== 'function'
503 ? Readable.wrap(pair.readable)
504 : pair.readable;
505 const w = pair.writable;
506 let readable = !!isReadable(r);
507 let writable = !!isWritable(w);
508 let ondrain;
509 let onfinish;
510 let onreadable;
511 let onclose;
512 let d;
513 function onfinished(err) {
514 const cb = onclose;
515 onclose = null;
516 if (cb) {
517 cb(err);
518 } else if (err) {
519 d.destroy(err);
520 } else if (!readable && !writable) {
521 d.destroy();
522 }
523 }
524 
525 // TODO(ronag): Avoid double buffering.
526 // Implement Writable/Readable/Duplex traits.
527 // See, https://github.com/nodejs/node/pull/33515.
528 d = new Duplexify({
529 // TODO (ronag): highWaterMark?
530 readableObjectMode: !!(
531 r !== null &&
532 r !== undefined &&
533 r.readableObjectMode
534 ),
535 writableObjectMode: !!(
536 w !== null &&
537 w !== undefined &&
538 w.writableObjectMode
539 ),
540 readable,
541 writable,
542 });
543 if (writable) {
544 eos(w, (err) => {
545 writable = false;
546 if (err) {
547 destroyer(r, err);
548 }
549 onfinished(err);
550 });
551 d._write = function (chunk, encoding, callback) {
552 if (w.write(chunk, encoding)) {
553 callback();
554 } else {
555 ondrain = callback;
556 }
557 };
558 d._final = function (callback) {
559 w.end();
560 onfinish = callback;
561 };
562 w.on('drain', function () {
563 if (ondrain) {
564 const cb = ondrain;
565 ondrain = null;
566 cb();
567 }
568 });
569 w.on('finish', function () {
570 if (onfinish) {
571 const cb = onfinish;
572 onfinish = null;
573 cb();
574 }
575 });
576 }
577 if (readable) {
578 eos(r, (err) => {
579 readable = false;
580 if (err) {
581 destroyer(r, err);
582 }
583 onfinished(err);
584 });
585 r.on('readable', function () {
586 if (onreadable) {
587 const cb = onreadable;
588 onreadable = null;
589 cb();
590 }
591 });
592 r.on('end', function () {
593 d.push(null);
594 });
595 d._read = function () {
596 while (true) {
597 const buf = r.read();
598 if (buf === null) {
599 onreadable = d._read;
600 return;
601 }
602 if (!d.push(buf)) {
603 return;
604 }
605 }
606 };
607 }
608 d._destroy = function (err, callback) {
609 if (!err && onclose !== null) {
610 err = new AbortError();
611 }
612 onreadable = null;
613 ondrain = null;
614 onfinish = null;
615 if (onclose === null) {
616 callback(err);
617 } else {
618 onclose = callback;
619 destroyer(w, err);
620 destroyer(r, err);
621 }
622 };
623 return d;
624}
625 
626const kCallback = Symbol('Callback');
627const kInitOtherSide = Symbol('InitOtherSide');
628 
629class DuplexSide extends Duplex {
630 #otherSide = null;
631 
632 constructor(options) {
633 super(options);
634 this[kCallback] = null;
635 this.#otherSide = null;
636 }
637 
638 [kInitOtherSide](otherSide) {
639 // Ensure this can only be set once, to enforce encapsulation.
640 if (this.#otherSide === null) {
641 this.#otherSide = otherSide;
642 } else {
643 assert(this.#otherSide === null);
644 }
645 }
646 
647 _read() {
648 const callback = this[kCallback];
649 if (callback) {
650 this[kCallback] = null;
651 callback();
652 }
653 }
654 
655 _write(chunk, encoding, callback) {
656 assert(this.#otherSide !== null);
657 assert(this.#otherSide[kCallback] === null);
658 if (chunk.length === 0) {
659 nextTick(callback);
660 } else {
661 this.#otherSide.push(chunk);
662 this.#otherSide[kCallback] = callback;
663 }
664 }
665 
666 _final(callback) {
667 this.#otherSide.on('end', callback);
668 this.#otherSide.push(null);
669 }
670}
671 
672export function duplexPair(options) {
673 const side0 = new DuplexSide(options);
674 const side1 = new DuplexSide(options);
675 side0[kInitOtherSide](side1);
676 side1[kInitOtherSide](side0);
677 return [side0, side1];
678}
679 
680/**
681 * @param {Duplex} duplex
682 * @returns {ReadableWritablePair}
683 */
684export function newReadableWritablePairFromDuplex(
685 duplex,
686 createTypeBytes = false
687) {
688 // Not using the internal/streams/utils isWritableNodeStream and
689 // isReadableNodeStream utilities here because they will return false
690 // if the duplex was created with writable or readable options set to
691 // false. Instead, we'll check the readable and writable state after
692 // and return closed WritableStream or closed ReadableStream as
693 // necessary.
694 if (
695 typeof duplex?._writableState !== 'object' ||
696 typeof duplex?._readableState !== 'object'
697 ) {
698 throw new ERR_INVALID_ARG_TYPE('duplex', 'stream.Duplex', duplex);
699 }
700 
701 if (isDestroyed(duplex)) {
702 const writable = new WritableStream();
703 const readable = new ReadableStream();
704 writable.close();
705 readable.cancel();
706 return { readable, writable };
707 }
708 
709 const writable = isWritable(duplex)
710 ? newWritableStreamFromStreamWritable(duplex)
711 : new WritableStream();
712 
713 if (!isWritable(duplex)) writable.close();
714 
715 const readableOptions = createTypeBytes ? { type: 'bytes' } : {};
716 const readable = isReadable(duplex)
717 ? newReadableStreamFromStreamReadable(duplex, {}, createTypeBytes)
718 : new ReadableStream(readableOptions);
719 
720 if (!isReadable(duplex)) readable.cancel();
721 
722 return { writable, readable };
723}
724 
725/**
726 * @param {ReadableWritablePair} pair
727 * @param {{
728 * allowHalfOpen? : boolean,
729 * decodeStrings? : boolean,
730 * encoding? : string,
731 * highWaterMark? : number,
732 * objectMode? : boolean,
733 * signal? : AbortSignal,
734 * }} [options]
735 * @returns {Duplex}
736 */
737export function newStreamDuplexFromReadableWritablePair(
738 pair = {},
739 options = {}
740) {
741 validateObject(pair, 'pair');
742 const { readable: readableStream, writable: writableStream } = pair;
743 
744 if (!(readableStream instanceof ReadableStream)) {
745 throw new ERR_INVALID_ARG_TYPE(
746 'pair.readable',
747 'ReadableStream',
748 readableStream
749 );
750 }
751 if (!(writableStream instanceof WritableStream)) {
752 throw new ERR_INVALID_ARG_TYPE(
753 'pair.writable',
754 'WritableStream',
755 writableStream
756 );
757 }
758 
759 validateObject(options, 'options');
760 const {
761 allowHalfOpen = false,
762 objectMode = false,
763 encoding,
764 decodeStrings = true,
765 highWaterMark,
766 signal,
767 } = options;
768 
769 validateBoolean(objectMode, 'options.objectMode');
770 if (encoding !== undefined && !Buffer.isEncoding(encoding))
771 throw new ERR_INVALID_ARG_VALUE(encoding, 'options.encoding');
772 
773 const writer = writableStream.getWriter();
774 const reader = readableStream.getReader();
775 let writableClosed = false;
776 let readableClosed = false;
777 
778 const duplex = new Duplex({
779 allowHalfOpen,
780 highWaterMark,
781 objectMode,
782 encoding,
783 decodeStrings,
784 signal,
785 
786 writev(chunks, callback) {
787 function done(error) {
788 error = error.filter((e) => e);
789 try {
790 callback(error.length === 0 ? undefined : error);
791 } catch (error) {
792 // In a next tick because this is happening within
793 // a promise context, and if there are any errors
794 // thrown we don't want those to cause an unhandled
795 // rejection. Let's just escape the promise and
796 // handle it separately.
797 nextTick(() => destroy.call(duplex, error));
798 }
799 }
800 
801 writer.ready.then(() => {
802 return Promise.all(
803 chunks.map((data) => {
804 return writer.write(data.chunk);
805 })
806 ).then(done, done);
807 }, done);
808 },
809 
810 write(chunk, encoding, callback) {
811 if (typeof chunk === 'string' && decodeStrings && !objectMode) {
812 const enc = normalizeEncoding(encoding);
813 
814 if (enc === 'utf8') {
815 chunk = encoder.encode(chunk);
816 } else {
817 chunk = Buffer.from(chunk, encoding);
818 chunk = new Uint8Array(
819 chunk.buffer,
820 chunk.byteOffset,
821 chunk.byteLength
822 );
823 }
824 }
825 
826 function done(error) {
827 try {
828 callback(error);
829 } catch (error) {
830 destroy.call(duplex, error);
831 }
832 }
833 
834 writer.ready.then(() => {
835 return writer.write(chunk).then(done, done);
836 }, done);
837 },
838 
839 final(callback) {
840 function done(error) {
841 try {
842 callback(error);
843 } catch (error) {
844 // In a next tick because this is happening within
845 // a promise context, and if there are any errors
846 // thrown we don't want those to cause an unhandled
847 // rejection. Let's just escape the promise and
848 // handle it separately.
849 nextTick(() => destroy.call(duplex, error));
850 }
851 }
852 
853 if (!writableClosed) {
854 writer.close().then(done, done);
855 }
856 },
857 
858 read() {
859 reader.read().then(
860 (chunk) => {
861 if (chunk.done) {
862 duplex.push(null);
863 } else {
864 duplex.push(chunk.value);
865 }
866 },
867 (error) => destroy.call(duplex, error)
868 );
869 },
870 
871 destroy(error, callback) {
872 function done() {
873 try {
874 callback(error);
875 } catch (error) {
876 // In a next tick because this is happening within
877 // a promise context, and if there are any errors
878 // thrown we don't want those to cause an unhandled
879 // rejection. Let's just escape the promise and
880 // handle it separately.
881 nextTick(() => {
882 throw error;
883 });
884 }
885 }
886 
887 async function closeWriter() {
888 if (!writableClosed) await writer.abort(error);
889 }
890 
891 async function closeReader() {
892 if (!readableClosed) await reader.cancel(error);
893 }
894 
895 if (!writableClosed || !readableClosed) {
896 Promise.all([closeWriter(), closeReader()]).then(done, done);
897 return;
898 }
899 
900 done();
901 },
902 });
903 
904 writer.closed.then(
905 () => {
906 writableClosed = true;
907 if (!isWritableEnded(duplex))
908 destroy.call(duplex, new ERR_STREAM_PREMATURE_CLOSE());
909 },
910 (error) => {
911 writableClosed = true;
912 readableClosed = true;
913 destroy.call(duplex, error);
914 }
915 );
916 
917 reader.closed.then(
918 () => {
919 readableClosed = true;
920 },
921 (error) => {
922 writableClosed = true;
923 readableClosed = true;
924 destroy.call(duplex, error);
925 }
926 );
927 
928 return duplex;
929}