( items: readonly Input[], process: (item: Input, index: number) => Promise<Output>, shouldStop: () => boolean = () => false, concurrency = processingConcurrency(), )
| 24 | }; |
| 25 | |
| 26 | export const processWithConcurrency = async <Input, Output>( |
| 27 | items: readonly Input[], |
| 28 | process: (item: Input, index: number) => Promise<Output>, |
| 29 | shouldStop: () => boolean = () => false, |
| 30 | concurrency = processingConcurrency(), |
| 31 | ) => { |
| 32 | const results: Array<{ value: Output } | undefined> = []; |
| 33 | let processedCount = 0; |
| 34 | let failure: { error: unknown } | undefined; |
| 35 | const entries = items.entries(); |
| 36 | |
| 37 | const worker = async () => { |
| 38 | while (!failure && !shouldStop()) { |
| 39 | const next = entries.next(); |
| 40 | if (next.done) return; |
| 41 | const [index, item] = next.value; |
| 42 | |
| 43 | try { |
| 44 | results[index] = { value: await process(item, index) }; |
| 45 | } catch (error) { |
| 46 | failure ??= { error }; |
| 47 | return; |
| 48 | } |
| 49 | processedCount++; |
| 50 | await yieldToMainThread(); |
| 51 | } |
| 52 | }; |
| 53 | |
| 54 | await Promise.all( |
| 55 | Array.from( |
| 56 | { length: Math.min(items.length, Math.max(1, Math.floor(concurrency))) }, |
| 57 | worker, |
| 58 | ), |
| 59 | ); |
| 60 | |
| 61 | if (failure) throw failure.error; |
| 62 | |
| 63 | return { |
| 64 | results: results.flatMap((result) => (result ? [result.value] : [])), |
| 65 | processedCount, |
| 66 | stopped: processedCount < items.length, |
| 67 | }; |
| 68 | }; |
no test coverage detected