codekingpro/portable-devtools
115k
1export default async function pMap(
2 iterable,
3 mapper,
4 {
5 concurrency = Number.POSITIVE_INFINITY,
6 stopOnError = true,
7 signal,
8 } = {},
9) {
10 return new Promise((resolve_, reject_) => {
11 if (iterable[Symbol.iterator] === undefined && iterable[Symbol.asyncIterator] === undefined) {
12 throw new TypeError(`Expected \`input\` to be either an \`Iterable\` or \`AsyncIterable\`, got (${typeof iterable})`);
13 }
14
15 if (typeof mapper !== 'function') {
16 throw new TypeError('Mapper function is required');
17 }
18
19 if (!((Number.isSafeInteger(concurrency) && concurrency >= 1) || concurrency === Number.POSITIVE_INFINITY)) {
20 throw new TypeError(`Expected \`concurrency\` to be an integer from 1 and up or \`Infinity\`, got \`${concurrency}\` (${typeof concurrency})`);
21 }
22
23 const result = [];
24 const errors = [];
25 const skippedIndexesMap = new Map();
26 let isRejected = false;
27 let isResolved = false;
28 let isIterableDone = false;
29 let resolvingCount = 0;
30 let currentIndex = 0;
31 const iterator = iterable[Symbol.iterator] === undefined ? iterable[Symbol.asyncIterator]() : iterable[Symbol.iterator]();
32
33 const signalListener = () => {
34 reject(signal.reason);
35 };
36
37 const cleanup = () => {
38 signal?.removeEventListener('abort', signalListener);
39 };
40
41 const resolve = value => {
42 resolve_(value);
43 cleanup();
44 };
45
46 const reject = reason => {
47 isRejected = true;
48 isResolved = true;
49 reject_(reason);
50 cleanup();
51 };
52
53 if (signal) {
54 if (signal.aborted) {
55 reject(signal.reason);
56 }
57
58 signal.addEventListener('abort', signalListener, {once: true});
59 }
60
61 const next = async () => {
62 if (isResolved) {
63 return;
64 }
65
66 const nextItem = await iterator.next();
67
68 const index = currentIndex;
69 currentIndex++;
70
71 // Note: `iterator.next()` can be called many times in parallel.
72 // This can cause multiple calls to this `next()` function to
73 // receive a `nextItem` with `done === true`.
74 // The shutdown logic that rejects/resolves must be protected
75 // so it runs only one time as the `skippedIndex` logic is
76 // non-idempotent.
77 if (nextItem.done) {
78 isIterableDone = true;
79
80 if (resolvingCount === 0 && !isResolved) {
81 if (!stopOnError && errors.length > 0) {
82 reject(new AggregateError(errors)); // eslint-disable-line unicorn/error-message
83 return;
84 }
85
86 isResolved = true;
87
88 if (skippedIndexesMap.size === 0) {
89 resolve(result);
90 return;
91 }
92
93 const pureResult = [];
94
95 // Support multiple `pMapSkip`'s.
96 for (const [index, value] of result.entries()) {
97 if (skippedIndexesMap.get(index) === pMapSkip) {
98 continue;
99 }
100
101 pureResult.push(value);
102 }
103
104 resolve(pureResult);
105 }
106
107 return;
108 }
109
110 resolvingCount++;
111
112 // Intentionally detached
113 (async () => {
114 try {
115 const element = await nextItem.value;
116
117 if (isResolved) {
118 return;
119 }
120
121 const value = await mapper(element, index);
122
123 // Use Map to stage the index of the element.
124 if (value === pMapSkip) {
125 skippedIndexesMap.set(index, value);
126 }
127
128 result[index] = value;
129
130 resolvingCount--;
131 await next();
132 } catch (error) {
133 if (stopOnError) {
134 reject(error);
135 } else {
136 errors.push(error);
137 resolvingCount--;
138
139 // In that case we can't really continue regardless of `stopOnError` state
140 // since an iterable is likely to continue throwing after it throws once.
141 // If we continue calling `next()` indefinitely we will likely end up
142 // in an infinite loop of failed iteration.
143 try {
144 await next();
145 } catch (error) {
146 reject(error);
147 }
148 }
149 }
150 })();
151 };
152
153 // Create the concurrent runners in a detached (non-awaited)
154 // promise. We need this so we can await the `next()` calls
155 // to stop creating runners before hitting the concurrency limit
156 // if the iterable has already been marked as done.
157 // NOTE: We *must* do this for async iterators otherwise we'll spin up
158 // infinite `next()` calls by default and never start the event loop.
159 (async () => {
160 for (let index = 0; index < concurrency; index++) {
161 try {
162 // eslint-disable-next-line no-await-in-loop
163 await next();
164 } catch (error) {
165 reject(error);
166 break;
167 }
168
169 if (isIterableDone || isRejected) {
170 break;
171 }
172 }
173 })();
174 });
175}
176
177export function pMapIterable(
178 iterable,
179 mapper,
180 {
181 concurrency = Number.POSITIVE_INFINITY,
182 backpressure = concurrency,
183 } = {},
184) {
185 if (iterable[Symbol.iterator] === undefined && iterable[Symbol.asyncIterator] === undefined) {
186 throw new TypeError(`Expected \`input\` to be either an \`Iterable\` or \`AsyncIterable\`, got (${typeof iterable})`);
187 }
188
189 if (typeof mapper !== 'function') {
190 throw new TypeError('Mapper function is required');
191 }
192
193 if (!((Number.isSafeInteger(concurrency) && concurrency >= 1) || concurrency === Number.POSITIVE_INFINITY)) {
194 throw new TypeError(`Expected \`concurrency\` to be an integer from 1 and up or \`Infinity\`, got \`${concurrency}\` (${typeof concurrency})`);
195 }
196
197 if (!((Number.isSafeInteger(backpressure) && backpressure >= concurrency) || backpressure === Number.POSITIVE_INFINITY)) {
198 throw new TypeError(`Expected \`backpressure\` to be an integer from \`concurrency\` (${concurrency}) and up or \`Infinity\`, got \`${backpressure}\` (${typeof backpressure})`);
199 }
200
201 return {
202 async * [Symbol.asyncIterator]() {
203 const iterator = iterable[Symbol.asyncIterator] === undefined ? iterable[Symbol.iterator]() : iterable[Symbol.asyncIterator]();
204
205 const promises = [];
206 let pendingPromisesCount = 0;
207 let isDone = false;
208 let index = 0;
209
210 function trySpawn() {
211 if (isDone || !(pendingPromisesCount < concurrency && promises.length < backpressure)) {
212 return;
213 }
214
215 pendingPromisesCount++;
216
217 const promise = (async () => {
218 const {done, value} = await iterator.next();
219
220 if (done) {
221 pendingPromisesCount--;
222 return {done: true};
223 }
224
225 // Spawn if still below concurrency and backpressure limit
226 trySpawn();
227
228 try {
229 const returnValue = await mapper(await value, index++);
230
231 pendingPromisesCount--;
232
233 if (returnValue === pMapSkip) {
234 const index = promises.indexOf(promise);
235
236 if (index > 0) {
237 promises.splice(index, 1);
238 }
239 }
240
241 // Spawn if still below backpressure limit and just dropped below concurrency limit
242 trySpawn();
243
244 return {done: false, value: returnValue};
245 } catch (error) {
246 pendingPromisesCount--;
247 isDone = true;
248 return {error};
249 }
250 })();
251
252 promises.push(promise);
253 }
254
255 trySpawn();
256
257 while (promises.length > 0) {
258 const {error, done, value} = await promises[0]; // eslint-disable-line no-await-in-loop
259
260 promises.shift();
261
262 if (error) {
263 throw error;
264 }
265
266 if (done) {
267 return;
268 }
269
270 // Spawn if just dropped below backpressure limit and below the concurrency limit
271 trySpawn();
272
273 if (value === pMapSkip) {
274 continue;
275 }
276
277 yield value;
278 }
279 },
280 };
281}
282
283export const pMapSkip = Symbol('skip');
284 