queue.ts 2.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687
  1. type QueueInput = {
  2. paused: () => boolean
  3. bootstrap: () => Promise<void>
  4. bootstrapInstance: (directory: string) => Promise<void> | void
  5. key?: (directory: string) => string
  6. }
  7. export function createRefreshQueue(input: QueueInput) {
  8. const queued = new Map<string, string>()
  9. let root = false
  10. let running = false
  11. let timer: ReturnType<typeof setTimeout> | undefined
  12. const key = input.key ?? ((directory: string) => directory)
  13. const tick = () => new Promise<void>((resolve) => setTimeout(resolve, 0))
  14. const take = (count: number) => {
  15. if (queued.size === 0) return [] as string[]
  16. const items: string[] = []
  17. for (const [id, directory] of queued) {
  18. queued.delete(id)
  19. items.push(directory)
  20. if (items.length >= count) break
  21. }
  22. return items
  23. }
  24. const schedule = () => {
  25. if (timer) return
  26. timer = setTimeout(() => {
  27. timer = undefined
  28. void drain()
  29. }, 0)
  30. }
  31. const push = (directory: string) => {
  32. if (!directory) return
  33. queued.set(key(directory), directory)
  34. if (input.paused()) return
  35. schedule()
  36. }
  37. const refresh = () => {
  38. root = true
  39. if (input.paused()) return
  40. schedule()
  41. }
  42. async function drain() {
  43. if (running) return
  44. running = true
  45. try {
  46. while (true) {
  47. if (input.paused()) return
  48. if (root) {
  49. root = false
  50. await input.bootstrap()
  51. await tick()
  52. continue
  53. }
  54. const dirs = take(2)
  55. if (dirs.length === 0) return
  56. await Promise.all(dirs.map((dir) => input.bootstrapInstance(dir)))
  57. await tick()
  58. }
  59. } finally {
  60. running = false
  61. // oxlint-disable-next-line no-unsafe-finally -- intentional: early return skips schedule() when paused
  62. if (input.paused()) return
  63. if (root || queued.size) schedule()
  64. }
  65. }
  66. return {
  67. push,
  68. refresh,
  69. clear(directory: string) {
  70. queued.delete(key(directory))
  71. },
  72. dispose() {
  73. if (!timer) return
  74. clearTimeout(timer)
  75. timer = undefined
  76. },
  77. }
  78. }