| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1,5 +1,9 @@ | |||
| 1 | 1 | 'use strict'; | |
| 2 | 2 | ||
| 3 | + const { | ||
| 4 | + SymbolDispose, | ||
| 5 | + } = primordials; | ||
| 6 | + | ||
| 3 | 7 | const { | |
| 4 | 8 | AbortError, | |
| 5 | 9 | codes, | |
@@ -13,6 +17,7 @@ const { | |||
| 13 | 17 | ||
| 14 | 18 | const eos = require('internal/streams/end-of-stream'); | |
| 15 | 19 | const { ERR_INVALID_ARG_TYPE } = codes; | |
| 20 | + let addAbortListener; | ||
| 16 | 21 | ||
| 17 | 22 | // This method is inlined here for readable-stream | |
| 18 | 23 | // It also does not allow for signal to not exist on the stream | |
@@ -46,8 +51,9 @@ module.exports.addAbortSignalNoValidate = function(signal, stream) { | |||
| 46 | 51 | if (signal.aborted) { | |
| 47 | 52 | onAbort(); | |
| 48 | 53 | } else { | |
| 49 | - signal.addEventListener('abort', onAbort); | ||
| 50 | - eos(stream, () => signal.removeEventListener('abort', onAbort)); | ||
| 54 | + addAbortListener ??= require('events').addAbortListener; | ||
| 55 | + const disposable = addAbortListener(signal, onAbort); | ||
| 56 | + eos(stream, disposable[SymbolDispose]); | ||
| 51 | 57 | } | |
| 52 | 58 | return stream; | |
| 53 | 59 | }; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -22,7 +22,11 @@ const { | |||
| 22 | 22 | validateBoolean, | |
| 23 | 23 | } = require('internal/validators'); | |
| 24 | 24 | ||
| 25 | - const { Promise, PromisePrototypeThen } = primordials; | ||
| 25 | + const { | ||
| 26 | + Promise, | ||
| 27 | + PromisePrototypeThen, | ||
| 28 | + SymbolDispose, | ||
| 29 | + } = primordials; | ||
| 26 | 30 | ||
| 27 | 31 | const { | |
| 28 | 32 | isClosed, | |
@@ -40,6 +44,7 @@ const { | |||
| 40 | 44 | willEmitClose: _willEmitClose, | |
| 41 | 45 | kIsClosedPromise, | |
| 42 | 46 | } = require('internal/streams/utils'); | |
| 47 | + let addAbortListener; | ||
| 43 | 48 | ||
| 44 | 49 | function isRequest(stream) { | |
| 45 | 50 | return stream.setHeader && typeof stream.abort === 'function'; | |
@@ -249,12 +254,13 @@ function eos(stream, options, callback) { | |||
| 249 | 254 | if (options.signal.aborted) { | |
| 250 | 255 | process.nextTick(abort); | |
| 251 | 256 | } else { | |
| 257 | + addAbortListener ??= require('events').addAbortListener; | ||
| 258 | + const disposable = addAbortListener(options.signal, abort); | ||
| 252 | 259 | const originalCallback = callback; | |
| 253 | 260 | callback = once((...args) => { | |
| 254 | - options.signal.removeEventListener('abort', abort); | ||
| 261 | + disposable[SymbolDispose](); | ||
| 255 | 262 | originalCallback.apply(stream, args); | |
| 256 | 263 | }); | |
| 257 | - options.signal.addEventListener('abort', abort); | ||
| 258 | 264 | } | |
| 259 | 265 | } | |
| 260 | 266 | ||
@@ -272,12 +278,13 @@ function eosWeb(stream, options, callback) { | |||
| 272 | 278 | if (options.signal.aborted) { | |
| 273 | 279 | process.nextTick(abort); | |
| 274 | 280 | } else { | |
| 281 | + addAbortListener ??= require('events').addAbortListener; | ||
| 282 | + const disposable = addAbortListener(options.signal, abort); | ||
| 275 | 283 | const originalCallback = callback; | |
| 276 | 284 | callback = once((...args) => { | |
| 277 | - options.signal.removeEventListener('abort', abort); | ||
| 285 | + disposable[SymbolDispose](); | ||
| 278 | 286 | originalCallback.apply(stream, args); | |
| 279 | 287 | }); | |
| 280 | - options.signal.addEventListener('abort', abort); | ||
| 281 | 288 | } | |
| 282 | 289 | } | |
| 283 | 290 | const resolverFn = (...args) => { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1,6 +1,6 @@ | |||
| 1 | 1 | 'use strict'; | |
| 2 | 2 | ||
| 3 | - const { AbortController } = require('internal/abort_controller'); | ||
| 3 | + const { AbortController, AbortSignal } = require('internal/abort_controller'); | ||
| 4 | 4 | ||
| 5 | 5 | const { | |
| 6 | 6 | codes: { | |
@@ -16,7 +16,7 @@ const { | |||
| 16 | 16 | validateInteger, | |
| 17 | 17 | validateObject, | |
| 18 | 18 | } = require('internal/validators'); | |
| 19 | - const { kWeakHandler } = require('internal/event_target'); | ||
| 19 | + const { kWeakHandler, kResistStopPropagation } = require('internal/event_target'); | ||
| 20 | 20 | const { finished } = require('internal/streams/end-of-stream'); | |
| 21 | 21 | const staticCompose = require('internal/streams/compose'); | |
| 22 | 22 | const { | |
@@ -27,6 +27,7 @@ const { deprecate } = require('internal/util'); | |||
| 27 | 27 | ||
| 28 | 28 | const { | |
| 29 | 29 | ArrayPrototypePush, | |
| 30 | + Boolean, | ||
| 30 | 31 | MathFloor, | |
| 31 | 32 | Number, | |
| 32 | 33 | NumberIsNaN, | |
@@ -84,19 +85,11 @@ function map(fn, options) { | |||
| 84 | 85 | validateInteger(concurrency, 'concurrency', 1); | |
| 85 | 86 | ||
| 86 | 87 | return async function* map() { | |
| 87 | - const ac = new AbortController(); | ||
| 88 | + const signal = AbortSignal.any([options?.signal].filter(Boolean)); | ||
| 88 | 89 | const stream = this; | |
| 89 | 90 | const queue = []; | |
| 90 | - const signal = ac.signal; | ||
| 91 | 91 | const signalOpt = { signal }; | |
| 92 | 92 | ||
| 93 | - const abort = () => ac.abort(); | ||
| 94 | - if (options?.signal?.aborted) { | ||
| 95 | - abort(); | ||
| 96 | - } | ||
| 97 | - | ||
| 98 | - options?.signal?.addEventListener('abort', abort); | ||
| 99 | - | ||
| 100 | 93 | let next; | |
| 101 | 94 | let resume; | |
| 102 | 95 | let done = false; | |
@@ -153,7 +146,6 @@ function map(fn, options) { | |||
| 153 | 146 | next(); | |
| 154 | 147 | next = null; | |
| 155 | 148 | } | |
| 156 | - options?.signal?.removeEventListener('abort', abort); | ||
| 157 | 149 | } | |
| 158 | 150 | } | |
| 159 | 151 | ||
@@ -188,8 +180,6 @@ function map(fn, options) { | |||
| 188 | 180 | }); | |
| 189 | 181 | } | |
| 190 | 182 | } finally { | |
| 191 | - ac.abort(); | ||
| 192 | - | ||
| 193 | 183 | done = true; | |
| 194 | 184 | if (resume) { | |
| 195 | 185 | resume(); | |
@@ -301,7 +291,7 @@ async function reduce(reducer, initialValue, options) { | |||
| 301 | 291 | const ac = new AbortController(); | |
| 302 | 292 | const signal = ac.signal; | |
| 303 | 293 | if (options?.signal) { | |
| 304 | - const opts = { once: true, [kWeakHandler]: this }; | ||
| 294 | + const opts = { once: true, [kWeakHandler]: this, [kResistStopPropagation]: true }; | ||
| 305 | 295 | options.signal.addEventListener('abort', () => ac.abort(), opts); | |
| 306 | 296 | } | |
| 307 | 297 | let gotAnyItemFromStream = false; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -7,6 +7,7 @@ const { | |||
| 7 | 7 | ArrayIsArray, | |
| 8 | 8 | Promise, | |
| 9 | 9 | SymbolAsyncIterator, | |
| 10 | + SymbolDispose, | ||
| 10 | 11 | } = primordials; | |
| 11 | 12 | ||
| 12 | 13 | const eos = require('internal/streams/end-of-stream'); | |
@@ -44,6 +45,7 @@ const { AbortController } = require('internal/abort_controller'); | |||
| 44 | 45 | ||
| 45 | 46 | let PassThrough; | |
| 46 | 47 | let Readable; | |
| 48 | + let addAbortListener; | ||
| 47 | 49 | ||
| 48 | 50 | function destroyer(stream, reading, writing) { | |
| 49 | 51 | let finished = false; | |
@@ -206,7 +208,11 @@ function pipelineImpl(streams, callback, opts) { | |||
| 206 | 208 | finishImpl(new AbortError()); | |
| 207 | 209 | } | |
| 208 | 210 | ||
| 209 | - outerSignal?.addEventListener('abort', abort); | ||
| 211 | + addAbortListener ??= require('events').addAbortListener; | ||
| 212 | + let disposable; | ||
| 213 | + if (outerSignal) { | ||
| 214 | + disposable = addAbortListener(outerSignal, abort); | ||
| 215 | + } | ||
| 210 | 216 | ||
| 211 | 217 | let error; | |
| 212 | 218 | let value; | |
@@ -231,7 +237,7 @@ function pipelineImpl(streams, callback, opts) { | |||
| 231 | 237 | destroys.shift()(error); | |
| 232 | 238 | } | |
| 233 | 239 | ||
| 234 | - outerSignal?.removeEventListener('abort', abort); | ||
| 240 | + disposable?.[SymbolDispose](); | ||
| 235 | 241 | ac.abort(); | |
| 236 | 242 | ||
| 237 | 243 | if (final) { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -22,6 +22,7 @@ const { | |||
| 22 | 22 | SafePromiseAll, | |
| 23 | 23 | Symbol, | |
| 24 | 24 | SymbolAsyncIterator, | |
| 25 | + SymbolDispose, | ||
| 25 | 26 | SymbolToStringTag, | |
| 26 | 27 | Uint8Array, | |
| 27 | 28 | } = primordials; | |
@@ -140,6 +141,7 @@ const kRelease = Symbol('kRelease'); | |||
| 140 | 141 | ||
| 141 | 142 | let releasedError; | |
| 142 | 143 | let releasingError; | |
| 144 | + let addAbortListener; | ||
| 143 | 145 | ||
| 144 | 146 | const userModuleRegExp = /^ {4}at (?:[^/\\(]+ \()(?!node:(.+):\d+:\d+\)$).*/gm; | |
| 145 | 147 | ||
@@ -1259,6 +1261,7 @@ function readableStreamPipeTo( | |||
| 1259 | 1261 | ||
| 1260 | 1262 | let reader; | |
| 1261 | 1263 | let writer; | |
| 1264 | + let disposable; | ||
| 1262 | 1265 | // Both of these can throw synchronously. We want to capture | |
| 1263 | 1266 | // the error and return a rejected promise instead. | |
| 1264 | 1267 | try { | |
@@ -1291,7 +1294,7 @@ function readableStreamPipeTo( | |||
| 1291 | 1294 | writableStreamDefaultWriterRelease(writer); | |
| 1292 | 1295 | readableStreamReaderGenericRelease(reader); | |
| 1293 | 1296 | if (signal !== undefined) | |
| 1294 | - signal.removeEventListener('abort', abortAlgorithm); | ||
| 1297 | + disposable?.[SymbolDispose](); | ||
| 1295 | 1298 | if (rejected) | |
| 1296 | 1299 | promise.reject(error); | |
| 1297 | 1300 | else | |
@@ -1418,7 +1421,8 @@ function readableStreamPipeTo( | |||
| 1418 | 1421 | abortAlgorithm(); | |
| 1419 | 1422 | return promise.promise; | |
| 1420 | 1423 | } | |
| 1421 | - signal.addEventListener('abort', abortAlgorithm, { once: true }); | ||
| 1424 | + addAbortListener ??= require('events').addAbortListener; | ||
| 1425 | + disposable = addAbortListener(signal, abortAlgorithm); | ||
| 1422 | 1426 | } | |
| 1423 | 1427 | ||
| 1424 | 1428 | setPromiseHandled(run()); | |
| Back | FazBrowse Home | New Git URL |
0 commit comments