-
Notifications
You must be signed in to change notification settings - Fork 8k
/
async.js
107 lines (93 loc) · 2.66 KB
/
async.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
/*
* Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one
* or more contributor license agreements. Licensed under the Elastic License
* 2.0 and the Server Side Public License, v 1; you may not use this file except
* in compliance with, at your election, the Elastic License 2.0 or the Server
* Side Public License, v 1.
*/
/**
*
* @template T
* @template T2
* @param {(v: T) => Promise<T2>} fn
* @param {T} item
* @returns {Promise<PromiseSettledResult<T2>>}
*/
const settle = async (fn, item) => {
const [result] = await Promise.allSettled([(async () => fn(item))()]);
return result;
};
/**
* @template T
* @template T2
* @param {Array<T>} source
* @param {number} limit
* @param {(v: T) => Promise<T2>} mapFn
* @returns {Promise<T2[]>}
*/
function asyncMapWithLimit(source, limit, mapFn) {
return new Promise((resolve, reject) => {
if (limit < 1) {
reject(new Error('invalid limit, must be greater than 0'));
return;
}
let failed = false;
let inProgress = 0;
const queue = [...source.entries()];
/** @type {T2[]} */
const results = new Array(source.length);
/**
* this is run for each item, manages the inProgress state,
* calls the mapFn with that item, writes the map result to
* the result array, and calls runMore() after each item
* completes to either start another item or resolve the
* returned promise.
*
* @param {number} index
* @param {T} item
*/
function run(index, item) {
inProgress += 1;
settle(mapFn, item).then((result) => {
inProgress -= 1;
if (failed) {
return;
}
if (result.status === 'fulfilled') {
results[index] = result.value;
runMore();
return;
}
// when an error occurs we update the state to prevent
// holding onto old results and ignore future results
// from in-progress promises
failed = true;
results.length = 0;
reject(result.reason);
});
}
/**
* If there is work in the queue, schedule it, if there isn't
* any work to be scheduled and there isn't anything in progress
* then we're done. This function is called every time a mapFn
* promise resolves and once after initialization
*/
function runMore() {
if (!queue.length) {
if (inProgress === 0) {
resolve(results);
}
return;
}
while (inProgress < limit) {
const entry = queue.shift();
if (!entry) {
break;
}
run(...entry);
}
}
runMore();
});
}
module.exports = { asyncMapWithLimit };