| 1 | # minipass
|
|---|
| 2 |
|
|---|
| 3 | A _very_ minimal implementation of a [PassThrough
|
|---|
| 4 | stream](https://nodejs.org/api/stream.html#stream_class_stream_passthrough)
|
|---|
| 5 |
|
|---|
| 6 | [It's very
|
|---|
| 7 | fast](https://docs.google.com/spreadsheets/d/1oObKSrVwLX_7Ut4Z6g3fZW-AX1j1-k6w-cDsrkaSbHM/edit#gid=0)
|
|---|
| 8 | for objects, strings, and buffers.
|
|---|
| 9 |
|
|---|
| 10 | Supports `pipe()`ing (including multi-`pipe()` and backpressure transmission),
|
|---|
| 11 | buffering data until either a `data` event handler or `pipe()` is added (so
|
|---|
| 12 | you don't lose the first chunk), and most other cases where PassThrough is
|
|---|
| 13 | a good idea.
|
|---|
| 14 |
|
|---|
| 15 | There is a `read()` method, but it's much more efficient to consume data
|
|---|
| 16 | from this stream via `'data'` events or by calling `pipe()` into some other
|
|---|
| 17 | stream. Calling `read()` requires the buffer to be flattened in some
|
|---|
| 18 | cases, which requires copying memory.
|
|---|
| 19 |
|
|---|
| 20 | If you set `objectMode: true` in the options, then whatever is written will
|
|---|
| 21 | be emitted. Otherwise, it'll do a minimal amount of Buffer copying to
|
|---|
| 22 | ensure proper Streams semantics when `read(n)` is called.
|
|---|
| 23 |
|
|---|
| 24 | `objectMode` can also be set by doing `stream.objectMode = true`, or by
|
|---|
| 25 | writing any non-string/non-buffer data. `objectMode` cannot be set to
|
|---|
| 26 | false once it is set.
|
|---|
| 27 |
|
|---|
| 28 | This is not a `through` or `through2` stream. It doesn't transform the
|
|---|
| 29 | data, it just passes it right through. If you want to transform the data,
|
|---|
| 30 | extend the class, and override the `write()` method. Once you're done
|
|---|
| 31 | transforming the data however you want, call `super.write()` with the
|
|---|
| 32 | transform output.
|
|---|
| 33 |
|
|---|
| 34 | For some examples of streams that extend Minipass in various ways, check
|
|---|
| 35 | out:
|
|---|
| 36 |
|
|---|
| 37 | - [minizlib](http://npm.im/minizlib)
|
|---|
| 38 | - [fs-minipass](http://npm.im/fs-minipass)
|
|---|
| 39 | - [tar](http://npm.im/tar)
|
|---|
| 40 | - [minipass-collect](http://npm.im/minipass-collect)
|
|---|
| 41 | - [minipass-flush](http://npm.im/minipass-flush)
|
|---|
| 42 | - [minipass-pipeline](http://npm.im/minipass-pipeline)
|
|---|
| 43 | - [tap](http://npm.im/tap)
|
|---|
| 44 | - [tap-parser](http://npm.im/tap-parser)
|
|---|
| 45 | - [treport](http://npm.im/treport)
|
|---|
| 46 | - [minipass-fetch](http://npm.im/minipass-fetch)
|
|---|
| 47 | - [pacote](http://npm.im/pacote)
|
|---|
| 48 | - [make-fetch-happen](http://npm.im/make-fetch-happen)
|
|---|
| 49 | - [cacache](http://npm.im/cacache)
|
|---|
| 50 | - [ssri](http://npm.im/ssri)
|
|---|
| 51 | - [npm-registry-fetch](http://npm.im/npm-registry-fetch)
|
|---|
| 52 | - [minipass-json-stream](http://npm.im/minipass-json-stream)
|
|---|
| 53 | - [minipass-sized](http://npm.im/minipass-sized)
|
|---|
| 54 |
|
|---|
| 55 | ## Differences from Node.js Streams
|
|---|
| 56 |
|
|---|
| 57 | There are several things that make Minipass streams different from (and in
|
|---|
| 58 | some ways superior to) Node.js core streams.
|
|---|
| 59 |
|
|---|
| 60 | Please read these caveats if you are familiar with node-core streams and
|
|---|
| 61 | intend to use Minipass streams in your programs.
|
|---|
| 62 |
|
|---|
| 63 | You can avoid most of these differences entirely (for a very
|
|---|
| 64 | small performance penalty) by setting `{async: true}` in the
|
|---|
| 65 | constructor options.
|
|---|
| 66 |
|
|---|
| 67 | ### Timing
|
|---|
| 68 |
|
|---|
| 69 | Minipass streams are designed to support synchronous use-cases. Thus, data
|
|---|
| 70 | is emitted as soon as it is available, always. It is buffered until read,
|
|---|
| 71 | but no longer. Another way to look at it is that Minipass streams are
|
|---|
| 72 | exactly as synchronous as the logic that writes into them.
|
|---|
| 73 |
|
|---|
| 74 | This can be surprising if your code relies on `PassThrough.write()` always
|
|---|
| 75 | providing data on the next tick rather than the current one, or being able
|
|---|
| 76 | to call `resume()` and not have the entire buffer disappear immediately.
|
|---|
| 77 |
|
|---|
| 78 | However, without this synchronicity guarantee, there would be no way for
|
|---|
| 79 | Minipass to achieve the speeds it does, or support the synchronous use
|
|---|
| 80 | cases that it does. Simply put, waiting takes time.
|
|---|
| 81 |
|
|---|
| 82 | This non-deferring approach makes Minipass streams much easier to reason
|
|---|
| 83 | about, especially in the context of Promises and other flow-control
|
|---|
| 84 | mechanisms.
|
|---|
| 85 |
|
|---|
| 86 | Example:
|
|---|
| 87 |
|
|---|
| 88 | ```js
|
|---|
| 89 | const Minipass = require('minipass')
|
|---|
| 90 | const stream = new Minipass({ async: true })
|
|---|
| 91 | stream.on('data', () => console.log('data event'))
|
|---|
| 92 | console.log('before write')
|
|---|
| 93 | stream.write('hello')
|
|---|
| 94 | console.log('after write')
|
|---|
| 95 | // output:
|
|---|
| 96 | // before write
|
|---|
| 97 | // data event
|
|---|
| 98 | // after write
|
|---|
| 99 | ```
|
|---|
| 100 |
|
|---|
| 101 | ### Exception: Async Opt-In
|
|---|
| 102 |
|
|---|
| 103 | If you wish to have a Minipass stream with behavior that more
|
|---|
| 104 | closely mimics Node.js core streams, you can set the stream in
|
|---|
| 105 | async mode either by setting `async: true` in the constructor
|
|---|
| 106 | options, or by setting `stream.async = true` later on.
|
|---|
| 107 |
|
|---|
| 108 | ```js
|
|---|
| 109 | const Minipass = require('minipass')
|
|---|
| 110 | const asyncStream = new Minipass({ async: true })
|
|---|
| 111 | asyncStream.on('data', () => console.log('data event'))
|
|---|
| 112 | console.log('before write')
|
|---|
| 113 | asyncStream.write('hello')
|
|---|
| 114 | console.log('after write')
|
|---|
| 115 | // output:
|
|---|
| 116 | // before write
|
|---|
| 117 | // after write
|
|---|
| 118 | // data event <-- this is deferred until the next tick
|
|---|
| 119 | ```
|
|---|
| 120 |
|
|---|
| 121 | Switching _out_ of async mode is unsafe, as it could cause data
|
|---|
| 122 | corruption, and so is not enabled. Example:
|
|---|
| 123 |
|
|---|
| 124 | ```js
|
|---|
| 125 | const Minipass = require('minipass')
|
|---|
| 126 | const stream = new Minipass({ encoding: 'utf8' })
|
|---|
| 127 | stream.on('data', chunk => console.log(chunk))
|
|---|
| 128 | stream.async = true
|
|---|
| 129 | console.log('before writes')
|
|---|
| 130 | stream.write('hello')
|
|---|
| 131 | setStreamSyncAgainSomehow(stream) // <-- this doesn't actually exist!
|
|---|
| 132 | stream.write('world')
|
|---|
| 133 | console.log('after writes')
|
|---|
| 134 | // hypothetical output would be:
|
|---|
| 135 | // before writes
|
|---|
| 136 | // world
|
|---|
| 137 | // after writes
|
|---|
| 138 | // hello
|
|---|
| 139 | // NOT GOOD!
|
|---|
| 140 | ```
|
|---|
| 141 |
|
|---|
| 142 | To avoid this problem, once set into async mode, any attempt to
|
|---|
| 143 | make the stream sync again will be ignored.
|
|---|
| 144 |
|
|---|
| 145 | ```js
|
|---|
| 146 | const Minipass = require('minipass')
|
|---|
| 147 | const stream = new Minipass({ encoding: 'utf8' })
|
|---|
| 148 | stream.on('data', chunk => console.log(chunk))
|
|---|
| 149 | stream.async = true
|
|---|
| 150 | console.log('before writes')
|
|---|
| 151 | stream.write('hello')
|
|---|
| 152 | stream.async = false // <-- no-op, stream already async
|
|---|
| 153 | stream.write('world')
|
|---|
| 154 | console.log('after writes')
|
|---|
| 155 | // actual output:
|
|---|
| 156 | // before writes
|
|---|
| 157 | // after writes
|
|---|
| 158 | // hello
|
|---|
| 159 | // world
|
|---|
| 160 | ```
|
|---|
| 161 |
|
|---|
| 162 | ### No High/Low Water Marks
|
|---|
| 163 |
|
|---|
| 164 | Node.js core streams will optimistically fill up a buffer, returning `true`
|
|---|
| 165 | on all writes until the limit is hit, even if the data has nowhere to go.
|
|---|
| 166 | Then, they will not attempt to draw more data in until the buffer size dips
|
|---|
| 167 | below a minimum value.
|
|---|
| 168 |
|
|---|
| 169 | Minipass streams are much simpler. The `write()` method will return `true`
|
|---|
| 170 | if the data has somewhere to go (which is to say, given the timing
|
|---|
| 171 | guarantees, that the data is already there by the time `write()` returns).
|
|---|
| 172 |
|
|---|
| 173 | If the data has nowhere to go, then `write()` returns false, and the data
|
|---|
| 174 | sits in a buffer, to be drained out immediately as soon as anyone consumes
|
|---|
| 175 | it.
|
|---|
| 176 |
|
|---|
| 177 | Since nothing is ever buffered unnecessarily, there is much less
|
|---|
| 178 | copying data, and less bookkeeping about buffer capacity levels.
|
|---|
| 179 |
|
|---|
| 180 | ### Hazards of Buffering (or: Why Minipass Is So Fast)
|
|---|
| 181 |
|
|---|
| 182 | Since data written to a Minipass stream is immediately written all the way
|
|---|
| 183 | through the pipeline, and `write()` always returns true/false based on
|
|---|
| 184 | whether the data was fully flushed, backpressure is communicated
|
|---|
| 185 | immediately to the upstream caller. This minimizes buffering.
|
|---|
| 186 |
|
|---|
| 187 | Consider this case:
|
|---|
| 188 |
|
|---|
| 189 | ```js
|
|---|
| 190 | const {PassThrough} = require('stream')
|
|---|
| 191 | const p1 = new PassThrough({ highWaterMark: 1024 })
|
|---|
| 192 | const p2 = new PassThrough({ highWaterMark: 1024 })
|
|---|
| 193 | const p3 = new PassThrough({ highWaterMark: 1024 })
|
|---|
| 194 | const p4 = new PassThrough({ highWaterMark: 1024 })
|
|---|
| 195 |
|
|---|
| 196 | p1.pipe(p2).pipe(p3).pipe(p4)
|
|---|
| 197 | p4.on('data', () => console.log('made it through'))
|
|---|
| 198 |
|
|---|
| 199 | // this returns false and buffers, then writes to p2 on next tick (1)
|
|---|
| 200 | // p2 returns false and buffers, pausing p1, then writes to p3 on next tick (2)
|
|---|
| 201 | // p3 returns false and buffers, pausing p2, then writes to p4 on next tick (3)
|
|---|
| 202 | // p4 returns false and buffers, pausing p3, then emits 'data' and 'drain'
|
|---|
| 203 | // on next tick (4)
|
|---|
| 204 | // p3 sees p4's 'drain' event, and calls resume(), emitting 'resume' and
|
|---|
| 205 | // 'drain' on next tick (5)
|
|---|
| 206 | // p2 sees p3's 'drain', calls resume(), emits 'resume' and 'drain' on next tick (6)
|
|---|
| 207 | // p1 sees p2's 'drain', calls resume(), emits 'resume' and 'drain' on next
|
|---|
| 208 | // tick (7)
|
|---|
| 209 |
|
|---|
| 210 | p1.write(Buffer.alloc(2048)) // returns false
|
|---|
| 211 | ```
|
|---|
| 212 |
|
|---|
| 213 | Along the way, the data was buffered and deferred at each stage, and
|
|---|
| 214 | multiple event deferrals happened, for an unblocked pipeline where it was
|
|---|
| 215 | perfectly safe to write all the way through!
|
|---|
| 216 |
|
|---|
| 217 | Furthermore, setting a `highWaterMark` of `1024` might lead someone reading
|
|---|
| 218 | the code to think an advisory maximum of 1KiB is being set for the
|
|---|
| 219 | pipeline. However, the actual advisory buffering level is the _sum_ of
|
|---|
| 220 | `highWaterMark` values, since each one has its own bucket.
|
|---|
| 221 |
|
|---|
| 222 | Consider the Minipass case:
|
|---|
| 223 |
|
|---|
| 224 | ```js
|
|---|
| 225 | const m1 = new Minipass()
|
|---|
| 226 | const m2 = new Minipass()
|
|---|
| 227 | const m3 = new Minipass()
|
|---|
| 228 | const m4 = new Minipass()
|
|---|
| 229 |
|
|---|
| 230 | m1.pipe(m2).pipe(m3).pipe(m4)
|
|---|
| 231 | m4.on('data', () => console.log('made it through'))
|
|---|
| 232 |
|
|---|
| 233 | // m1 is flowing, so it writes the data to m2 immediately
|
|---|
| 234 | // m2 is flowing, so it writes the data to m3 immediately
|
|---|
| 235 | // m3 is flowing, so it writes the data to m4 immediately
|
|---|
| 236 | // m4 is flowing, so it fires the 'data' event immediately, returns true
|
|---|
| 237 | // m4's write returned true, so m3 is still flowing, returns true
|
|---|
| 238 | // m3's write returned true, so m2 is still flowing, returns true
|
|---|
| 239 | // m2's write returned true, so m1 is still flowing, returns true
|
|---|
| 240 | // No event deferrals or buffering along the way!
|
|---|
| 241 |
|
|---|
| 242 | m1.write(Buffer.alloc(2048)) // returns true
|
|---|
| 243 | ```
|
|---|
| 244 |
|
|---|
| 245 | It is extremely unlikely that you _don't_ want to buffer any data written,
|
|---|
| 246 | or _ever_ buffer data that can be flushed all the way through. Neither
|
|---|
| 247 | node-core streams nor Minipass ever fail to buffer written data, but
|
|---|
| 248 | node-core streams do a lot of unnecessary buffering and pausing.
|
|---|
| 249 |
|
|---|
| 250 | As always, the faster implementation is the one that does less stuff and
|
|---|
| 251 | waits less time to do it.
|
|---|
| 252 |
|
|---|
| 253 | ### Immediately emit `end` for empty streams (when not paused)
|
|---|
| 254 |
|
|---|
| 255 | If a stream is not paused, and `end()` is called before writing any data
|
|---|
| 256 | into it, then it will emit `end` immediately.
|
|---|
| 257 |
|
|---|
| 258 | If you have logic that occurs on the `end` event which you don't want to
|
|---|
| 259 | potentially happen immediately (for example, closing file descriptors,
|
|---|
| 260 | moving on to the next entry in an archive parse stream, etc.) then be sure
|
|---|
| 261 | to call `stream.pause()` on creation, and then `stream.resume()` once you
|
|---|
| 262 | are ready to respond to the `end` event.
|
|---|
| 263 |
|
|---|
| 264 | However, this is _usually_ not a problem because:
|
|---|
| 265 |
|
|---|
| 266 | ### Emit `end` When Asked
|
|---|
| 267 |
|
|---|
| 268 | One hazard of immediately emitting `'end'` is that you may not yet have had
|
|---|
| 269 | a chance to add a listener. In order to avoid this hazard, Minipass
|
|---|
| 270 | streams safely re-emit the `'end'` event if a new listener is added after
|
|---|
| 271 | `'end'` has been emitted.
|
|---|
| 272 |
|
|---|
| 273 | Ie, if you do `stream.on('end', someFunction)`, and the stream has already
|
|---|
| 274 | emitted `end`, then it will call the handler right away. (You can think of
|
|---|
| 275 | this somewhat like attaching a new `.then(fn)` to a previously-resolved
|
|---|
| 276 | Promise.)
|
|---|
| 277 |
|
|---|
| 278 | To prevent calling handlers multiple times who would not expect multiple
|
|---|
| 279 | ends to occur, all listeners are removed from the `'end'` event whenever it
|
|---|
| 280 | is emitted.
|
|---|
| 281 |
|
|---|
| 282 | ### Emit `error` When Asked
|
|---|
| 283 |
|
|---|
| 284 | The most recent error object passed to the `'error'` event is
|
|---|
| 285 | stored on the stream. If a new `'error'` event handler is added,
|
|---|
| 286 | and an error was previously emitted, then the event handler will
|
|---|
| 287 | be called immediately (or on `process.nextTick` in the case of
|
|---|
| 288 | async streams).
|
|---|
| 289 |
|
|---|
| 290 | This makes it much more difficult to end up trying to interact
|
|---|
| 291 | with a broken stream, if the error handler is added after an
|
|---|
| 292 | error was previously emitted.
|
|---|
| 293 |
|
|---|
| 294 | ### Impact of "immediate flow" on Tee-streams
|
|---|
| 295 |
|
|---|
| 296 | A "tee stream" is a stream piping to multiple destinations:
|
|---|
| 297 |
|
|---|
| 298 | ```js
|
|---|
| 299 | const tee = new Minipass()
|
|---|
| 300 | t.pipe(dest1)
|
|---|
| 301 | t.pipe(dest2)
|
|---|
| 302 | t.write('foo') // goes to both destinations
|
|---|
| 303 | ```
|
|---|
| 304 |
|
|---|
| 305 | Since Minipass streams _immediately_ process any pending data through the
|
|---|
| 306 | pipeline when a new pipe destination is added, this can have surprising
|
|---|
| 307 | effects, especially when a stream comes in from some other function and may
|
|---|
| 308 | or may not have data in its buffer.
|
|---|
| 309 |
|
|---|
| 310 | ```js
|
|---|
| 311 | // WARNING! WILL LOSE DATA!
|
|---|
| 312 | const src = new Minipass()
|
|---|
| 313 | src.write('foo')
|
|---|
| 314 | src.pipe(dest1) // 'foo' chunk flows to dest1 immediately, and is gone
|
|---|
| 315 | src.pipe(dest2) // gets nothing!
|
|---|
| 316 | ```
|
|---|
| 317 |
|
|---|
| 318 | One solution is to create a dedicated tee-stream junction that pipes to
|
|---|
| 319 | both locations, and then pipe to _that_ instead.
|
|---|
| 320 |
|
|---|
| 321 | ```js
|
|---|
| 322 | // Safe example: tee to both places
|
|---|
| 323 | const src = new Minipass()
|
|---|
| 324 | src.write('foo')
|
|---|
| 325 | const tee = new Minipass()
|
|---|
| 326 | tee.pipe(dest1)
|
|---|
| 327 | tee.pipe(dest2)
|
|---|
| 328 | src.pipe(tee) // tee gets 'foo', pipes to both locations
|
|---|
| 329 | ```
|
|---|
| 330 |
|
|---|
| 331 | The same caveat applies to `on('data')` event listeners. The first one
|
|---|
| 332 | added will _immediately_ receive all of the data, leaving nothing for the
|
|---|
| 333 | second:
|
|---|
| 334 |
|
|---|
| 335 | ```js
|
|---|
| 336 | // WARNING! WILL LOSE DATA!
|
|---|
| 337 | const src = new Minipass()
|
|---|
| 338 | src.write('foo')
|
|---|
| 339 | src.on('data', handler1) // receives 'foo' right away
|
|---|
| 340 | src.on('data', handler2) // nothing to see here!
|
|---|
| 341 | ```
|
|---|
| 342 |
|
|---|
| 343 | Using a dedicated tee-stream can be used in this case as well:
|
|---|
| 344 |
|
|---|
| 345 | ```js
|
|---|
| 346 | // Safe example: tee to both data handlers
|
|---|
| 347 | const src = new Minipass()
|
|---|
| 348 | src.write('foo')
|
|---|
| 349 | const tee = new Minipass()
|
|---|
| 350 | tee.on('data', handler1)
|
|---|
| 351 | tee.on('data', handler2)
|
|---|
| 352 | src.pipe(tee)
|
|---|
| 353 | ```
|
|---|
| 354 |
|
|---|
| 355 | All of the hazards in this section are avoided by setting `{
|
|---|
| 356 | async: true }` in the Minipass constructor, or by setting
|
|---|
| 357 | `stream.async = true` afterwards. Note that this does add some
|
|---|
| 358 | overhead, so should only be done in cases where you are willing
|
|---|
| 359 | to lose a bit of performance in order to avoid having to refactor
|
|---|
| 360 | program logic.
|
|---|
| 361 |
|
|---|
| 362 | ## USAGE
|
|---|
| 363 |
|
|---|
| 364 | It's a stream! Use it like a stream and it'll most likely do what you
|
|---|
| 365 | want.
|
|---|
| 366 |
|
|---|
| 367 | ```js
|
|---|
| 368 | const Minipass = require('minipass')
|
|---|
| 369 | const mp = new Minipass(options) // optional: { encoding, objectMode }
|
|---|
| 370 | mp.write('foo')
|
|---|
| 371 | mp.pipe(someOtherStream)
|
|---|
| 372 | mp.end('bar')
|
|---|
| 373 | ```
|
|---|
| 374 |
|
|---|
| 375 | ### OPTIONS
|
|---|
| 376 |
|
|---|
| 377 | * `encoding` How would you like the data coming _out_ of the stream to be
|
|---|
| 378 | encoded? Accepts any values that can be passed to `Buffer.toString()`.
|
|---|
| 379 | * `objectMode` Emit data exactly as it comes in. This will be flipped on
|
|---|
| 380 | by default if you write() something other than a string or Buffer at any
|
|---|
| 381 | point. Setting `objectMode: true` will prevent setting any encoding
|
|---|
| 382 | value.
|
|---|
| 383 | * `async` Defaults to `false`. Set to `true` to defer data
|
|---|
| 384 | emission until next tick. This reduces performance slightly,
|
|---|
| 385 | but makes Minipass streams use timing behavior closer to Node
|
|---|
| 386 | core streams. See [Timing](#timing) for more details.
|
|---|
| 387 |
|
|---|
| 388 | ### API
|
|---|
| 389 |
|
|---|
| 390 | Implements the user-facing portions of Node.js's `Readable` and `Writable`
|
|---|
| 391 | streams.
|
|---|
| 392 |
|
|---|
| 393 | ### Methods
|
|---|
| 394 |
|
|---|
| 395 | * `write(chunk, [encoding], [callback])` - Put data in. (Note that, in the
|
|---|
| 396 | base Minipass class, the same data will come out.) Returns `false` if
|
|---|
| 397 | the stream will buffer the next write, or true if it's still in "flowing"
|
|---|
| 398 | mode.
|
|---|
| 399 | * `end([chunk, [encoding]], [callback])` - Signal that you have no more
|
|---|
| 400 | data to write. This will queue an `end` event to be fired when all the
|
|---|
| 401 | data has been consumed.
|
|---|
| 402 | * `setEncoding(encoding)` - Set the encoding for data coming of the stream.
|
|---|
| 403 | This can only be done once.
|
|---|
| 404 | * `pause()` - No more data for a while, please. This also prevents `end`
|
|---|
| 405 | from being emitted for empty streams until the stream is resumed.
|
|---|
| 406 | * `resume()` - Resume the stream. If there's data in the buffer, it is all
|
|---|
| 407 | discarded. Any buffered events are immediately emitted.
|
|---|
| 408 | * `pipe(dest)` - Send all output to the stream provided. When
|
|---|
| 409 | data is emitted, it is immediately written to any and all pipe
|
|---|
| 410 | destinations. (Or written on next tick in `async` mode.)
|
|---|
| 411 | * `unpipe(dest)` - Stop piping to the destination stream. This
|
|---|
| 412 | is immediate, meaning that any asynchronously queued data will
|
|---|
| 413 | _not_ make it to the destination when running in `async` mode.
|
|---|
| 414 | * `options.end` - Boolean, end the destination stream when
|
|---|
| 415 | the source stream ends. Default `true`.
|
|---|
| 416 | * `options.proxyErrors` - Boolean, proxy `error` events from
|
|---|
| 417 | the source stream to the destination stream. Note that
|
|---|
| 418 | errors are _not_ proxied after the pipeline terminates,
|
|---|
| 419 | either due to the source emitting `'end'` or manually
|
|---|
| 420 | unpiping with `src.unpipe(dest)`. Default `false`.
|
|---|
| 421 | * `on(ev, fn)`, `emit(ev, fn)` - Minipass streams are EventEmitters. Some
|
|---|
| 422 | events are given special treatment, however. (See below under "events".)
|
|---|
| 423 | * `promise()` - Returns a Promise that resolves when the stream emits
|
|---|
| 424 | `end`, or rejects if the stream emits `error`.
|
|---|
| 425 | * `collect()` - Return a Promise that resolves on `end` with an array
|
|---|
| 426 | containing each chunk of data that was emitted, or rejects if the stream
|
|---|
| 427 | emits `error`. Note that this consumes the stream data.
|
|---|
| 428 | * `concat()` - Same as `collect()`, but concatenates the data into a single
|
|---|
| 429 | Buffer object. Will reject the returned promise if the stream is in
|
|---|
| 430 | objectMode, or if it goes into objectMode by the end of the data.
|
|---|
| 431 | * `read(n)` - Consume `n` bytes of data out of the buffer. If `n` is not
|
|---|
| 432 | provided, then consume all of it. If `n` bytes are not available, then
|
|---|
| 433 | it returns null. **Note** consuming streams in this way is less
|
|---|
| 434 | efficient, and can lead to unnecessary Buffer copying.
|
|---|
| 435 | * `destroy([er])` - Destroy the stream. If an error is provided, then an
|
|---|
| 436 | `'error'` event is emitted. If the stream has a `close()` method, and
|
|---|
| 437 | has not emitted a `'close'` event yet, then `stream.close()` will be
|
|---|
| 438 | called. Any Promises returned by `.promise()`, `.collect()` or
|
|---|
| 439 | `.concat()` will be rejected. After being destroyed, writing to the
|
|---|
| 440 | stream will emit an error. No more data will be emitted if the stream is
|
|---|
| 441 | destroyed, even if it was previously buffered.
|
|---|
| 442 |
|
|---|
| 443 | ### Properties
|
|---|
| 444 |
|
|---|
| 445 | * `bufferLength` Read-only. Total number of bytes buffered, or in the case
|
|---|
| 446 | of objectMode, the total number of objects.
|
|---|
| 447 | * `encoding` The encoding that has been set. (Setting this is equivalent
|
|---|
| 448 | to calling `setEncoding(enc)` and has the same prohibition against
|
|---|
| 449 | setting multiple times.)
|
|---|
| 450 | * `flowing` Read-only. Boolean indicating whether a chunk written to the
|
|---|
| 451 | stream will be immediately emitted.
|
|---|
| 452 | * `emittedEnd` Read-only. Boolean indicating whether the end-ish events
|
|---|
| 453 | (ie, `end`, `prefinish`, `finish`) have been emitted. Note that
|
|---|
| 454 | listening on any end-ish event will immediateyl re-emit it if it has
|
|---|
| 455 | already been emitted.
|
|---|
| 456 | * `writable` Whether the stream is writable. Default `true`. Set to
|
|---|
| 457 | `false` when `end()`
|
|---|
| 458 | * `readable` Whether the stream is readable. Default `true`.
|
|---|
| 459 | * `buffer` A [yallist](http://npm.im/yallist) linked list of chunks written
|
|---|
| 460 | to the stream that have not yet been emitted. (It's probably a bad idea
|
|---|
| 461 | to mess with this.)
|
|---|
| 462 | * `pipes` A [yallist](http://npm.im/yallist) linked list of streams that
|
|---|
| 463 | this stream is piping into. (It's probably a bad idea to mess with
|
|---|
| 464 | this.)
|
|---|
| 465 | * `destroyed` A getter that indicates whether the stream was destroyed.
|
|---|
| 466 | * `paused` True if the stream has been explicitly paused, otherwise false.
|
|---|
| 467 | * `objectMode` Indicates whether the stream is in `objectMode`. Once set
|
|---|
| 468 | to `true`, it cannot be set to `false`.
|
|---|
| 469 |
|
|---|
| 470 | ### Events
|
|---|
| 471 |
|
|---|
| 472 | * `data` Emitted when there's data to read. Argument is the data to read.
|
|---|
| 473 | This is never emitted while not flowing. If a listener is attached, that
|
|---|
| 474 | will resume the stream.
|
|---|
| 475 | * `end` Emitted when there's no more data to read. This will be emitted
|
|---|
| 476 | immediately for empty streams when `end()` is called. If a listener is
|
|---|
| 477 | attached, and `end` was already emitted, then it will be emitted again.
|
|---|
| 478 | All listeners are removed when `end` is emitted.
|
|---|
| 479 | * `prefinish` An end-ish event that follows the same logic as `end` and is
|
|---|
| 480 | emitted in the same conditions where `end` is emitted. Emitted after
|
|---|
| 481 | `'end'`.
|
|---|
| 482 | * `finish` An end-ish event that follows the same logic as `end` and is
|
|---|
| 483 | emitted in the same conditions where `end` is emitted. Emitted after
|
|---|
| 484 | `'prefinish'`.
|
|---|
| 485 | * `close` An indication that an underlying resource has been released.
|
|---|
| 486 | Minipass does not emit this event, but will defer it until after `end`
|
|---|
| 487 | has been emitted, since it throws off some stream libraries otherwise.
|
|---|
| 488 | * `drain` Emitted when the internal buffer empties, and it is again
|
|---|
| 489 | suitable to `write()` into the stream.
|
|---|
| 490 | * `readable` Emitted when data is buffered and ready to be read by a
|
|---|
| 491 | consumer.
|
|---|
| 492 | * `resume` Emitted when stream changes state from buffering to flowing
|
|---|
| 493 | mode. (Ie, when `resume` is called, `pipe` is called, or a `data` event
|
|---|
| 494 | listener is added.)
|
|---|
| 495 |
|
|---|
| 496 | ### Static Methods
|
|---|
| 497 |
|
|---|
| 498 | * `Minipass.isStream(stream)` Returns `true` if the argument is a stream,
|
|---|
| 499 | and false otherwise. To be considered a stream, the object must be
|
|---|
| 500 | either an instance of Minipass, or an EventEmitter that has either a
|
|---|
| 501 | `pipe()` method, or both `write()` and `end()` methods. (Pretty much any
|
|---|
| 502 | stream in node-land will return `true` for this.)
|
|---|
| 503 |
|
|---|
| 504 | ## EXAMPLES
|
|---|
| 505 |
|
|---|
| 506 | Here are some examples of things you can do with Minipass streams.
|
|---|
| 507 |
|
|---|
| 508 | ### simple "are you done yet" promise
|
|---|
| 509 |
|
|---|
| 510 | ```js
|
|---|
| 511 | mp.promise().then(() => {
|
|---|
| 512 | // stream is finished
|
|---|
| 513 | }, er => {
|
|---|
| 514 | // stream emitted an error
|
|---|
| 515 | })
|
|---|
| 516 | ```
|
|---|
| 517 |
|
|---|
| 518 | ### collecting
|
|---|
| 519 |
|
|---|
| 520 | ```js
|
|---|
| 521 | mp.collect().then(all => {
|
|---|
| 522 | // all is an array of all the data emitted
|
|---|
| 523 | // encoding is supported in this case, so
|
|---|
| 524 | // so the result will be a collection of strings if
|
|---|
| 525 | // an encoding is specified, or buffers/objects if not.
|
|---|
| 526 | //
|
|---|
| 527 | // In an async function, you may do
|
|---|
| 528 | // const data = await stream.collect()
|
|---|
| 529 | })
|
|---|
| 530 | ```
|
|---|
| 531 |
|
|---|
| 532 | ### collecting into a single blob
|
|---|
| 533 |
|
|---|
| 534 | This is a bit slower because it concatenates the data into one chunk for
|
|---|
| 535 | you, but if you're going to do it yourself anyway, it's convenient this
|
|---|
| 536 | way:
|
|---|
| 537 |
|
|---|
| 538 | ```js
|
|---|
| 539 | mp.concat().then(onebigchunk => {
|
|---|
| 540 | // onebigchunk is a string if the stream
|
|---|
| 541 | // had an encoding set, or a buffer otherwise.
|
|---|
| 542 | })
|
|---|
| 543 | ```
|
|---|
| 544 |
|
|---|
| 545 | ### iteration
|
|---|
| 546 |
|
|---|
| 547 | You can iterate over streams synchronously or asynchronously in platforms
|
|---|
| 548 | that support it.
|
|---|
| 549 |
|
|---|
| 550 | Synchronous iteration will end when the currently available data is
|
|---|
| 551 | consumed, even if the `end` event has not been reached. In string and
|
|---|
| 552 | buffer mode, the data is concatenated, so unless multiple writes are
|
|---|
| 553 | occurring in the same tick as the `read()`, sync iteration loops will
|
|---|
| 554 | generally only have a single iteration.
|
|---|
| 555 |
|
|---|
| 556 | To consume chunks in this way exactly as they have been written, with no
|
|---|
| 557 | flattening, create the stream with the `{ objectMode: true }` option.
|
|---|
| 558 |
|
|---|
| 559 | ```js
|
|---|
| 560 | const mp = new Minipass({ objectMode: true })
|
|---|
| 561 | mp.write('a')
|
|---|
| 562 | mp.write('b')
|
|---|
| 563 | for (let letter of mp) {
|
|---|
| 564 | console.log(letter) // a, b
|
|---|
| 565 | }
|
|---|
| 566 | mp.write('c')
|
|---|
| 567 | mp.write('d')
|
|---|
| 568 | for (let letter of mp) {
|
|---|
| 569 | console.log(letter) // c, d
|
|---|
| 570 | }
|
|---|
| 571 | mp.write('e')
|
|---|
| 572 | mp.end()
|
|---|
| 573 | for (let letter of mp) {
|
|---|
| 574 | console.log(letter) // e
|
|---|
| 575 | }
|
|---|
| 576 | for (let letter of mp) {
|
|---|
| 577 | console.log(letter) // nothing
|
|---|
| 578 | }
|
|---|
| 579 | ```
|
|---|
| 580 |
|
|---|
| 581 | Asynchronous iteration will continue until the end event is reached,
|
|---|
| 582 | consuming all of the data.
|
|---|
| 583 |
|
|---|
| 584 | ```js
|
|---|
| 585 | const mp = new Minipass({ encoding: 'utf8' })
|
|---|
| 586 |
|
|---|
| 587 | // some source of some data
|
|---|
| 588 | let i = 5
|
|---|
| 589 | const inter = setInterval(() => {
|
|---|
| 590 | if (i-- > 0)
|
|---|
| 591 | mp.write(Buffer.from('foo\n', 'utf8'))
|
|---|
| 592 | else {
|
|---|
| 593 | mp.end()
|
|---|
| 594 | clearInterval(inter)
|
|---|
| 595 | }
|
|---|
| 596 | }, 100)
|
|---|
| 597 |
|
|---|
| 598 | // consume the data with asynchronous iteration
|
|---|
| 599 | async function consume () {
|
|---|
| 600 | for await (let chunk of mp) {
|
|---|
| 601 | console.log(chunk)
|
|---|
| 602 | }
|
|---|
| 603 | return 'ok'
|
|---|
| 604 | }
|
|---|
| 605 |
|
|---|
| 606 | consume().then(res => console.log(res))
|
|---|
| 607 | // logs `foo\n` 5 times, and then `ok`
|
|---|
| 608 | ```
|
|---|
| 609 |
|
|---|
| 610 | ### subclass that `console.log()`s everything written into it
|
|---|
| 611 |
|
|---|
| 612 | ```js
|
|---|
| 613 | class Logger extends Minipass {
|
|---|
| 614 | write (chunk, encoding, callback) {
|
|---|
| 615 | console.log('WRITE', chunk, encoding)
|
|---|
| 616 | return super.write(chunk, encoding, callback)
|
|---|
| 617 | }
|
|---|
| 618 | end (chunk, encoding, callback) {
|
|---|
| 619 | console.log('END', chunk, encoding)
|
|---|
| 620 | return super.end(chunk, encoding, callback)
|
|---|
| 621 | }
|
|---|
| 622 | }
|
|---|
| 623 |
|
|---|
| 624 | someSource.pipe(new Logger()).pipe(someDest)
|
|---|
| 625 | ```
|
|---|
| 626 |
|
|---|
| 627 | ### same thing, but using an inline anonymous class
|
|---|
| 628 |
|
|---|
| 629 | ```js
|
|---|
| 630 | // js classes are fun
|
|---|
| 631 | someSource
|
|---|
| 632 | .pipe(new (class extends Minipass {
|
|---|
| 633 | emit (ev, ...data) {
|
|---|
| 634 | // let's also log events, because debugging some weird thing
|
|---|
| 635 | console.log('EMIT', ev)
|
|---|
| 636 | return super.emit(ev, ...data)
|
|---|
| 637 | }
|
|---|
| 638 | write (chunk, encoding, callback) {
|
|---|
| 639 | console.log('WRITE', chunk, encoding)
|
|---|
| 640 | return super.write(chunk, encoding, callback)
|
|---|
| 641 | }
|
|---|
| 642 | end (chunk, encoding, callback) {
|
|---|
| 643 | console.log('END', chunk, encoding)
|
|---|
| 644 | return super.end(chunk, encoding, callback)
|
|---|
| 645 | }
|
|---|
| 646 | }))
|
|---|
| 647 | .pipe(someDest)
|
|---|
| 648 | ```
|
|---|
| 649 |
|
|---|
| 650 | ### subclass that defers 'end' for some reason
|
|---|
| 651 |
|
|---|
| 652 | ```js
|
|---|
| 653 | class SlowEnd extends Minipass {
|
|---|
| 654 | emit (ev, ...args) {
|
|---|
| 655 | if (ev === 'end') {
|
|---|
| 656 | console.log('going to end, hold on a sec')
|
|---|
| 657 | setTimeout(() => {
|
|---|
| 658 | console.log('ok, ready to end now')
|
|---|
| 659 | super.emit('end', ...args)
|
|---|
| 660 | }, 100)
|
|---|
| 661 | } else {
|
|---|
| 662 | return super.emit(ev, ...args)
|
|---|
| 663 | }
|
|---|
| 664 | }
|
|---|
| 665 | }
|
|---|
| 666 | ```
|
|---|
| 667 |
|
|---|
| 668 | ### transform that creates newline-delimited JSON
|
|---|
| 669 |
|
|---|
| 670 | ```js
|
|---|
| 671 | class NDJSONEncode extends Minipass {
|
|---|
| 672 | write (obj, cb) {
|
|---|
| 673 | try {
|
|---|
| 674 | // JSON.stringify can throw, emit an error on that
|
|---|
| 675 | return super.write(JSON.stringify(obj) + '\n', 'utf8', cb)
|
|---|
| 676 | } catch (er) {
|
|---|
| 677 | this.emit('error', er)
|
|---|
| 678 | }
|
|---|
| 679 | }
|
|---|
| 680 | end (obj, cb) {
|
|---|
| 681 | if (typeof obj === 'function') {
|
|---|
| 682 | cb = obj
|
|---|
| 683 | obj = undefined
|
|---|
| 684 | }
|
|---|
| 685 | if (obj !== undefined) {
|
|---|
| 686 | this.write(obj)
|
|---|
| 687 | }
|
|---|
| 688 | return super.end(cb)
|
|---|
| 689 | }
|
|---|
| 690 | }
|
|---|
| 691 | ```
|
|---|
| 692 |
|
|---|
| 693 | ### transform that parses newline-delimited JSON
|
|---|
| 694 |
|
|---|
| 695 | ```js
|
|---|
| 696 | class NDJSONDecode extends Minipass {
|
|---|
| 697 | constructor (options) {
|
|---|
| 698 | // always be in object mode, as far as Minipass is concerned
|
|---|
| 699 | super({ objectMode: true })
|
|---|
| 700 | this._jsonBuffer = ''
|
|---|
| 701 | }
|
|---|
| 702 | write (chunk, encoding, cb) {
|
|---|
| 703 | if (typeof chunk === 'string' &&
|
|---|
| 704 | typeof encoding === 'string' &&
|
|---|
| 705 | encoding !== 'utf8') {
|
|---|
| 706 | chunk = Buffer.from(chunk, encoding).toString()
|
|---|
| 707 | } else if (Buffer.isBuffer(chunk))
|
|---|
| 708 | chunk = chunk.toString()
|
|---|
| 709 | }
|
|---|
| 710 | if (typeof encoding === 'function') {
|
|---|
| 711 | cb = encoding
|
|---|
| 712 | }
|
|---|
| 713 | const jsonData = (this._jsonBuffer + chunk).split('\n')
|
|---|
| 714 | this._jsonBuffer = jsonData.pop()
|
|---|
| 715 | for (let i = 0; i < jsonData.length; i++) {
|
|---|
| 716 | try {
|
|---|
| 717 | // JSON.parse can throw, emit an error on that
|
|---|
| 718 | super.write(JSON.parse(jsonData[i]))
|
|---|
| 719 | } catch (er) {
|
|---|
| 720 | this.emit('error', er)
|
|---|
| 721 | continue
|
|---|
| 722 | }
|
|---|
| 723 | }
|
|---|
| 724 | if (cb)
|
|---|
| 725 | cb()
|
|---|
| 726 | }
|
|---|
| 727 | }
|
|---|
| 728 | ```
|
|---|