source: node_modules/minipass/index.js@ 81bc7da

finki-main main
Last change on this file since 81bc7da was 81bc7da, checked in by Klimentina Efremova <klimentina08642@…>, 3 months ago

Initial commit

  • Property mode set to 100644
File size: 16.2 KB
Line 
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}
Note: See TracBrowser for help on using the repository browser.