(t *testing.T)
| 78 | } |
| 79 | |
| 80 | func TestConcurrencySafety(t *testing.T) { |
| 81 | q := queue.NewTaskQueue[int]() |
| 82 | var wg sync.WaitGroup |
| 83 | n := 1000 |
| 84 | // producers |
| 85 | wg.Go(func() { |
| 86 | for i := range n { |
| 87 | q.Add(newTask(fmt.Sprintf("p%d", i))) |
| 88 | } |
| 89 | }) |
| 90 | // consumers |
| 91 | wg.Go(func() { |
| 92 | count := 0 |
| 93 | for count < n { |
| 94 | _, err := q.Get() |
| 95 | if err != nil { |
| 96 | continue |
| 97 | } |
| 98 | count++ |
| 99 | } |
| 100 | }) |
| 101 | wg.Wait() |
| 102 | } |