From 3f0a6a64e3a9d34fa3fe41aac1f8950632083dc2 Mon Sep 17 00:00:00 2001 From: Ray Liu Date: Sun, 20 Sep 2026 19:59:49 -0400 Subject: [PATCH] [6/9][test][worker] test concurrent worker lifecycle --- test/integration/run.sh | 17 +++ test/integration/worker_concurrency_test.go | 121 ++++++++++++++++++++ 2 files changed, 138 insertions(+) create mode 100644 test/integration/worker_concurrency_test.go diff --git a/test/integration/run.sh b/test/integration/run.sh index 78523b0..e96f999 100644 --- a/test/integration/run.sh +++ b/test/integration/run.sh @@ -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 \ @@ -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 \ diff --git a/test/integration/worker_concurrency_test.go b/test/integration/worker_concurrency_test.go new file mode 100644 index 0000000..96bac88 --- /dev/null +++ b/test/integration/worker_concurrency_test.go @@ -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) + } +}