package utils
import (
"fmt"
"sync"
"testing"
"go.uber.org/atomic"
)
func TestParallelExecConcurrent(t *testing.T) {
const n = 137
var wg sync.WaitGroup
for g := 0; g < 8; g++ {
wg.Add(1)
go func() {
defer wg.Done()
for iter := 0; iter < 300; iter++ {
vals := make([]int, n)
counts := make([]atomic.Uint32, n)
for i := range vals {
vals[i] = i
}
ParallelExec(vals, 1, 2, func(i int) { counts[i].Add(1) })
for i := 0; i < n; i++ {
if got := counts[i].Load(); got != 1 {
t.Errorf("index %d visited %d times", i, got)
return
}
}
}
}()
}
wg.Wait()
}
var benchSink atomic.Uint64
func perItemWork() {
var acc uint64
for i := 0; i < 250; i++ {
acc = acc*1099511628211 + uint64(i)
}
benchSink.Add(acc)
}
func benchmarkParallelExec(b *testing.B, fanout int, concurrent bool) {
vals := make([]int, fanout)
fn := func(i int) { perItemWork() }
b.ResetTimer()
b.ReportAllocs()
if concurrent {
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
ParallelExec(vals, 1, 2, fn)
}
})
return
}
for n := 0; n < b.N; n++ {
ParallelExec(vals, 1, 2, fn)
}
}
func BenchmarkParallelExec(b *testing.B) {
for _, fanout := range []int{50, 200} {
b.Run(fmt.Sprintf("single/f%d", fanout), func(b *testing.B) {
benchmarkParallelExec(b, fanout, false)
})
b.Run(fmt.Sprintf("concurrent/f%d", fanout), func(b *testing.B) {
benchmarkParallelExec(b, fanout, true)
})
}
}