File
Blob: src/node/internal/streams_duplex.js
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | // |
| 5 | // Copyright Joyent, Inc. and other Node contributors. |
| 6 | // |
| 7 | // Permission is hereby granted, free of charge, to any person obtaining a |
| 8 | // copy of this software and associated documentation files (the |
| 9 | // "Software"), to deal in the Software without restriction, including |
| 10 | // without limitation the rights to use, copy, modify, merge, publish, |
| 11 | // distribute, sublicense, and/or sell copies of the Software, and to permit |
| 12 | // persons to whom the Software is furnished to do so, subject to the |
| 13 | // following conditions: |
| 14 | // |
| 15 | // The above copyright notice and this permission notice shall be included |
| 16 | // in all copies or substantial portions of the Software. |
| 17 | // |
| 18 | // THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS |
| 19 | // OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF |
| 20 | // MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN |
| 21 | // NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, |
| 22 | // DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR |
| 23 | // OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE |
| 24 | // USE OR OTHER DEALINGS IN THE SOFTWARE. |
| 25 | |
| 26 | /* TODO: the following is adopted code, enabling linting one day */ |
| 27 | /* eslint-disable */ |
| 28 | |
| 29 | import { Buffer } from 'node-internal:internal_buffer'; |
| 30 | import { |
| 31 | Readable, |
| 32 | newReadableStreamFromStreamReadable, |
| 33 | } from 'node-internal:streams_readable'; |
| 34 | import { |
| 35 | Writable, |
| 36 | newWritableStreamFromStreamWritable, |
| 37 | } from 'node-internal:streams_writable'; |
| 38 | import { ok as assert } from 'node-internal:internal_assert'; |
| 39 | import { Stream } from 'node-internal:streams_legacy'; |
| 40 | import { nextTick } from 'node-internal:internal_process'; |
| 41 | import { validateBoolean, validateObject } from 'node-internal:validators'; |
| 42 | import { normalizeEncoding } from 'node-internal:internal_utils'; |
| 43 | import { addAbortSignal } from 'node-internal:streams_add_abort_signal'; |
| 44 | |
| 45 | import { |
| 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'; |
| 57 | import { |
| 58 | construct as destroyConstruct, |
| 59 | destroyer, |
| 60 | } from 'node-internal:streams_destroy'; |
| 61 | import { eos } from 'node-internal:streams_end_of_stream'; |
| 62 | |
| 63 | import { |
| 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 | */ |
| 76 | const encoder = new TextEncoder(); |
| 77 | |
| 78 | Object.setPrototypeOf(Duplex.prototype, Readable.prototype); |
| 79 | Object.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`. |
| 90 | Duplex.prototype.destroy = Writable.prototype.destroy; |
| 91 | |
| 92 | export 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`. |
| 163 | Duplex.prototype.destroy = Writable.prototype.destroy; |
| 164 | |
| 165 | Object.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 | |
| 231 | export function fromWeb(pair, options) { |
| 232 | return newStreamDuplexFromReadableWritablePair(pair, options); |
| 233 | } |
| 234 | |
| 235 | export function toWeb(duplex) { |
| 236 | return newReadableWritablePairFromDuplex(duplex); |
| 237 | } |
| 238 | |
| 239 | export function toBYOBWeb(duplex) { |
| 240 | return newReadableWritablePairFromDuplex(duplex, true /* createTypeBytes */); |
| 241 | } |
| 242 | |
| 243 | export function from(body) { |
| 244 | return duplexify(body, 'body'); |
| 245 | } |
| 246 | |
| 247 | Duplex.fromWeb = fromWeb; |
| 248 | Duplex.toWeb = toWeb; |
| 249 | Duplex.from = from; |
| 250 | |
| 251 | // ====================================================================================== |
| 252 | |
| 253 | function isBlob(b) { |
| 254 | return b instanceof Blob; |
| 255 | } |
| 256 | |
| 257 | // This is needed for pre node 17. |
| 258 | class 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 | |
| 285 | function 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 | |
| 450 | function 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 | |
| 500 | function _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 | |
| 626 | const kCallback = Symbol('Callback'); |
| 627 | const kInitOtherSide = Symbol('InitOtherSide'); |
| 628 | |
| 629 | class 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 | |
| 672 | export 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 | */ |
| 684 | export 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 | */ |
| 737 | export 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 | } |