async-pool.ts 958 B

12345678910111213141516171819202122232425262728293031323334
  1. interface MapWithConcurrencyOptions {
  2. signal?: AbortSignal
  3. onItemComplete?: (index: number) => void
  4. }
  5. /**
  6. * Map items with a fixed concurrency limit. Result order matches input order.
  7. */
  8. export async function mapWithConcurrency<T, R>(
  9. items: readonly T[],
  10. concurrency: number,
  11. fn: (item: T, index: number) => Promise<R>,
  12. options?: MapWithConcurrencyOptions,
  13. ): Promise<R[]> {
  14. if (items.length === 0) return []
  15. const results: R[] = new Array(items.length)
  16. let nextIndex = 0
  17. async function worker(): Promise<void> {
  18. while (true) {
  19. if (options?.signal?.aborted) return
  20. const index = nextIndex
  21. nextIndex += 1
  22. if (index >= items.length) return
  23. results[index] = await fn(items[index], index)
  24. options?.onItemComplete?.(index)
  25. }
  26. }
  27. const workerCount = Math.max(1, Math.min(concurrency, items.length))
  28. await Promise.all(Array.from({ length: workerCount }, () => worker()))
  29. return results
  30. }