basant307/AI_Governance_Project
048
1const { EventEmitter } = require('events-universal')2const STREAM_DESTROYED = new Error('Stream was destroyed')3const PREMATURE_CLOSE = new Error('Premature close')4 5const FIFO = require('fast-fifo')6const TextDecoder = require('text-decoder')7 8// if we do a future major, expect queue microtask to be there always, for now a bit defensive9const qmt = typeof queueMicrotask === 'undefined' ? fn => global.process.nextTick(fn) : queueMicrotask10 11/* eslint-disable no-multi-spaces */12 13// 29 bits used total (4 from shared, 14 from read, and 11 from write)14const MAX = ((1 << 29) - 1)15 16// Shared state17const OPENING = 0b000118const PREDESTROYING = 0b001019const DESTROYING = 0b010020const DESTROYED = 0b100021 22const NOT_OPENING = MAX ^ OPENING23const NOT_PREDESTROYING = MAX ^ PREDESTROYING24 25// Read state (4 bit offset from shared state)26const READ_ACTIVE = 0b00000000000001 << 427const READ_UPDATING = 0b00000000000010 << 428const READ_PRIMARY = 0b00000000000100 << 429const READ_QUEUED = 0b00000000001000 << 430const READ_RESUMED = 0b00000000010000 << 431const READ_PIPE_DRAINED = 0b00000000100000 << 432const READ_ENDING = 0b00000001000000 << 433const READ_EMIT_DATA = 0b00000010000000 << 434const READ_EMIT_READABLE = 0b00000100000000 << 435const READ_EMITTED_READABLE = 0b00001000000000 << 436const READ_DONE = 0b00010000000000 << 437const READ_NEXT_TICK = 0b00100000000000 << 438const READ_NEEDS_PUSH = 0b01000000000000 << 439const READ_READ_AHEAD = 0b10000000000000 << 440 41// Combined read state42const READ_FLOWING = READ_RESUMED | READ_PIPE_DRAINED43const READ_ACTIVE_AND_NEEDS_PUSH = READ_ACTIVE | READ_NEEDS_PUSH44const READ_PRIMARY_AND_ACTIVE = READ_PRIMARY | READ_ACTIVE45const READ_EMIT_READABLE_AND_QUEUED = READ_EMIT_READABLE | READ_QUEUED46const READ_RESUMED_READ_AHEAD = READ_RESUMED | READ_READ_AHEAD47 48const READ_NOT_ACTIVE = MAX ^ READ_ACTIVE49const READ_NON_PRIMARY = MAX ^ READ_PRIMARY50const READ_NON_PRIMARY_AND_PUSHED = MAX ^ (READ_PRIMARY | READ_NEEDS_PUSH)51const READ_PUSHED = MAX ^ READ_NEEDS_PUSH52const READ_PAUSED = MAX ^ READ_RESUMED53const READ_NOT_QUEUED = MAX ^ (READ_QUEUED | READ_EMITTED_READABLE)54const READ_NOT_ENDING = MAX ^ READ_ENDING55const READ_PIPE_NOT_DRAINED = MAX ^ READ_FLOWING56const READ_NOT_NEXT_TICK = MAX ^ READ_NEXT_TICK57const READ_NOT_UPDATING = MAX ^ READ_UPDATING58const READ_NO_READ_AHEAD = MAX ^ READ_READ_AHEAD59const READ_PAUSED_NO_READ_AHEAD = MAX ^ READ_RESUMED_READ_AHEAD60 61// Write state (18 bit offset, 4 bit offset from shared state and 14 from read state)62const WRITE_ACTIVE = 0b00000000001 << 1863const WRITE_UPDATING = 0b00000000010 << 1864const WRITE_PRIMARY = 0b00000000100 << 1865const WRITE_QUEUED = 0b00000001000 << 1866const WRITE_UNDRAINED = 0b00000010000 << 1867const WRITE_DONE = 0b00000100000 << 1868const WRITE_EMIT_DRAIN = 0b00001000000 << 1869const WRITE_NEXT_TICK = 0b00010000000 << 1870const WRITE_WRITING = 0b00100000000 << 1871const WRITE_FINISHING = 0b01000000000 << 1872const WRITE_CORKED = 0b10000000000 << 1873 74const WRITE_NOT_ACTIVE = MAX ^ (WRITE_ACTIVE | WRITE_WRITING)75const WRITE_NON_PRIMARY = MAX ^ WRITE_PRIMARY76const WRITE_NOT_FINISHING = MAX ^ (WRITE_ACTIVE | WRITE_FINISHING)77const WRITE_DRAINED = MAX ^ WRITE_UNDRAINED78const WRITE_NOT_QUEUED = MAX ^ WRITE_QUEUED79const WRITE_NOT_NEXT_TICK = MAX ^ WRITE_NEXT_TICK80const WRITE_NOT_UPDATING = MAX ^ WRITE_UPDATING81const WRITE_NOT_CORKED = MAX ^ WRITE_CORKED82 83// Combined shared state84const ACTIVE = READ_ACTIVE | WRITE_ACTIVE85const NOT_ACTIVE = MAX ^ ACTIVE86const DONE = READ_DONE | WRITE_DONE87const DESTROY_STATUS = DESTROYING | DESTROYED | PREDESTROYING88const OPEN_STATUS = DESTROY_STATUS | OPENING89const AUTO_DESTROY = DESTROY_STATUS | DONE90const NON_PRIMARY = WRITE_NON_PRIMARY & READ_NON_PRIMARY91const ACTIVE_OR_TICKING = WRITE_NEXT_TICK | READ_NEXT_TICK92const TICKING = ACTIVE_OR_TICKING & NOT_ACTIVE93const IS_OPENING = OPEN_STATUS | TICKING94 95// Combined shared state and read state96const READ_PRIMARY_STATUS = OPEN_STATUS | READ_ENDING | READ_DONE97const READ_STATUS = OPEN_STATUS | READ_DONE | READ_QUEUED98const READ_ENDING_STATUS = OPEN_STATUS | READ_ENDING | READ_QUEUED99const READ_READABLE_STATUS = OPEN_STATUS | READ_EMIT_READABLE | READ_QUEUED | READ_EMITTED_READABLE100const SHOULD_NOT_READ = OPEN_STATUS | READ_ACTIVE | READ_ENDING | READ_DONE | READ_NEEDS_PUSH | READ_READ_AHEAD101const READ_BACKPRESSURE_STATUS = DESTROY_STATUS | READ_ENDING | READ_DONE102const READ_UPDATE_SYNC_STATUS = READ_UPDATING | OPEN_STATUS | READ_NEXT_TICK | READ_PRIMARY103const READ_NEXT_TICK_OR_OPENING = READ_NEXT_TICK | OPENING104 105// Combined write state106const WRITE_PRIMARY_STATUS = OPEN_STATUS | WRITE_FINISHING | WRITE_DONE107const WRITE_QUEUED_AND_UNDRAINED = WRITE_QUEUED | WRITE_UNDRAINED108const WRITE_QUEUED_AND_ACTIVE = WRITE_QUEUED | WRITE_ACTIVE109const WRITE_DRAIN_STATUS = WRITE_QUEUED | WRITE_UNDRAINED | OPEN_STATUS | WRITE_ACTIVE110const WRITE_STATUS = OPEN_STATUS | WRITE_ACTIVE | WRITE_QUEUED | WRITE_CORKED111const WRITE_PRIMARY_AND_ACTIVE = WRITE_PRIMARY | WRITE_ACTIVE112const WRITE_ACTIVE_AND_WRITING = WRITE_ACTIVE | WRITE_WRITING113const WRITE_FINISHING_STATUS = OPEN_STATUS | WRITE_FINISHING | WRITE_QUEUED_AND_ACTIVE | WRITE_DONE114const WRITE_BACKPRESSURE_STATUS = WRITE_UNDRAINED | DESTROY_STATUS | WRITE_FINISHING | WRITE_DONE115const WRITE_UPDATE_SYNC_STATUS = WRITE_UPDATING | OPEN_STATUS | WRITE_NEXT_TICK | WRITE_PRIMARY116const WRITE_DROP_DATA = WRITE_FINISHING | WRITE_DONE | DESTROY_STATUS117 118const asyncIterator = Symbol.asyncIterator || Symbol('asyncIterator')119 120class WritableState {121 constructor (stream, { highWaterMark = 16384, map = null, mapWritable, byteLength, byteLengthWritable } = {}) {122 this.stream = stream123 this.queue = new FIFO()124 this.highWaterMark = highWaterMark125 this.buffered = 0126 this.error = null127 this.pipeline = null128 this.drains = null // if we add more seldomly used helpers we might them into a subobject so its a single ptr129 this.byteLength = byteLengthWritable || byteLength || defaultByteLength130 this.map = mapWritable || map131 this.afterWrite = afterWrite.bind(this)132 this.afterUpdateNextTick = updateWriteNT.bind(this)133 }134 135 get ended () {136 return (this.stream._duplexState & WRITE_DONE) !== 0137 }138 139 push (data) {140 if ((this.stream._duplexState & WRITE_DROP_DATA) !== 0) return false141 if (this.map !== null) data = this.map(data)142 143 this.buffered += this.byteLength(data)144 this.queue.push(data)145 146 if (this.buffered < this.highWaterMark) {147 this.stream._duplexState |= WRITE_QUEUED148 return true149 }150 151 this.stream._duplexState |= WRITE_QUEUED_AND_UNDRAINED152 return false153 }154 155 shift () {156 const data = this.queue.shift()157 158 this.buffered -= this.byteLength(data)159 if (this.buffered === 0) this.stream._duplexState &= WRITE_NOT_QUEUED160 161 return data162 }163 164 end (data) {165 if (typeof data === 'function') this.stream.once('finish', data)166 else if (data !== undefined && data !== null) this.push(data)167 this.stream._duplexState = (this.stream._duplexState | WRITE_FINISHING) & WRITE_NON_PRIMARY168 }169 170 autoBatch (data, cb) {171 const buffer = []172 const stream = this.stream173 174 buffer.push(data)175 while ((stream._duplexState & WRITE_STATUS) === WRITE_QUEUED_AND_ACTIVE) {176 buffer.push(stream._writableState.shift())177 }178 179 if ((stream._duplexState & OPEN_STATUS) !== 0) return cb(null)180 stream._writev(buffer, cb)181 }182 183 update () {184 const stream = this.stream185 186 stream._duplexState |= WRITE_UPDATING187 188 do {189 while ((stream._duplexState & WRITE_STATUS) === WRITE_QUEUED) {190 const data = this.shift()191 stream._duplexState |= WRITE_ACTIVE_AND_WRITING192 stream._write(data, this.afterWrite)193 }194 195 if ((stream._duplexState & WRITE_PRIMARY_AND_ACTIVE) === 0) this.updateNonPrimary()196 } while (this.continueUpdate() === true)197 198 stream._duplexState &= WRITE_NOT_UPDATING199 }200 201 updateNonPrimary () {202 const stream = this.stream203 204 if ((stream._duplexState & WRITE_FINISHING_STATUS) === WRITE_FINISHING) {205 stream._duplexState = stream._duplexState | WRITE_ACTIVE206 stream._final(afterFinal.bind(this))207 return208 }209 210 if ((stream._duplexState & DESTROY_STATUS) === DESTROYING) {211 if ((stream._duplexState & ACTIVE_OR_TICKING) === 0) {212 stream._duplexState |= ACTIVE213 stream._destroy(afterDestroy.bind(this))214 }215 return216 }217 218 if ((stream._duplexState & IS_OPENING) === OPENING) {219 stream._duplexState = (stream._duplexState | ACTIVE) & NOT_OPENING220 stream._open(afterOpen.bind(this))221 }222 }223 224 continueUpdate () {225 if ((this.stream._duplexState & WRITE_NEXT_TICK) === 0) return false226 this.stream._duplexState &= WRITE_NOT_NEXT_TICK227 return true228 }229 230 updateCallback () {231 if ((this.stream._duplexState & WRITE_UPDATE_SYNC_STATUS) === WRITE_PRIMARY) this.update()232 else this.updateNextTick()233 }234 235 updateNextTick () {236 if ((this.stream._duplexState & WRITE_NEXT_TICK) !== 0) return237 this.stream._duplexState |= WRITE_NEXT_TICK238 if ((this.stream._duplexState & WRITE_UPDATING) === 0) qmt(this.afterUpdateNextTick)239 }240}241 242class ReadableState {243 constructor (stream, { highWaterMark = 16384, map = null, mapReadable, byteLength, byteLengthReadable } = {}) {244 this.stream = stream245 this.queue = new FIFO()246 this.highWaterMark = highWaterMark === 0 ? 1 : highWaterMark247 this.buffered = 0248 this.readAhead = highWaterMark > 0249 this.error = null250 this.pipeline = null251 this.byteLength = byteLengthReadable || byteLength || defaultByteLength252 this.map = mapReadable || map253 this.pipeTo = null254 this.afterRead = afterRead.bind(this)255 this.afterUpdateNextTick = updateReadNT.bind(this)256 }257 258 get ended () {259 return (this.stream._duplexState & READ_DONE) !== 0260 }261 262 pipe (pipeTo, cb) {263 if (this.pipeTo !== null) throw new Error('Can only pipe to one destination')264 if (typeof cb !== 'function') cb = null265 266 this.stream._duplexState |= READ_PIPE_DRAINED267 this.pipeTo = pipeTo268 this.pipeline = new Pipeline(this.stream, pipeTo, cb)269 270 if (cb) this.stream.on('error', noop) // We already error handle this so supress crashes271 272 if (isStreamx(pipeTo)) {273 pipeTo._writableState.pipeline = this.pipeline274 if (cb) pipeTo.on('error', noop) // We already error handle this so supress crashes275 pipeTo.on('finish', this.pipeline.finished.bind(this.pipeline)) // TODO: just call finished from pipeTo itself276 } else {277 const onerror = this.pipeline.done.bind(this.pipeline, pipeTo)278 const onclose = this.pipeline.done.bind(this.pipeline, pipeTo, null) // onclose has a weird bool arg279 pipeTo.on('error', onerror)280 pipeTo.on('close', onclose)281 pipeTo.on('finish', this.pipeline.finished.bind(this.pipeline))282 }283 284 pipeTo.on('drain', afterDrain.bind(this))285 this.stream.emit('piping', pipeTo)286 pipeTo.emit('pipe', this.stream)287 }288 289 push (data) {290 const stream = this.stream291 292 if (data === null) {293 this.highWaterMark = 0294 stream._duplexState = (stream._duplexState | READ_ENDING) & READ_NON_PRIMARY_AND_PUSHED295 return false296 }297 298 if (this.map !== null) {299 data = this.map(data)300 if (data === null) {301 stream._duplexState &= READ_PUSHED302 return this.buffered < this.highWaterMark303 }304 }305 306 this.buffered += this.byteLength(data)307 this.queue.push(data)308 309 stream._duplexState = (stream._duplexState | READ_QUEUED) & READ_PUSHED310 311 return this.buffered < this.highWaterMark312 }313 314 shift () {315 const data = this.queue.shift()316 317 this.buffered -= this.byteLength(data)318 if (this.buffered === 0) this.stream._duplexState &= READ_NOT_QUEUED319 return data320 }321 322 unshift (data) {323 const pending = [this.map !== null ? this.map(data) : data]324 while (this.buffered > 0) pending.push(this.shift())325 326 for (let i = 0; i < pending.length - 1; i++) {327 const data = pending[i]328 this.buffered += this.byteLength(data)329 this.queue.push(data)330 }331 332 this.push(pending[pending.length - 1])333 }334 335 read () {336 const stream = this.stream337 338 if ((stream._duplexState & READ_STATUS) === READ_QUEUED) {339 const data = this.shift()340 if (this.pipeTo !== null && this.pipeTo.write(data) === false) stream._duplexState &= READ_PIPE_NOT_DRAINED341 if ((stream._duplexState & READ_EMIT_DATA) !== 0) stream.emit('data', data)342 return data343 }344 345 if (this.readAhead === false) {346 stream._duplexState |= READ_READ_AHEAD347 this.updateNextTick()348 }349 350 return null351 }352 353 drain () {354 const stream = this.stream355 356 while ((stream._duplexState & READ_STATUS) === READ_QUEUED && (stream._duplexState & READ_FLOWING) !== 0) {357 const data = this.shift()358 if (this.pipeTo !== null && this.pipeTo.write(data) === false) stream._duplexState &= READ_PIPE_NOT_DRAINED359 if ((stream._duplexState & READ_EMIT_DATA) !== 0) stream.emit('data', data)360 }361 }362 363 update () {364 const stream = this.stream365 366 stream._duplexState |= READ_UPDATING367 368 do {369 this.drain()370 371 while (this.buffered < this.highWaterMark && (stream._duplexState & SHOULD_NOT_READ) === READ_READ_AHEAD) {372 stream._duplexState |= READ_ACTIVE_AND_NEEDS_PUSH373 stream._read(this.afterRead)374 this.drain()375 }376 377 if ((stream._duplexState & READ_READABLE_STATUS) === READ_EMIT_READABLE_AND_QUEUED) {378 stream._duplexState |= READ_EMITTED_READABLE379 stream.emit('readable')380 }381 382 if ((stream._duplexState & READ_PRIMARY_AND_ACTIVE) === 0) this.updateNonPrimary()383 } while (this.continueUpdate() === true)384 385 stream._duplexState &= READ_NOT_UPDATING386 }387 388 updateNonPrimary () {389 const stream = this.stream390 391 if ((stream._duplexState & READ_ENDING_STATUS) === READ_ENDING) {392 stream._duplexState = (stream._duplexState | READ_DONE) & READ_NOT_ENDING393 stream.emit('end')394 if ((stream._duplexState & AUTO_DESTROY) === DONE) stream._duplexState |= DESTROYING395 if (this.pipeTo !== null) this.pipeTo.end()396 }397 398 if ((stream._duplexState & DESTROY_STATUS) === DESTROYING) {399 if ((stream._duplexState & ACTIVE_OR_TICKING) === 0) {400 stream._duplexState |= ACTIVE401 stream._destroy(afterDestroy.bind(this))402 }403 return404 }405 406 if ((stream._duplexState & IS_OPENING) === OPENING) {407 stream._duplexState = (stream._duplexState | ACTIVE) & NOT_OPENING408 stream._open(afterOpen.bind(this))409 }410 }411 412 continueUpdate () {413 if ((this.stream._duplexState & READ_NEXT_TICK) === 0) return false414 this.stream._duplexState &= READ_NOT_NEXT_TICK415 return true416 }417 418 updateCallback () {419 if ((this.stream._duplexState & READ_UPDATE_SYNC_STATUS) === READ_PRIMARY) this.update()420 else this.updateNextTick()421 }422 423 updateNextTickIfOpen () {424 if ((this.stream._duplexState & READ_NEXT_TICK_OR_OPENING) !== 0) return425 this.stream._duplexState |= READ_NEXT_TICK426 if ((this.stream._duplexState & READ_UPDATING) === 0) qmt(this.afterUpdateNextTick)427 }428 429 updateNextTick () {430 if ((this.stream._duplexState & READ_NEXT_TICK) !== 0) return431 this.stream._duplexState |= READ_NEXT_TICK432 if ((this.stream._duplexState & READ_UPDATING) === 0) qmt(this.afterUpdateNextTick)433 }434}435 436class TransformState {437 constructor (stream) {438 this.data = null439 this.afterTransform = afterTransform.bind(stream)440 this.afterFinal = null441 }442}443 444class Pipeline {445 constructor (src, dst, cb) {446 this.from = src447 this.to = dst448 this.afterPipe = cb449 this.error = null450 this.pipeToFinished = false451 }452 453 finished () {454 this.pipeToFinished = true455 }456 457 done (stream, err) {458 if (err) this.error = err459 460 if (stream === this.to) {461 this.to = null462 463 if (this.from !== null) {464 if ((this.from._duplexState & READ_DONE) === 0 || !this.pipeToFinished) {465 this.from.destroy(this.error || new Error('Writable stream closed prematurely'))466 }467 return468 }469 }470 471 if (stream === this.from) {472 this.from = null473 474 if (this.to !== null) {475 if ((stream._duplexState & READ_DONE) === 0) {476 this.to.destroy(this.error || new Error('Readable stream closed before ending'))477 }478 return479 }480 }481 482 if (this.afterPipe !== null) this.afterPipe(this.error)483 this.to = this.from = this.afterPipe = null484 }485}486 487function afterDrain () {488 this.stream._duplexState |= READ_PIPE_DRAINED489 this.updateCallback()490}491 492function afterFinal (err) {493 const stream = this.stream494 if (err) stream.destroy(err)495 if ((stream._duplexState & DESTROY_STATUS) === 0) {496 stream._duplexState |= WRITE_DONE497 stream.emit('finish')498 }499 if ((stream._duplexState & AUTO_DESTROY) === DONE) {500 stream._duplexState |= DESTROYING501 }502 503 stream._duplexState &= WRITE_NOT_FINISHING504 505 // no need to wait the extra tick here, so we short circuit that506 if ((stream._duplexState & WRITE_UPDATING) === 0) this.update()507 else this.updateNextTick()508}509 510function afterDestroy (err) {511 const stream = this.stream512 513 if (!err && this.error !== STREAM_DESTROYED) err = this.error514 if (err) stream.emit('error', err)515 stream._duplexState |= DESTROYED516 stream.emit('close')517 518 const rs = stream._readableState519 const ws = stream._writableState520 521 if (rs !== null && rs.pipeline !== null) rs.pipeline.done(stream, err)522 523 if (ws !== null) {524 while (ws.drains !== null && ws.drains.length > 0) ws.drains.shift().resolve(false)525 if (ws.pipeline !== null) ws.pipeline.done(stream, err)526 }527}528 529function afterWrite (err) {530 const stream = this.stream531 532 if (err) stream.destroy(err)533 stream._duplexState &= WRITE_NOT_ACTIVE534 535 if (this.drains !== null) tickDrains(this.drains)536 537 if ((stream._duplexState & WRITE_DRAIN_STATUS) === WRITE_UNDRAINED) {538 stream._duplexState &= WRITE_DRAINED539 if ((stream._duplexState & WRITE_EMIT_DRAIN) === WRITE_EMIT_DRAIN) {540 stream.emit('drain')541 }542 }543 544 this.updateCallback()545}546 547function afterRead (err) {548 if (err) this.stream.destroy(err)549 this.stream._duplexState &= READ_NOT_ACTIVE550 if (this.readAhead === false && (this.stream._duplexState & READ_RESUMED) === 0) this.stream._duplexState &= READ_NO_READ_AHEAD551 this.updateCallback()552}553 554function updateReadNT () {555 if ((this.stream._duplexState & READ_UPDATING) === 0) {556 this.stream._duplexState &= READ_NOT_NEXT_TICK557 this.update()558 }559}560 561function updateWriteNT () {562 if ((this.stream._duplexState & WRITE_UPDATING) === 0) {563 this.stream._duplexState &= WRITE_NOT_NEXT_TICK564 this.update()565 }566}567 568function tickDrains (drains) {569 for (let i = 0; i < drains.length; i++) {570 // drains.writes are monotonic, so if one is 0 its always the first one571 if (--drains[i].writes === 0) {572 drains.shift().resolve(true)573 i--574 }575 }576}577 578function afterOpen (err) {579 const stream = this.stream580 581 if (err) stream.destroy(err)582 583 if ((stream._duplexState & DESTROYING) === 0) {584 if ((stream._duplexState & READ_PRIMARY_STATUS) === 0) stream._duplexState |= READ_PRIMARY585 if ((stream._duplexState & WRITE_PRIMARY_STATUS) === 0) stream._duplexState |= WRITE_PRIMARY586 stream.emit('open')587 }588 589 stream._duplexState &= NOT_ACTIVE590 591 if (stream._writableState !== null) {592 stream._writableState.updateCallback()593 }594 595 if (stream._readableState !== null) {596 stream._readableState.updateCallback()597 }598}599 600function afterTransform (err, data) {601 if (data !== undefined && data !== null) this.push(data)602 this._writableState.afterWrite(err)603}604 605function newListener (name) {606 if (this._readableState !== null) {607 if (name === 'data') {608 this._duplexState |= (READ_EMIT_DATA | READ_RESUMED_READ_AHEAD)609 this._readableState.updateNextTick()610 }611 if (name === 'readable') {612 this._duplexState |= READ_EMIT_READABLE613 this._readableState.updateNextTick()614 }615 }616 617 if (this._writableState !== null) {618 if (name === 'drain') {619 this._duplexState |= WRITE_EMIT_DRAIN620 this._writableState.updateNextTick()621 }622 }623}624 625class Stream extends EventEmitter {626 constructor (opts) {627 super()628 629 this._duplexState = 0630 this._readableState = null631 this._writableState = null632 633 if (opts) {634 if (opts.open) this._open = opts.open635 if (opts.destroy) this._destroy = opts.destroy636 if (opts.predestroy) this._predestroy = opts.predestroy637 if (opts.signal) {638 opts.signal.addEventListener('abort', abort.bind(this))639 }640 }641 642 this.on('newListener', newListener)643 }644 645 _open (cb) {646 cb(null)647 }648 649 _destroy (cb) {650 cb(null)651 }652 653 _predestroy () {654 // does nothing655 }656 657 get readable () {658 return this._readableState !== null ? true : undefined659 }660 661 get writable () {662 return this._writableState !== null ? true : undefined663 }664 665 get destroyed () {666 return (this._duplexState & DESTROYED) !== 0667 }668 669 get destroying () {670 return (this._duplexState & DESTROY_STATUS) !== 0671 }672 673 destroy (err) {674 if ((this._duplexState & DESTROY_STATUS) === 0) {675 if (!err) err = STREAM_DESTROYED676 this._duplexState = (this._duplexState | DESTROYING) & NON_PRIMARY677 678 if (this._readableState !== null) {679 this._readableState.highWaterMark = 0680 this._readableState.error = err681 }682 if (this._writableState !== null) {683 this._writableState.highWaterMark = 0684 this._writableState.error = err685 }686 687 this._duplexState |= PREDESTROYING688 this._predestroy()689 this._duplexState &= NOT_PREDESTROYING690 691 if (this._readableState !== null) this._readableState.updateNextTick()692 if (this._writableState !== null) this._writableState.updateNextTick()693 }694 }695}696 697class Readable extends Stream {698 constructor (opts) {699 super(opts)700 701 this._duplexState |= OPENING | WRITE_DONE | READ_READ_AHEAD702 this._readableState = new ReadableState(this, opts)703 704 if (opts) {705 if (this._readableState.readAhead === false) this._duplexState &= READ_NO_READ_AHEAD706 if (opts.read) this._read = opts.read707 if (opts.eagerOpen) this._readableState.updateNextTick()708 if (opts.encoding) this.setEncoding(opts.encoding)709 }710 }711 712 setEncoding (encoding) {713 const dec = new TextDecoder(encoding)714 const map = this._readableState.map || echo715 this._readableState.map = mapOrSkip716 return this717 718 function mapOrSkip (data) {719 const next = dec.push(data)720 return next === '' && (data.byteLength !== 0 || dec.remaining > 0) ? null : map(next)721 }722 }723 724 _read (cb) {725 cb(null)726 }727 728 pipe (dest, cb) {729 this._readableState.updateNextTick()730 this._readableState.pipe(dest, cb)731 return dest732 }733 734 read () {735 this._readableState.updateNextTick()736 return this._readableState.read()737 }738 739 push (data) {740 this._readableState.updateNextTickIfOpen()741 return this._readableState.push(data)742 }743 744 unshift (data) {745 this._readableState.updateNextTickIfOpen()746 return this._readableState.unshift(data)747 }748 749 resume () {750 this._duplexState |= READ_RESUMED_READ_AHEAD751 this._readableState.updateNextTick()752 return this753 }754 755 pause () {756 this._duplexState &= (this._readableState.readAhead === false ? READ_PAUSED_NO_READ_AHEAD : READ_PAUSED)757 return this758 }759 760 static _fromAsyncIterator (ite, opts) {761 let destroy762 763 const rs = new Readable({764 ...opts,765 read (cb) {766 ite.next().then(push).then(cb.bind(null, null)).catch(cb)767 },768 predestroy () {769 destroy = ite.return()770 },771 destroy (cb) {772 if (!destroy) return cb(null)773 destroy.then(cb.bind(null, null)).catch(cb)774 }775 })776 777 return rs778 779 function push (data) {780 if (data.done) rs.push(null)781 else rs.push(data.value)782 }783 }784 785 static from (data, opts) {786 if (isReadStreamx(data)) return data787 if (data[asyncIterator]) return this._fromAsyncIterator(data[asyncIterator](), opts)788 if (!Array.isArray(data)) data = data === undefined ? [] : [data]789 790 let i = 0791 return new Readable({792 ...opts,793 read (cb) {794 this.push(i === data.length ? null : data[i++])795 cb(null)796 }797 })798 }799 800 static isBackpressured (rs) {801 return (rs._duplexState & READ_BACKPRESSURE_STATUS) !== 0 || rs._readableState.buffered >= rs._readableState.highWaterMark802 }803 804 static isPaused (rs) {805 return (rs._duplexState & READ_RESUMED) === 0806 }807 808 [asyncIterator] () {809 const stream = this810 811 let error = null812 let promiseResolve = null813 let promiseReject = null814 815 this.on('error', (err) => { error = err })816 this.on('readable', onreadable)817 this.on('close', onclose)818 819 return {820 [asyncIterator] () {821 return this822 },823 next () {824 return new Promise(function (resolve, reject) {825 promiseResolve = resolve826 promiseReject = reject827 const data = stream.read()828 if (data !== null) ondata(data)829 else if ((stream._duplexState & DESTROYED) !== 0) ondata(null)830 })831 },832 return () {833 return destroy(null)834 },835 throw (err) {836 return destroy(err)837 }838 }839 840 function onreadable () {841 if (promiseResolve !== null) ondata(stream.read())842 }843 844 function onclose () {845 if (promiseResolve !== null) ondata(null)846 }847 848 function ondata (data) {849 if (promiseReject === null) return850 if (error) promiseReject(error)851 else if (data === null && (stream._duplexState & READ_DONE) === 0) promiseReject(STREAM_DESTROYED)852 else promiseResolve({ value: data, done: data === null })853 promiseReject = promiseResolve = null854 }855 856 function destroy (err) {857 stream.destroy(err)858 return new Promise((resolve, reject) => {859 if (stream._duplexState & DESTROYED) return resolve({ value: undefined, done: true })860 stream.once('close', function () {861 if (err) reject(err)862 else resolve({ value: undefined, done: true })863 })864 })865 }866 }867}868 869class Writable extends Stream {870 constructor (opts) {871 super(opts)872 873 this._duplexState |= OPENING | READ_DONE874 this._writableState = new WritableState(this, opts)875 876 if (opts) {877 if (opts.writev) this._writev = opts.writev878 if (opts.write) this._write = opts.write879 if (opts.final) this._final = opts.final880 if (opts.eagerOpen) this._writableState.updateNextTick()881 }882 }883 884 cork () {885 this._duplexState |= WRITE_CORKED886 }887 888 uncork () {889 this._duplexState &= WRITE_NOT_CORKED890 this._writableState.updateNextTick()891 }892 893 _writev (batch, cb) {894 cb(null)895 }896 897 _write (data, cb) {898 this._writableState.autoBatch(data, cb)899 }900 901 _final (cb) {902 cb(null)903 }904 905 static isBackpressured (ws) {906 return (ws._duplexState & WRITE_BACKPRESSURE_STATUS) !== 0907 }908 909 static drained (ws) {910 if (ws.destroyed) return Promise.resolve(false)911 const state = ws._writableState912 const pending = (isWritev(ws) ? Math.min(1, state.queue.length) : state.queue.length)913 const writes = pending + ((ws._duplexState & WRITE_WRITING) ? 1 : 0)914 if (writes === 0) return Promise.resolve(true)915 if (state.drains === null) state.drains = []916 return new Promise((resolve) => {917 state.drains.push({ writes, resolve })918 })919 }920 921 write (data) {922 this._writableState.updateNextTick()923 return this._writableState.push(data)924 }925 926 end (data) {927 this._writableState.updateNextTick()928 this._writableState.end(data)929 return this930 }931}932 933class Duplex extends Readable { // and Writable934 constructor (opts) {935 super(opts)936 937 this._duplexState = OPENING | (this._duplexState & READ_READ_AHEAD)938 this._writableState = new WritableState(this, opts)939 940 if (opts) {941 if (opts.writev) this._writev = opts.writev942 if (opts.write) this._write = opts.write943 if (opts.final) this._final = opts.final944 }945 }946 947 cork () {948 this._duplexState |= WRITE_CORKED949 }950 951 uncork () {952 this._duplexState &= WRITE_NOT_CORKED953 this._writableState.updateNextTick()954 }955 956 _writev (batch, cb) {957 cb(null)958 }959 960 _write (data, cb) {961 this._writableState.autoBatch(data, cb)962 }963 964 _final (cb) {965 cb(null)966 }967 968 write (data) {969 this._writableState.updateNextTick()970 return this._writableState.push(data)971 }972 973 end (data) {974 this._writableState.updateNextTick()975 this._writableState.end(data)976 return this977 }978}979 980class Transform extends Duplex {981 constructor (opts) {982 super(opts)983 this._transformState = new TransformState(this)984 985 if (opts) {986 if (opts.transform) this._transform = opts.transform987 if (opts.flush) this._flush = opts.flush988 }989 }990 991 _write (data, cb) {992 if (this._readableState.buffered >= this._readableState.highWaterMark) {993 this._transformState.data = data994 } else {995 this._transform(data, this._transformState.afterTransform)996 }997 }998 999 _read (cb) {1000 if (this._transformState.data !== null) {1001 const data = this._transformState.data1002 this._transformState.data = null1003 cb(null)1004 this._transform(data, this._transformState.afterTransform)1005 } else {1006 cb(null)1007 }1008 }1009 1010 destroy (err) {1011 super.destroy(err)1012 if (this._transformState.data !== null) {1013 this._transformState.data = null1014 this._transformState.afterTransform()1015 }1016 }1017 1018 _transform (data, cb) {1019 cb(null, data)1020 }1021 1022 _flush (cb) {1023 cb(null)1024 }1025 1026 _final (cb) {1027 this._transformState.afterFinal = cb1028 this._flush(transformAfterFlush.bind(this))1029 }1030}1031 1032class PassThrough extends Transform {}1033 1034function transformAfterFlush (err, data) {1035 const cb = this._transformState.afterFinal1036 if (err) return cb(err)1037 if (data !== null && data !== undefined) this.push(data)1038 this.push(null)1039 cb(null)1040}1041 1042function pipelinePromise (...streams) {1043 return new Promise((resolve, reject) => {1044 return pipeline(...streams, (err) => {1045 if (err) return reject(err)1046 resolve()1047 })1048 })1049}1050 1051function pipeline (stream, ...streams) {1052 const all = Array.isArray(stream) ? [...stream, ...streams] : [stream, ...streams]1053 const done = (all.length && typeof all[all.length - 1] === 'function') ? all.pop() : null1054 1055 if (all.length < 2) throw new Error('Pipeline requires at least 2 streams')1056 1057 let src = all[0]1058 let dest = null1059 let error = null1060 1061 for (let i = 1; i < all.length; i++) {1062 dest = all[i]1063 1064 if (isStreamx(src)) {1065 src.pipe(dest, onerror)1066 } else {1067 errorHandle(src, true, i > 1, onerror)1068 src.pipe(dest)1069 }1070 1071 src = dest1072 }1073 1074 if (done) {1075 let fin = false1076 1077 const autoDestroy = isStreamx(dest) || !!(dest._writableState && dest._writableState.autoDestroy)1078 1079 dest.on('error', (err) => {1080 if (error === null) error = err1081 })1082 1083 dest.on('finish', () => {1084 fin = true1085 if (!autoDestroy) done(error)1086 })1087 1088 if (autoDestroy) {1089 dest.on('close', () => done(error || (fin ? null : PREMATURE_CLOSE)))1090 }1091 }1092 1093 return dest1094 1095 function errorHandle (s, rd, wr, onerror) {1096 s.on('error', onerror)1097 s.on('close', onclose)1098 1099 function onclose () {1100 if (rd && s._readableState && !s._readableState.ended) return onerror(PREMATURE_CLOSE)1101 if (wr && s._writableState && !s._writableState.ended) return onerror(PREMATURE_CLOSE)1102 }1103 }1104 1105 function onerror (err) {1106 if (!err || error) return1107 error = err1108 1109 for (const s of all) {1110 s.destroy(err)1111 }1112 }1113}1114 1115function echo (s) {1116 return s1117}1118 1119function isStream (stream) {1120 return !!stream._readableState || !!stream._writableState1121}1122 1123function isStreamx (stream) {1124 return typeof stream._duplexState === 'number' && isStream(stream)1125}1126 1127function isEnded (stream) {1128 return !!stream._readableState && stream._readableState.ended1129}1130 1131function isFinished (stream) {1132 return !!stream._writableState && stream._writableState.ended1133}1134 1135function getStreamError (stream, opts = {}) {1136 const err = (stream._readableState && stream._readableState.error) || (stream._writableState && stream._writableState.error)1137 1138 // avoid implicit errors by default1139 return (!opts.all && err === STREAM_DESTROYED) ? null : err1140}1141 1142function isReadStreamx (stream) {1143 return isStreamx(stream) && stream.readable1144}1145 1146function isDisturbed (stream) {1147 return (stream._duplexState & OPENING) !== OPENING || (stream._duplexState & ACTIVE_OR_TICKING) !== 01148}1149 1150function isTypedArray (data) {1151 return typeof data === 'object' && data !== null && typeof data.byteLength === 'number'1152}1153 1154function defaultByteLength (data) {1155 return isTypedArray(data) ? data.byteLength : 10241156}1157 1158function noop () {}1159 1160function abort () {1161 this.destroy(new Error('Stream aborted.'))1162}1163 1164function isWritev (s) {1165 return s._writev !== Writable.prototype._writev && s._writev !== Duplex.prototype._writev1166}1167 1168module.exports = {1169 pipeline,1170 pipelinePromise,1171 isStream,1172 isStreamx,1173 isEnded,1174 isFinished,1175 isDisturbed,1176 getStreamError,1177 Stream,1178 Writable,1179 Readable,1180 Duplex,1181 Transform,1182 // Export PassThrough for compatibility with Node.js core's stream module1183 PassThrough1184}1185 