a9e4824691f12a7c1e03161e9d60425a2d5ca637

Author
Ayman Bagabas <ayman.bagabas@gmail.com>
Committer
Ayman Bagabas <ayman.bagabas@gmail.com>
Date

Message

fix(server): possible panic using maps

use a sync.Map to avoid possible panic on concurrent map iteration and
map write

Fixes: b07de747b54a ("fix(log): respect log config (#280)")

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