CoolFace
Datasetpublic

basant307/AI_Governance_Project

sourceHugging Faceapache-2.0updated 2mo agoView on Hugging Face
0likes48downloads
index.js1185 linesDownload Raw Back to streamx
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 
basant307/AI_Governance_Project · CoolFace