Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
index.js650 linesDownload Raw Back to minipass
1'use strict'
2const proc = typeof process === 'object' && process ? process : {
3  stdout: null,
4  stderr: null,
5}
6const EE = require('events')
7const Stream = require('stream')
8const SD = require('string_decoder').StringDecoder
9
10const EOF = Symbol('EOF')
11const MAYBE_EMIT_END = Symbol('maybeEmitEnd')
12const EMITTED_END = Symbol('emittedEnd')
13const EMITTING_END = Symbol('emittingEnd')
14const EMITTED_ERROR = Symbol('emittedError')
15const CLOSED = Symbol('closed')
16const READ = Symbol('read')
17const FLUSH = Symbol('flush')
18const FLUSHCHUNK = Symbol('flushChunk')
19const ENCODING = Symbol('encoding')
20const DECODER = Symbol('decoder')
21const FLOWING = Symbol('flowing')
22const PAUSED = Symbol('paused')
23const RESUME = Symbol('resume')
24const BUFFERLENGTH = Symbol('bufferLength')
25const BUFFERPUSH = Symbol('bufferPush')
26const BUFFERSHIFT = Symbol('bufferShift')
27const OBJECTMODE = Symbol('objectMode')
28const DESTROYED = Symbol('destroyed')
29const EMITDATA = Symbol('emitData')
30const EMITEND = Symbol('emitEnd')
31const EMITEND2 = Symbol('emitEnd2')
32const ASYNC = Symbol('async')
33
34const defer = fn => Promise.resolve().then(fn)
35
36// TODO remove when Node v8 support drops
37const doIter = global._MP_NO_ITERATOR_SYMBOLS_  !== '1'
38const ASYNCITERATOR = doIter && Symbol.asyncIterator
39  || Symbol('asyncIterator not implemented')
40const ITERATOR = doIter && Symbol.iterator
41  || Symbol('iterator not implemented')
42
43// events that mean 'the stream is over'
44// these are treated specially, and re-emitted
45// if they are listened for after emitting.
46const isEndish = ev =>
47  ev === 'end' ||
48  ev === 'finish' ||
49  ev === 'prefinish'
50
51const isArrayBuffer = b => b instanceof ArrayBuffer ||
52  typeof b === 'object' &&
53  b.constructor &&
54  b.constructor.name === 'ArrayBuffer' &&
55  b.byteLength >= 0
56
57const isArrayBufferView = b => !Buffer.isBuffer(b) && ArrayBuffer.isView(b)
58
59class Pipe {
60  constructor (src, dest, opts) {
61    this.src = src
62    this.dest = dest
63    this.opts = opts
64    this.ondrain = () => src[RESUME]()
65    dest.on('drain', this.ondrain)
66  }
67  unpipe () {
68    this.dest.removeListener('drain', this.ondrain)
69  }
70  // istanbul ignore next - only here for the prototype
71  proxyErrors () {}
72  end () {
73    this.unpipe()
74    if (this.opts.end)
75      this.dest.end()
76  }
77}
78
79class PipeProxyErrors extends Pipe {
80  unpipe () {
81    this.src.removeListener('error', this.proxyErrors)
82    super.unpipe()
83  }
84  constructor (src, dest, opts) {
85    super(src, dest, opts)
86    this.proxyErrors = er => dest.emit('error', er)
87    src.on('error', this.proxyErrors)
88  }
89}
90
91module.exports = class Minipass extends Stream {
92  constructor (options) {
93    super()
94    this[FLOWING] = false
95    // whether we're explicitly paused
96    this[PAUSED] = false
97    this.pipes = []
98    this.buffer = []
99    this[OBJECTMODE] = options && options.objectMode || false
100    if (this[OBJECTMODE])
101      this[ENCODING] = null
102    else
103      this[ENCODING] = options && options.encoding || null
104    if (this[ENCODING] === 'buffer')
105      this[ENCODING] = null
106    this[ASYNC] = options && !!options.async || false
107    this[DECODER] = this[ENCODING] ? new SD(this[ENCODING]) : null
108    this[EOF] = false
109    this[EMITTED_END] = false
110    this[EMITTING_END] = false
111    this[CLOSED] = false
112    this[EMITTED_ERROR] = null
113    this.writable = true
114    this.readable = true
115    this[BUFFERLENGTH] = 0
116    this[DESTROYED] = false
117  }
118
119  get bufferLength () { return this[BUFFERLENGTH] }
120
121  get encoding () { return this[ENCODING] }
122  set encoding (enc) {
123    if (this[OBJECTMODE])
124      throw new Error('cannot set encoding in objectMode')
125
126    if (this[ENCODING] && enc !== this[ENCODING] &&
127        (this[DECODER] && this[DECODER].lastNeed || this[BUFFERLENGTH]))
128      throw new Error('cannot change encoding')
129
130    if (this[ENCODING] !== enc) {
131      this[DECODER] = enc ? new SD(enc) : null
132      if (this.buffer.length)
133        this.buffer = this.buffer.map(chunk => this[DECODER].write(chunk))
134    }
135
136    this[ENCODING] = enc
137  }
138
139  setEncoding (enc) {
140    this.encoding = enc
141  }
142
143  get objectMode () { return this[OBJECTMODE] }
144  set objectMode (om) { this[OBJECTMODE] = this[OBJECTMODE] || !!om }
145
146  get ['async'] () { return this[ASYNC] }
147  set ['async'] (a) { this[ASYNC] = this[ASYNC] || !!a }
148
149  write (chunk, encoding, cb) {
150    if (this[EOF])
151      throw new Error('write after end')
152
153    if (this[DESTROYED]) {
154      this.emit('error', Object.assign(
155        new Error('Cannot call write after a stream was destroyed'),
156        { code: 'ERR_STREAM_DESTROYED' }
157      ))
158      return true
159    }
160
161    if (typeof encoding === 'function')
162      cb = encoding, encoding = 'utf8'
163
164    if (!encoding)
165      encoding = 'utf8'
166
167    const fn = this[ASYNC] ? defer : f => f()
168
169    // convert array buffers and typed array views into buffers
170    // at some point in the future, we may want to do the opposite!
171    // leave strings and buffers as-is
172    // anything else switches us into object mode
173    if (!this[OBJECTMODE] && !Buffer.isBuffer(chunk)) {
174      if (isArrayBufferView(chunk))
175        chunk = Buffer.from(chunk.buffer, chunk.byteOffset, chunk.byteLength)
176      else if (isArrayBuffer(chunk))
177        chunk = Buffer.from(chunk)
178      else if (typeof chunk !== 'string')
179        // use the setter so we throw if we have encoding set
180        this.objectMode = true
181    }
182
183    // handle object mode up front, since it's simpler
184    // this yields better performance, fewer checks later.
185    if (this[OBJECTMODE]) {
186      /* istanbul ignore if - maybe impossible? */
187      if (this.flowing && this[BUFFERLENGTH] !== 0)
188        this[FLUSH](true)
189
190      if (this.flowing)
191        this.emit('data', chunk)
192      else
193        this[BUFFERPUSH](chunk)
194
195      if (this[BUFFERLENGTH] !== 0)
196        this.emit('readable')
197
198      if (cb)
199        fn(cb)
200
201      return this.flowing
202    }
203
204    // at this point the chunk is a buffer or string
205    // don't buffer it up or send it to the decoder
206    if (!chunk.length) {
207      if (this[BUFFERLENGTH] !== 0)
208        this.emit('readable')
209      if (cb)
210        fn(cb)
211      return this.flowing
212    }
213
214    // fast-path writing strings of same encoding to a stream with
215    // an empty buffer, skipping the buffer/decoder dance
216    if (typeof chunk === 'string' &&
217        // unless it is a string already ready for us to use
218        !(encoding === this[ENCODING] && !this[DECODER].lastNeed)) {
219      chunk = Buffer.from(chunk, encoding)
220    }
221
222    if (Buffer.isBuffer(chunk) && this[ENCODING])
223      chunk = this[DECODER].write(chunk)
224
225    // Note: flushing CAN potentially switch us into not-flowing mode
226    if (this.flowing && this[BUFFERLENGTH] !== 0)
227      this[FLUSH](true)
228
229    if (this.flowing)
230      this.emit('data', chunk)
231    else
232      this[BUFFERPUSH](chunk)
233
234    if (this[BUFFERLENGTH] !== 0)
235      this.emit('readable')
236
237    if (cb)
238      fn(cb)
239
240    return this.flowing
241  }
242
243  read (n) {
244    if (this[DESTROYED])
245      return null
246
247    if (this[BUFFERLENGTH] === 0 || n === 0 || n > this[BUFFERLENGTH]) {
248      this[MAYBE_EMIT_END]()
249      return null
250    }
251
252    if (this[OBJECTMODE])
253      n = null
254
255    if (this.buffer.length > 1 && !this[OBJECTMODE]) {
256      if (this.encoding)
257        this.buffer = [this.buffer.join('')]
258      else
259        this.buffer = [Buffer.concat(this.buffer, this[BUFFERLENGTH])]
260    }
261
262    const ret = this[READ](n || null, this.buffer[0])
263    this[MAYBE_EMIT_END]()
264    return ret
265  }
266
267  [READ] (n, chunk) {
268    if (n === chunk.length || n === null)
269      this[BUFFERSHIFT]()
270    else {
271      this.buffer[0] = chunk.slice(n)
272      chunk = chunk.slice(0, n)
273      this[BUFFERLENGTH] -= n
274    }
275
276    this.emit('data', chunk)
277
278    if (!this.buffer.length && !this[EOF])
279      this.emit('drain')
280
281    return chunk
282  }
283
284  end (chunk, encoding, cb) {
285    if (typeof chunk === 'function')
286      cb = chunk, chunk = null
287    if (typeof encoding === 'function')
288      cb = encoding, encoding = 'utf8'
289    if (chunk)
290      this.write(chunk, encoding)
291    if (cb)
292      this.once('end', cb)
293    this[EOF] = true
294    this.writable = false
295
296    // if we haven't written anything, then go ahead and emit,
297    // even if we're not reading.
298    // we'll re-emit if a new 'end' listener is added anyway.
299    // This makes MP more suitable to write-only use cases.
300    if (this.flowing || !this[PAUSED])
301      this[MAYBE_EMIT_END]()
302    return this
303  }
304
305  // don't let the internal resume be overwritten
306  [RESUME] () {
307    if (this[DESTROYED])
308      return
309
310    this[PAUSED] = false
311    this[FLOWING] = true
312    this.emit('resume')
313    if (this.buffer.length)
314      this[FLUSH]()
315    else if (this[EOF])
316      this[MAYBE_EMIT_END]()
317    else
318      this.emit('drain')
319  }
320
321  resume () {
322    return this[RESUME]()
323  }
324
325  pause () {
326    this[FLOWING] = false
327    this[PAUSED] = true
328  }
329
330  get destroyed () {
331    return this[DESTROYED]
332  }
333
334  get flowing () {
335    return this[FLOWING]
336  }
337
338  get paused () {
339    return this[PAUSED]
340  }
341
342  [BUFFERPUSH] (chunk) {
343    if (this[OBJECTMODE])
344      this[BUFFERLENGTH] += 1
345    else
346      this[BUFFERLENGTH] += chunk.length
347    this.buffer.push(chunk)
348  }
349
350  [BUFFERSHIFT] () {
351    if (this.buffer.length) {
352      if (this[OBJECTMODE])
353        this[BUFFERLENGTH] -= 1
354      else
355        this[BUFFERLENGTH] -= this.buffer[0].length
356    }
357    return this.buffer.shift()
358  }
359
360  [FLUSH] (noDrain) {
361    do {} while (this[FLUSHCHUNK](this[BUFFERSHIFT]()))
362
363    if (!noDrain && !this.buffer.length && !this[EOF])
364      this.emit('drain')
365  }
366
367  [FLUSHCHUNK] (chunk) {
368    return chunk ? (this.emit('data', chunk), this.flowing) : false
369  }
370
371  pipe (dest, opts) {
372    if (this[DESTROYED])
373      return
374
375    const ended = this[EMITTED_END]
376    opts = opts || {}
377    if (dest === proc.stdout || dest === proc.stderr)
378      opts.end = false
379    else
380      opts.end = opts.end !== false
381    opts.proxyErrors = !!opts.proxyErrors
382
383    // piping an ended stream ends immediately
384    if (ended) {
385      if (opts.end)
386        dest.end()
387    } else {
388      this.pipes.push(!opts.proxyErrors ? new Pipe(this, dest, opts)
389        : new PipeProxyErrors(this, dest, opts))
390      if (this[ASYNC])
391        defer(() => this[RESUME]())
392      else
393        this[RESUME]()
394    }
395
396    return dest
397  }
398
399  unpipe (dest) {
400    const p = this.pipes.find(p => p.dest === dest)
401    if (p) {
402      this.pipes.splice(this.pipes.indexOf(p), 1)
403      p.unpipe()
404    }
405  }
406
407  addListener (ev, fn) {
408    return this.on(ev, fn)
409  }
410
411  on (ev, fn) {
412    const ret = super.on(ev, fn)
413    if (ev === 'data' && !this.pipes.length && !this.flowing)
414      this[RESUME]()
415    else if (ev === 'readable' && this[BUFFERLENGTH] !== 0)
416      super.emit('readable')
417    else if (isEndish(ev) && this[EMITTED_END]) {
418      super.emit(ev)
419      this.removeAllListeners(ev)
420    } else if (ev === 'error' && this[EMITTED_ERROR]) {
421      if (this[ASYNC])
422        defer(() => fn.call(this, this[EMITTED_ERROR]))
423      else
424        fn.call(this, this[EMITTED_ERROR])
425    }
426    return ret
427  }
428
429  get emittedEnd () {
430    return this[EMITTED_END]
431  }
432
433  [MAYBE_EMIT_END] () {
434    if (!this[EMITTING_END] &&
435        !this[EMITTED_END] &&
436        !this[DESTROYED] &&
437        this.buffer.length === 0 &&
438        this[EOF]) {
439      this[EMITTING_END] = true
440      this.emit('end')
441      this.emit('prefinish')
442      this.emit('finish')
443      if (this[CLOSED])
444        this.emit('close')
445      this[EMITTING_END] = false
446    }
447  }
448
449  emit (ev, data, ...extra) {
450    // error and close are only events allowed after calling destroy()
451    if (ev !== 'error' && ev !== 'close' && ev !== DESTROYED && this[DESTROYED])
452      return
453    else if (ev === 'data') {
454      return !data ? false
455        : this[ASYNC] ? defer(() => this[EMITDATA](data))
456        : this[EMITDATA](data)
457    } else if (ev === 'end') {
458      return this[EMITEND]()
459    } else if (ev === 'close') {
460      this[CLOSED] = true
461      // don't emit close before 'end' and 'finish'
462      if (!this[EMITTED_END] && !this[DESTROYED])
463        return
464      const ret = super.emit('close')
465      this.removeAllListeners('close')
466      return ret
467    } else if (ev === 'error') {
468      this[EMITTED_ERROR] = data
469      const ret = super.emit('error', data)
470      this[MAYBE_EMIT_END]()
471      return ret
472    } else if (ev === 'resume') {
473      const ret = super.emit('resume')
474      this[MAYBE_EMIT_END]()
475      return ret
476    } else if (ev === 'finish' || ev === 'prefinish') {
477      const ret = super.emit(ev)
478      this.removeAllListeners(ev)
479      return ret
480    }
481
482    // Some other unknown event
483    const ret = super.emit(ev, data, ...extra)
484    this[MAYBE_EMIT_END]()
485    return ret
486  }
487
488  [EMITDATA] (data) {
489    for (const p of this.pipes) {
490      if (p.dest.write(data) === false)
491        this.pause()
492    }
493    const ret = super.emit('data', data)
494    this[MAYBE_EMIT_END]()
495    return ret
496  }
497
498  [EMITEND] () {
499    if (this[EMITTED_END])
500      return
501
502    this[EMITTED_END] = true
503    this.readable = false
504    if (this[ASYNC])
505      defer(() => this[EMITEND2]())
506    else
507      this[EMITEND2]()
508  }
509
510  [EMITEND2] () {
511    if (this[DECODER]) {
512      const data = this[DECODER].end()
513      if (data) {
514        for (const p of this.pipes) {
515          p.dest.write(data)
516        }
517        super.emit('data', data)
518      }
519    }
520
521    for (const p of this.pipes) {
522      p.end()
523    }
524    const ret = super.emit('end')
525    this.removeAllListeners('end')
526    return ret
527  }
528
529  // const all = await stream.collect()
530  collect () {
531    const buf = []
532    if (!this[OBJECTMODE])
533      buf.dataLength = 0
534    // set the promise first, in case an error is raised
535    // by triggering the flow here.
536    const p = this.promise()
537    this.on('data', c => {
538      buf.push(c)
539      if (!this[OBJECTMODE])
540        buf.dataLength += c.length
541    })
542    return p.then(() => buf)
543  }
544
545  // const data = await stream.concat()
546  concat () {
547    return this[OBJECTMODE]
548      ? Promise.reject(new Error('cannot concat in objectMode'))
549      : this.collect().then(buf =>
550          this[OBJECTMODE]
551            ? Promise.reject(new Error('cannot concat in objectMode'))
552            : this[ENCODING] ? buf.join('') : Buffer.concat(buf, buf.dataLength))
553  }
554
555  // stream.promise().then(() => done, er => emitted error)
556  promise () {
557    return new Promise((resolve, reject) => {
558      this.on(DESTROYED, () => reject(new Error('stream destroyed')))
559      this.on('error', er => reject(er))
560      this.on('end', () => resolve())
561    })
562  }
563
564  // for await (let chunk of stream)
565  [ASYNCITERATOR] () {
566    const next = () => {
567      const res = this.read()
568      if (res !== null)
569        return Promise.resolve({ done: false, value: res })
570
571      if (this[EOF])
572        return Promise.resolve({ done: true })
573
574      let resolve = null
575      let reject = null
576      const onerr = er => {
577        this.removeListener('data', ondata)
578        this.removeListener('end', onend)
579        reject(er)
580      }
581      const ondata = value => {
582        this.removeListener('error', onerr)
583        this.removeListener('end', onend)
584        this.pause()
585        resolve({ value: value, done: !!this[EOF] })
586      }
587      const onend = () => {
588        this.removeListener('error', onerr)
589        this.removeListener('data', ondata)
590        resolve({ done: true })
591      }
592      const ondestroy = () => onerr(new Error('stream destroyed'))
593      return new Promise((res, rej) => {
594        reject = rej
595        resolve = res
596        this.once(DESTROYED, ondestroy)
597        this.once('error', onerr)
598        this.once('end', onend)
599        this.once('data', ondata)
600      })
601    }
602
603    return { next }
604  }
605
606  // for (let chunk of stream)
607  [ITERATOR] () {
608    const next = () => {
609      const value = this.read()
610      const done = value === null
611      return { value, done }
612    }
613    return { next }
614  }
615
616  destroy (er) {
617    if (this[DESTROYED]) {
618      if (er)
619        this.emit('error', er)
620      else
621        this.emit(DESTROYED)
622      return this
623    }
624
625    this[DESTROYED] = true
626
627    // throw away all buffered data, it's never coming out
628    this.buffer.length = 0
629    this[BUFFERLENGTH] = 0
630
631    if (typeof this.close === 'function' && !this[CLOSED])
632      this.close()
633
634    if (er)
635      this.emit('error', er)
636    else // if no error to emit, still reject pending promises
637      this.emit(DESTROYED)
638
639    return this
640  }
641
642  static isStream (s) {
643    return !!s && (s instanceof Minipass || s instanceof Stream ||
644      s instanceof EE && (
645        typeof s.pipe === 'function' || // readable
646        (typeof s.write === 'function' && typeof s.end === 'function') // writable
647      ))
648  }
649}
650 
codekingpro/portable-devtools · Team Ai