Diff
1diff --git a/internal/sync/workqueue.go b/internal/sync/workqueue.go
2index 16eb442a3c5588ca7be41cfd47281c1f90109629..1c07e74b125746ad8491e3d96766d96de1587bb4 100644
3--- a/internal/sync/workqueue.go
4+++ b/internal/sync/workqueue.go
5@@ -10,8 +10,7 @@ import (
6 // WorkPool is a pool of work to be done.
7 type WorkPool struct {
8 workers int
9- work map[string]func()
10- mu sync.RWMutex
11+ work sync.Map
12 sem *semaphore.Weighted
13 ctx context.Context
14 logger func(string, ...interface{})
15@@ -33,7 +32,6 @@ func WithWorkPoolLogger(logger func(string, ...interface{})) WorkPoolOption {
16 func NewWorkPool(ctx context.Context, workers int, opts ...WorkPoolOption) *WorkPool {
17 wq := &WorkPool{
18 workers: workers,
19- work: make(map[string]func()),
20 ctx: ctx,
21 }
22
23@@ -52,20 +50,22 @@ func NewWorkPool(ctx context.Context, workers int, opts ...WorkPoolOption) *Work
24
25 // Run starts the workers and waits for them to finish.
26 func (wq *WorkPool) Run() {
27- for id, fn := range wq.work {
28+ wq.work.Range(func(key, value any) bool {
29+ id := key.(string)
30+ fn := value.(func())
31 if err := wq.sem.Acquire(wq.ctx, 1); err != nil {
32 wq.logf("workpool: %v", err)
33- return
34+ return false
35 }
36
37 go func(id string, fn func()) {
38 defer wq.sem.Release(1)
39 fn()
40- wq.mu.Lock()
41- delete(wq.work, id)
42- wq.mu.Unlock()
43+ wq.work.Delete(id)
44 }(id, fn)
45- }
46+
47+ return true
48+ })
49
50 if err := wq.sem.Acquire(wq.ctx, int64(wq.workers)); err != nil {
51 wq.logf("workpool: %v", err)
52@@ -75,19 +75,15 @@ func (wq *WorkPool) Run() {
53 // Add adds a new job to the pool.
54 // If the job already exists, it is a no-op.
55 func (wq *WorkPool) Add(id string, fn func()) {
56- wq.mu.Lock()
57- defer wq.mu.Unlock()
58- if _, ok := wq.work[id]; ok {
59+ if _, ok := wq.work.Load(id); ok {
60 return
61 }
62- wq.work[id] = fn
63+ wq.work.Store(id, fn)
64 }
65
66 // Status checks if a job is in the queue.
67 func (wq *WorkPool) Status(id string) bool {
68- wq.mu.RLock()
69- defer wq.mu.RUnlock()
70- _, ok := wq.work[id]
71+ _, ok := wq.work.Load(id)
72 return ok
73 }
74