| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent a7dfa43 commit 1787bfa
2 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -33,52 +33,26 @@ const { | |||
| 33 | 33 | isIterable, | |
| 34 | 34 | isReadableNodeStream, | |
| 35 | 35 | isNodeStream, | |
| 36 | - isReadableFinished, | ||
| 37 | 36 | } = require('internal/streams/utils'); | |
| 38 | 37 | const { AbortController } = require('internal/abort_controller'); | |
| 39 | 38 | ||
| 40 | 39 | let PassThrough; | |
| 41 | 40 | let Readable; | |
| 42 | 41 | ||
| 43 | - function destroyer(stream, reading, writing, callback) { | ||
| 44 | - callback = once(callback); | ||
| 45 | - | ||
| 42 | + function destroyer(stream, reading, writing) { | ||
| 46 | 43 | let finished = false; | |
| 47 | 44 | stream.on('close', () => { | |
| 48 | 45 | finished = true; | |
| 49 | 46 | }); | |
| 50 | 47 | ||
| 51 | 48 | eos(stream, { readable: reading, writable: writing }, (err) => { | |
| 52 | 49 | finished = !err; | |
| 53 | - | ||
| 54 | - const rState = stream._readableState; | ||
| 55 | - if ( | ||
| 56 | - err && | ||
| 57 | - err.code === 'ERR_STREAM_PREMATURE_CLOSE' && | ||
| 58 | - reading && | ||
| 59 | - (rState && rState.ended && !rState.errored && !rState.errorEmitted) | ||
| 60 | - ) { | ||
| 61 | - // Some readable streams will emit 'close' before 'end'. However, since | ||
| 62 | - // this is on the readable side 'end' should still be emitted if the | ||
| 63 | - // stream has been ended and no error emitted. This should be allowed in | ||
| 64 | - // favor of backwards compatibility. Since the stream is piped to a | ||
| 65 | - // destination this should not result in any observable difference. | ||
| 66 | - // We don't need to check if this is a writable premature close since | ||
| 67 | - // eos will only fail with premature close on the reading side for | ||
| 68 | - // duplex streams. | ||
| 69 | - stream | ||
| 70 | - .once('end', callback) | ||
| 71 | - .once('error', callback); | ||
| 72 | - } else { | ||
| 73 | - callback(err); | ||
| 74 | - } | ||
| 75 | 50 | }); | |
| 76 | 51 | ||
| 77 | 52 | return (err) => { | |
| 78 | 53 | if (finished) return; | |
| 79 | 54 | finished = true; | |
| 80 | - destroyImpl.destroyer(stream, err); | ||
| 81 | - callback(err || new ERR_STREAM_DESTROYED('pipe')); | ||
| 55 | + destroyImpl.destroyer(stream, err || new ERR_STREAM_DESTROYED('pipe')); | ||
| 82 | 56 | }; | |
| 83 | 57 | } | |
| 84 | 58 | ||
@@ -109,7 +83,7 @@ async function* fromReadable(val) { | |||
| 109 | 83 | yield* Readable.prototype[SymbolAsyncIterator].call(val); | |
| 110 | 84 | } | |
| 111 | 85 | ||
| 112 | - async function pump(iterable, writable, finish, opts) { | ||
| 86 | + async function pump(iterable, writable, finish, { end }) { | ||
| 113 | 87 | let error; | |
| 114 | 88 | let onresolve = null; | |
| 115 | 89 | ||
@@ -153,7 +127,7 @@ async function pump(iterable, writable, finish, opts) { | |||
| 153 | 127 | } | |
| 154 | 128 | } | |
| 155 | 129 | ||
| 156 | - if (opts?.end !== false) { | ||
| 130 | + if (end) { | ||
| 157 | 131 | writable.end(); | |
| 158 | 132 | } | |
| 159 | 133 | ||
@@ -220,7 +194,7 @@ function pipelineImpl(streams, callback, opts) { | |||
| 220 | 194 | ac.abort(); | |
| 221 | 195 | ||
| 222 | 196 | if (final) { | |
| 223 | - callback(error, value); | ||
| 197 | + process.nextTick(callback, error, value); | ||
| 224 | 198 | } | |
| 225 | 199 | } | |
| 226 | 200 | ||
@@ -233,18 +207,19 @@ function pipelineImpl(streams, callback, opts) { | |||
| 233 | 207 | ||
| 234 | 208 | if (isNodeStream(stream)) { | |
| 235 | 209 | if (end) { | |
| 236 | - finishCount++; | ||
| 237 | - destroys.push(destroyer(stream, reading, writing, (err) => { | ||
| 238 | - if (!err && !reading && isReadableFinished(stream, false)) { | ||
| 239 | - stream.read(0); | ||
| 240 | - destroyer(stream, true, writing, finish); | ||
| 241 | - } else { | ||
| 242 | - finish(err); | ||
| 243 | - } | ||
| 244 | - })); | ||
| 245 | - } else { | ||
| 246 | - stream.on('error', finish); | ||
| 210 | + destroys.push(destroyer(stream, reading, writing)); | ||
| 247 | 211 | } | |
| 212 | + | ||
| 213 | + // Catch stream errors that occur after pipe/pump has completed. | ||
| 214 | + stream.on('error', (err) => { | ||
| 215 | + if ( | ||
| 216 | + err && | ||
| 217 | + err.name !== 'AbortError' && | ||
| 218 | + err.code !== 'ERR_STREAM_PREMATURE_CLOSE' | ||
| 219 | + ) { | ||
| 220 | + finish(err); | ||
| 221 | + } | ||
| 222 | + }); | ||
| 248 | 223 | } | |
| 249 | 224 | ||
| 250 | 225 | if (i === 0) { | |
@@ -286,15 +261,18 @@ function pipelineImpl(streams, callback, opts) { | |||
| 286 | 261 | // second use. | |
| 287 | 262 | const then = ret?.then; | |
| 288 | 263 | if (typeof then === 'function') { | |
| 264 | + finishCount++; | ||
| 289 | 265 | then.call(ret, | |
| 290 | 266 | (val) => { | |
| 291 | 267 | value = val; | |
| 292 | 268 | pt.write(val); | |
| 293 | 269 | if (end) { | |
| 294 | 270 | pt.end(); | |
| 295 | 271 | } | |
| 272 | + process.nextTick(finish); | ||
| 296 | 273 | }, (err) => { | |
| 297 | 274 | pt.destroy(err); | |
| 275 | + process.nextTick(finish, err); | ||
| 298 | 276 | }, | |
| 299 | 277 | ); | |
| 300 | 278 | } else if (isIterable(ret, true)) { | |
@@ -307,24 +285,18 @@ function pipelineImpl(streams, callback, opts) { | |||
| 307 | 285 | ||
| 308 | 286 | ret = pt; | |
| 309 | 287 | ||
| 310 | - finishCount++; | ||
| 311 | - destroys.push(destroyer(ret, false, true, finish)); | ||
| 288 | + destroys.push(destroyer(ret, false, true)); | ||
| 312 | 289 | } | |
| 313 | 290 | } else if (isNodeStream(stream)) { | |
| 314 | 291 | if (isReadableNodeStream(ret)) { | |
| 315 | - ret.pipe(stream, { end }); | ||
| 316 | - | ||
| 317 | - // Compat. Before node v10.12.0 stdio used to throw an error so | ||
| 318 | - // pipe() did/does not end() stdio destinations. | ||
| 319 | - // Now they allow it but "secretly" don't close the underlying fd. | ||
| 320 | - if (stream === process.stdout || stream === process.stderr) { | ||
| 321 | - ret.on('end', () => stream.end()); | ||
| 322 | - } | ||
| 323 | - } else { | ||
| 324 | - ret = makeAsyncIterable(ret); | ||
| 325 | - | ||
| 292 | + finishCount += 2; | ||
| 293 | + pipe(ret, stream, finish, { end }); | ||
| 294 | + } else if (isIterable(ret)) { | ||
| 326 | 295 | finishCount++; | |
| 327 | 296 | pump(ret, stream, finish, { end }); | |
| 297 | + } else { | ||
| 298 | + throw new ERR_INVALID_ARG_TYPE( | ||
| 299 | + 'val', ['Readable', 'Iterable', 'AsyncIterable'], ret); | ||
| 328 | 300 | } | |
| 329 | 301 | ret = stream; | |
| 330 | 302 | } else { | |
@@ -339,4 +311,41 @@ function pipelineImpl(streams, callback, opts) { | |||
| 339 | 311 | return ret; | |
| 340 | 312 | } | |
| 341 | 313 | ||
| 314 | + function pipe(src, dst, finish, { end }) { | ||
| 315 | + src.pipe(dst, { end }); | ||
| 316 | + | ||
| 317 | + if (end) { | ||
| 318 | + // Compat. Before node v10.12.0 stdio used to throw an error so | ||
| 319 | + // pipe() did/does not end() stdio destinations. | ||
| 320 | + // Now they allow it but "secretly" don't close the underlying fd. | ||
| 321 | + src.once('end', () => dst.end()); | ||
| 322 | + } else { | ||
| 323 | + finish(); | ||
| 324 | + } | ||
| 325 | + | ||
| 326 | + eos(src, { readable: true, writable: false }, (err) => { | ||
| 327 | + const rState = src._readableState; | ||
| 328 | + if ( | ||
| 329 | + err && | ||
| 330 | + err.code === 'ERR_STREAM_PREMATURE_CLOSE' && | ||
| 331 | + (rState && rState.ended && !rState.errored && !rState.errorEmitted) | ||
| 332 | + ) { | ||
| 333 | + // Some readable streams will emit 'close' before 'end'. However, since | ||
| 334 | + // this is on the readable side 'end' should still be emitted if the | ||
| 335 | + // stream has been ended and no error emitted. This should be allowed in | ||
| 336 | + // favor of backwards compatibility. Since the stream is piped to a | ||
| 337 | + // destination this should not result in any observable difference. | ||
| 338 | + // We don't need to check if this is a writable premature close since | ||
| 339 | + // eos will only fail with premature close on the reading side for | ||
| 340 | + // duplex streams. | ||
| 341 | + src | ||
| 342 | + .once('end', finish) | ||
| 343 | + .once('error', finish); | ||
| 344 | + } else { | ||
| 345 | + finish(err); | ||
| 346 | + } | ||
| 347 | + }); | ||
| 348 | + eos(dst, { readable: false, writable: true }, finish); | ||
| 349 | + } | ||
| 350 | + | ||
| 342 | 351 | module.exports = { pipelineImpl, pipeline }; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1027,7 +1027,7 @@ const tsp = require('timers/promises'); | |||
| 1027 | 1027 | const src = new PassThrough(); | |
| 1028 | 1028 | const dst = new PassThrough(); | |
| 1029 | 1029 | pipeline(src, dst, common.mustSucceed(() => { | |
| 1030 | - assert.strictEqual(dst.destroyed, true); | ||
| 1030 | + assert.strictEqual(dst.destroyed, false); | ||
| 1031 | 1031 | })); | |
| 1032 | 1032 | src.end(); | |
| 1033 | 1033 | } | |
@@ -1462,7 +1462,7 @@ const tsp = require('timers/promises'); | |||
| 1462 | 1462 | ||
| 1463 | 1463 | await pipelinePromise(read, duplex); | |
| 1464 | 1464 | ||
| 1465 | - assert.strictEqual(duplex.destroyed, true); | ||
| 1465 | + assert.strictEqual(duplex.destroyed, false); | ||
| 1466 | 1466 | } | |
| 1467 | 1467 | ||
| 1468 | 1468 | run().then(common.mustCall()); | |
@@ -1488,3 +1488,27 @@ const tsp = require('timers/promises'); | |||
| 1488 | 1488 | ||
| 1489 | 1489 | run().then(common.mustCall()); | |
| 1490 | 1490 | } | |
| 1491 | + | ||
| 1492 | + { | ||
| 1493 | + const s = new PassThrough({ objectMode: true }); | ||
| 1494 | + pipeline(async function*() { | ||
| 1495 | + await Promise.resolve(); | ||
| 1496 | + yield 'hello'; | ||
| 1497 | + yield 'world'; | ||
| 1498 | + yield 'world'; | ||
| 1499 | + }, s, async function(source) { | ||
| 1500 | + let ret = ''; | ||
| 1501 | + let n = 0; | ||
| 1502 | + for await (const chunk of source) { | ||
| 1503 | + if (n++ > 1) { | ||
| 1504 | + break; | ||
| 1505 | + } | ||
| 1506 | + ret += chunk; | ||
| 1507 | + } | ||
| 1508 | + return ret; | ||
| 1509 | + }, common.mustCall((err, val) => { | ||
| 1510 | + assert.strictEqual(err, undefined); | ||
| 1511 | + assert.strictEqual(val, 'helloworld'); | ||
| 1512 | + assert.strictEqual(s.destroyed, true); | ||
| 1513 | + })); | ||
| 1514 | + } | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments