Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
index.js284 linesDownload Raw Back to p-map
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 
codekingpro/portable-devtools · Team Ai