/var/www/flurhymez/node_modules/webpack/lib/util
NameSizeModeActions
hash/-0755rm
ArrayHelpers.js2800644editdlrm
ArrayQueue.js21980644editdlrm
AsyncQueue.js96310644editdlrm
binarySearchBounds.js19070644editdlrm
cleverMerge.js165340644editdlrm
comparators.js124460644editdlrm
compileBooleanMatcher.js58220644editdlrm
create-schema-validation.js4770644editdlrm
createHash.js48190644editdlrm
deprecation.js64540644editdlrm
deterministicGrouping.js137640644editdlrm
extractUrlAndGlobal.js4160644editdlrm
findGraphRoots.js61110644editdlrm
fs.js107760644editdlrm
Hash.js9250644editdlrm
identifier.js103350644editdlrm
internalSerializables.js99730644editdlrm
IterableHelpers.js9620644editdlrm
LazyBucketSortedSet.js57320644editdlrm
LazySet.js45690644editdlrm
makeSerializable.js6400644editdlrm
MapHelpers.js4720644editdlrm
memoize.js6040644editdlrm
numberHash.js10600644editdlrm
objectToMap.js3460644editdlrm
ParallelismFactorCalculator.js15280644editdlrm
processAsyncTree.js14830644editdlrm
propertyAccess.js11900644editdlrm
Queue.js10480644editdlrm
registerExternalSerializer.js79130644editdlrm
runtime.js146050644editdlrm
Semaphore.js10080644editdlrm
semver.js155460644editdlrm
serialization.js40180644editdlrm
SetHelpers.js23160644editdlrm
smartGrouping.js52710644editdlrm
SortableSet.js36350644editdlrm
source.js17590644editdlrm
StackedCacheMap.js22840644editdlrm
StackedMap.js34530644editdlrm
StringXor.js11200644editdlrm
TupleQueue.js13170644editdlrm
TupleSet.js29090644editdlrm
URLAbsoluteSpecifier.js25530644editdlrm
WeakTupleMap.js34430644editdlrm
Edit: /var/www/flurhymez/node_modules/webpack/lib/util/AsyncQueue.js (9631B)
/* MIT License http://www.opensource.org/licenses/mit-license.php Author Tobias Koppers @sokra */ "use strict"; const { SyncHook, AsyncSeriesHook } = require("tapable"); const { makeWebpackError } = require("../HookWebpackError"); const WebpackError = require("../WebpackError"); const ArrayQueue = require("./ArrayQueue"); const QUEUED_STATE = 0; const PROCESSING_STATE = 1; const DONE_STATE = 2; let inHandleResult = 0; /** * @template T * @callback Callback * @param {WebpackError=} err * @param {T=} result */ /** * @template T * @template K * @template R */ class AsyncQueueEntry { /** * @param {T} item the item * @param {Callback} callback the callback */ constructor(item, callback) { this.item = item; /** @type {typeof QUEUED_STATE | typeof PROCESSING_STATE | typeof DONE_STATE} */ this.state = QUEUED_STATE; this.callback = callback; /** @type {Callback[] | undefined} */ this.callbacks = undefined; this.result = undefined; /** @type {WebpackError | undefined} */ this.error = undefined; } } /** * @template T * @template K * @template R */ class AsyncQueue { /** * @param {Object} options options object * @param {string=} options.name name of the queue * @param {number=} options.parallelism how many items should be processed at once * @param {AsyncQueue=} options.parent parent queue, which will have priority over this queue and with shared parallelism * @param {function(T): K=} options.getKey extract key from item * @param {function(T, Callback): void} options.processor async function to process items */ constructor({ name, parallelism, parent, processor, getKey }) { this._name = name; this._parallelism = parallelism || 1; this._processor = processor; this._getKey = getKey || /** @type {(T) => K} */ (item => /** @type {any} */ (item)); /** @type {Map>} */ this._entries = new Map(); /** @type {ArrayQueue>} */ this._queued = new ArrayQueue(); /** @type {AsyncQueue[]} */ this._children = undefined; this._activeTasks = 0; this._willEnsureProcessing = false; this._needProcessing = false; this._stopped = false; this._root = parent ? parent._root : this; if (parent) { if (this._root._children === undefined) { this._root._children = [this]; } else { this._root._children.push(this); } } this.hooks = { /** @type {AsyncSeriesHook<[T]>} */ beforeAdd: new AsyncSeriesHook(["item"]), /** @type {SyncHook<[T]>} */ added: new SyncHook(["item"]), /** @type {AsyncSeriesHook<[T]>} */ beforeStart: new AsyncSeriesHook(["item"]), /** @type {SyncHook<[T]>} */ started: new SyncHook(["item"]), /** @type {SyncHook<[T, Error, R]>} */ result: new SyncHook(["item", "error", "result"]) }; this._ensureProcessing = this._ensureProcessing.bind(this); } /** * @param {T} item an item * @param {Callback} callback callback function * @returns {void} */ add(item, callback) { if (this._stopped) return callback(new WebpackError("Queue was stopped")); this.hooks.beforeAdd.callAsync(item, err => { if (err) { callback( makeWebpackError(err, `AsyncQueue(${this._name}).hooks.beforeAdd`) ); return; } const key = this._getKey(item); const entry = this._entries.get(key); if (entry !== undefined) { if (entry.state === DONE_STATE) { if (inHandleResult++ > 3) { process.nextTick(() => callback(entry.error, entry.result)); } else { callback(entry.error, entry.result); } inHandleResult--; } else if (entry.callbacks === undefined) { entry.callbacks = [callback]; } else { entry.callbacks.push(callback); } return; } const newEntry = new AsyncQueueEntry(item, callback); if (this._stopped) { this.hooks.added.call(item); this._root._activeTasks++; process.nextTick(() => this._handleResult(newEntry, new WebpackError("Queue was stopped")) ); } else { this._entries.set(key, newEntry); this._queued.enqueue(newEntry); const root = this._root; root._needProcessing = true; if (root._willEnsureProcessing === false) { root._willEnsureProcessing = true; setImmediate(root._ensureProcessing); } this.hooks.added.call(item); } }); } /** * @param {T} item an item * @returns {void} */ invalidate(item) { const key = this._getKey(item); const entry = this._entries.get(key); this._entries.delete(key); if (entry.state === QUEUED_STATE) { this._queued.delete(entry); } } /** * Waits for an already started item * @param {T} item an item * @param {Callback} callback callback function * @returns {void} */ waitFor(item, callback) { const key = this._getKey(item); const entry = this._entries.get(key); if (entry === undefined) { return callback( new WebpackError( "waitFor can only be called for an already started item" ) ); } if (entry.state === DONE_STATE) { process.nextTick(() => callback(entry.error, entry.result)); } else if (entry.callbacks === undefined) { entry.callbacks = [callback]; } else { entry.callbacks.push(callback); } } /** * @returns {void} */ stop() { this._stopped = true; const queue = this._queued; this._queued = new ArrayQueue(); const root = this._root; for (const entry of queue) { this._entries.delete(this._getKey(entry.item)); root._activeTasks++; this._handleResult(entry, new WebpackError("Queue was stopped")); } } /** * @returns {void} */ increaseParallelism() { const root = this._root; root._parallelism++; /* istanbul ignore next */ if (root._willEnsureProcessing === false && root._needProcessing) { root._willEnsureProcessing = true; setImmediate(root._ensureProcessing); } } /** * @returns {void} */ decreaseParallelism() { const root = this._root; root._parallelism--; } /** * @param {T} item an item * @returns {boolean} true, if the item is currently being processed */ isProcessing(item) { const key = this._getKey(item); const entry = this._entries.get(key); return entry !== undefined && entry.state === PROCESSING_STATE; } /** * @param {T} item an item * @returns {boolean} true, if the item is currently queued */ isQueued(item) { const key = this._getKey(item); const entry = this._entries.get(key); return entry !== undefined && entry.state === QUEUED_STATE; } /** * @param {T} item an item * @returns {boolean} true, if the item is currently queued */ isDone(item) { const key = this._getKey(item); const entry = this._entries.get(key); return entry !== undefined && entry.state === DONE_STATE; } /** * @returns {void} */ _ensureProcessing() { while (this._activeTasks < this._parallelism) { const entry = this._queued.dequeue(); if (entry === undefined) break; this._activeTasks++; entry.state = PROCESSING_STATE; this._startProcessing(entry); } this._willEnsureProcessing = false; if (this._queued.length > 0) return; if (this._children !== undefined) { for (const child of this._children) { while (this._activeTasks < this._parallelism) { const entry = child._queued.dequeue(); if (entry === undefined) break; this._activeTasks++; entry.state = PROCESSING_STATE; child._startProcessing(entry); } if (child._queued.length > 0) return; } } if (!this._willEnsureProcessing) this._needProcessing = false; } /** * @param {AsyncQueueEntry} entry the entry * @returns {void} */ _startProcessing(entry) { this.hooks.beforeStart.callAsync(entry.item, err => { if (err) { this._handleResult( entry, makeWebpackError(err, `AsyncQueue(${this._name}).hooks.beforeStart`) ); return; } let inCallback = false; try { this._processor(entry.item, (e, r) => { inCallback = true; this._handleResult(entry, e, r); }); } catch (err) { if (inCallback) throw err; this._handleResult(entry, err, null); } this.hooks.started.call(entry.item); }); } /** * @param {AsyncQueueEntry} entry the entry * @param {WebpackError=} err error, if any * @param {R=} result result, if any * @returns {void} */ _handleResult(entry, err, result) { this.hooks.result.callAsync(entry.item, err, result, hookError => { const error = hookError ? makeWebpackError(hookError, `AsyncQueue(${this._name}).hooks.result`) : err; const callback = entry.callback; const callbacks = entry.callbacks; entry.state = DONE_STATE; entry.callback = undefined; entry.callbacks = undefined; entry.result = result; entry.error = error; const root = this._root; root._activeTasks--; if (root._willEnsureProcessing === false && root._needProcessing) { root._willEnsureProcessing = true; setImmediate(root._ensureProcessing); } if (inHandleResult++ > 3) { process.nextTick(() => { callback(error, result); if (callbacks !== undefined) { for (const callback of callbacks) { callback(error, result); } } }); } else { callback(error, result); if (callbacks !== undefined) { for (const callback of callbacks) { callback(error, result); } } } inHandleResult--; }); } clear() { this._entries.clear(); this._queued.clear(); this._activeTasks = 0; this._willEnsureProcessing = false; this._needProcessing = false; this._stopped = false; } } module.exports = AsyncQueue;