Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions test/integration/run.sh
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,15 @@ docker compose --file "$compose_file" exec -T kafka \
--partitions 1 \
--replication-factor 1

docker compose --file "$compose_file" exec -T kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:19092 \
--create \
--if-not-exists \
--topic concurrency-test-ready \
--partitions 1 \
--replication-factor 1

docker compose --file "$compose_file" exec -T kafka \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:19092 \
Expand Down Expand Up @@ -76,6 +85,14 @@ docker compose --file "$compose_file" exec -T kafka \
--entity-name kq.retry-test.workers \
--add-config share.auto.offset.reset=earliest

docker compose --file "$compose_file" exec -T kafka \
/opt/kafka/bin/kafka-configs.sh \
--bootstrap-server localhost:19092 \
--alter \
--entity-type groups \
--entity-name kq.concurrency-test.workers \
--add-config share.auto.offset.reset=earliest

docker compose --file "$compose_file" exec -T kafka \
/opt/kafka/bin/kafka-configs.sh \
--bootstrap-server localhost:19092 \
Expand Down
121 changes: 121 additions & 0 deletions test/integration/worker_concurrency_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
//go:build integration

package integration_test

import (
"context"
"fmt"
"os"
"sync/atomic"
"testing"
"time"

"github.com/raydatray/kq"
)

func TestWorkerConcurrency(t *testing.T) {
broker := os.Getenv("KQ_TEST_BROKERS")
if broker == "" {
t.Skip("KQ_TEST_BROKERS is not set")
}

const (
concurrency = 4
taskCount = 9
)
config := kq.Config{
Brokers: []string{broker},
Queue: "concurrency-test",
}
testCtx, cancelTest := context.WithTimeout(context.Background(), 20*time.Second)
defer cancelTest()

client, err := kq.NewClient(config)
if err != nil {
t.Fatal(err)
}
for i := range taskCount {
if _, err := client.Enqueue(testCtx, kq.Task{
Type: "concurrency",
Payload: []byte(fmt.Sprintf("task-%d", i)),
}); err != nil {
t.Fatal(err)
}
}
client.Close()

started := make(chan string, taskCount)
completed := make(chan string, taskCount)
release := make(chan struct{})
var active atomic.Int32
var maximum atomic.Int32
worker, err := kq.NewWorker(kq.WorkerConfig{
Config: config,
Concurrency: concurrency,
}, func(ctx context.Context, task kq.Task) error {
current := active.Add(1)
defer active.Add(-1)
for {
prior := maximum.Load()
if current <= prior || maximum.CompareAndSwap(prior, current) {
break
}
}

id := string(task.Payload)
started <- id
select {
case <-ctx.Done():
return ctx.Err()
case <-release:
completed <- id
return nil
}
})
if err != nil {
t.Fatal(err)
}
defer worker.Close()

workerCtx, cancelWorker := context.WithCancel(testCtx)
done := make(chan error, 1)
go func() { done <- worker.Run(workerCtx) }()

firstGeneration := make(map[string]bool, concurrency)
for range concurrency {
select {
case id := <-started:
if firstGeneration[id] {
t.Fatalf("task %q started twice", id)
}
firstGeneration[id] = true
case <-testCtx.Done():
t.Fatal("worker did not fill its concurrency pool")
}
}
if got := maximum.Load(); got != concurrency {
t.Fatalf("maximum concurrency = %d, want %d", got, concurrency)
}
close(release)

seen := make(map[string]bool, taskCount)
for len(seen) < taskCount {
select {
case id := <-completed:
if seen[id] {
t.Fatalf("task %q completed twice", id)
}
seen[id] = true
case <-testCtx.Done():
t.Fatalf("completed %d/%d tasks", len(seen), taskCount)
}
}

cancelWorker()
if err := <-done; err != nil {
t.Fatal(err)
}
if got := maximum.Load(); got > concurrency {
t.Fatalf("maximum concurrency = %d, want at most %d", got, concurrency)
}
}
Loading