From 9d15a6aa4c44a8f7cc06c6bf6b2caced8eb64636 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Paul=20Gro=C3=9Fmann?= Date: Tue, 11 Aug 2026 12:51:10 +0200 Subject: [PATCH 1/4] fix: data race by not using default http handler --- pkg/pubsub/publisher.go | 2 +- pkg/pubsub/subscriber.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/pubsub/publisher.go b/pkg/pubsub/publisher.go index 09a259c..f270592 100644 --- a/pkg/pubsub/publisher.go +++ b/pkg/pubsub/publisher.go @@ -25,7 +25,7 @@ type Publisher struct { // API dataplane client fails to initialize. func NewPublisher(topicID uuid.UUID, opts ...Option) *Publisher { cfg := &clientConfig{ - httpClient: http.DefaultClient, + httpClient: &http.Client{}, host: "pubsub.eu01.onstackit.cloud", logger: logr.FromSlogHandler(slog.Default().Handler()), } diff --git a/pkg/pubsub/subscriber.go b/pkg/pubsub/subscriber.go index e4aea9b..8ee3374 100644 --- a/pkg/pubsub/subscriber.go +++ b/pkg/pubsub/subscriber.go @@ -27,7 +27,7 @@ type Subscriber struct { // API dataplane client fails to initialize. func NewSubscriber(topicID uuid.UUID, subscriptionID uuid.UUID, opts ...Option) *Subscriber { cfg := &clientConfig{ - httpClient: http.DefaultClient, + httpClient: &http.Client{}, host: "pubsub.eu01.onstackit.cloud", logger: logr.FromSlogHandler(slog.Default().Handler()), } From 4e37b1b7bec4cb8d181b5fe5449ad318fbf35e57 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Paul=20Gro=C3=9Fmann?= Date: Tue, 11 Aug 2026 13:01:15 +0200 Subject: [PATCH 2/4] fix: corrected comment that functions do not return an error --- pkg/pubsub/publisher.go | 3 +-- pkg/pubsub/subscriber.go | 3 +-- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/pkg/pubsub/publisher.go b/pkg/pubsub/publisher.go index f270592..68d7b6d 100644 --- a/pkg/pubsub/publisher.go +++ b/pkg/pubsub/publisher.go @@ -21,8 +21,7 @@ type Publisher struct { httpClient *http.Client } -// NewPublisher instantiates a new Publisher. It returns an error if the underlying -// API dataplane client fails to initialize. +// NewPublisher instantiates a new Publisher. func NewPublisher(topicID uuid.UUID, opts ...Option) *Publisher { cfg := &clientConfig{ httpClient: &http.Client{}, diff --git a/pkg/pubsub/subscriber.go b/pkg/pubsub/subscriber.go index 8ee3374..c97653c 100644 --- a/pkg/pubsub/subscriber.go +++ b/pkg/pubsub/subscriber.go @@ -23,8 +23,7 @@ type Subscriber struct { wg sync.WaitGroup } -// NewSubscriber instantiates a new Subscriber. It returns an error if the underlying -// API dataplane client fails to initialize. +// NewSubscriber instantiates a new Subscriber. func NewSubscriber(topicID uuid.UUID, subscriptionID uuid.UUID, opts ...Option) *Subscriber { cfg := &clientConfig{ httpClient: &http.Client{}, From 21bb1d4ed0cc7bef02aacdb7b426ce9c52ed584b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Paul=20Gro=C3=9Fmann?= Date: Tue, 11 Aug 2026 13:48:18 +0200 Subject: [PATCH 3/4] refactor: aligned longPullDuration settings to other parameters --- pkg/pubsub/pulljob.go | 41 +++++++++++++++++++++++------------ pkg/pubsub/subscriber.go | 20 +++++++---------- pkg/pubsub/subscriber_test.go | 1 - 3 files changed, 35 insertions(+), 27 deletions(-) diff --git a/pkg/pubsub/pulljob.go b/pkg/pubsub/pulljob.go index 48bc52a..2a14dcb 100644 --- a/pkg/pubsub/pulljob.go +++ b/pkg/pubsub/pulljob.go @@ -3,13 +3,14 @@ package pubsub import ( "context" "errors" + "fmt" "time" ) type pullJob struct { subscription *Subscriber maxPullMessages int32 - longPullDuration *int32 + longPullDuration int32 interval time.Duration bufferSize int errHandler func(err error) bool @@ -27,7 +28,7 @@ func WithPullMaxMessages(maximum int32) PullJobOption { func WithPullLongPullDuration(milliseconds int32) PullJobOption { return func(b *pullJob) { - b.longPullDuration = &milliseconds + b.longPullDuration = milliseconds } } @@ -53,12 +54,13 @@ func WithErrorHandler(handler func(err error) bool) PullJobOption { } } -func newPullJob(s *Subscriber, opts []PullJobOption) *pullJob { +func newPullJob(s *Subscriber, opts []PullJobOption) (*pullJob, error) { b := &pullJob{ - subscription: s, - maxPullMessages: 10, - interval: time.Second * 5, - bufferSize: 0, + subscription: s, + maxPullMessages: 10, + interval: 1, + longPullDuration: 5000, + bufferSize: 0, errHandler: func(err error) bool { s.logger.Error(err, "fatal background error") return true @@ -69,7 +71,13 @@ func newPullJob(s *Subscriber, opts []PullJobOption) *pullJob { opt(b) } - return b + if b.interval < 1 { + return nil, &ConfigurationError{ + Msg: fmt.Sprintf("interval must be at least set to 1, got %d", b.interval), + } + } + + return b, nil } func (b *pullJob) runLoop(ctx context.Context, handler func(context.Context, PullMessages)) { @@ -81,10 +89,7 @@ func (b *pullJob) runLoop(ctx context.Context, handler func(context.Context, Pul case <-ctx.Done(): return case <-ticker.C: - pullOpts := []PullOption{WithMaxMessages(b.maxPullMessages)} - if b.longPullDuration != nil { - pullOpts = append(pullOpts, WithLongPullDuration(*b.longPullDuration)) - } + pullOpts := []PullOption{WithMaxMessages(b.maxPullMessages), WithLongPullDuration(b.longPullDuration)} messages, err := b.subscription.Pull(ctx, pullOpts...) if err != nil { //nolint:nestif var sdkErr SDKError // Declare the target variable @@ -121,7 +126,11 @@ func (s *Subscriber) PullJobCallback( return ErrMissingCallback } - job := newPullJob(s, opts) + job, err := newPullJob(s, opts) + if err != nil { + return err + } + s.wg.Go(func() { job.runLoop(ctx, callback) }) @@ -135,7 +144,11 @@ func (s *Subscriber) PullJobCallback( } func (s *Subscriber) PullJobChan(ctx context.Context, opts ...PullJobOption) (<-chan PullMessages, error) { - job := newPullJob(s, opts) + job, err := newPullJob(s, opts) + if err != nil { + return nil, err + } + out := make(chan PullMessages, job.bufferSize) s.wg.Go(func() { diff --git a/pkg/pubsub/subscriber.go b/pkg/pubsub/subscriber.go index c97653c..8c75654 100644 --- a/pkg/pubsub/subscriber.go +++ b/pkg/pubsub/subscriber.go @@ -111,7 +111,7 @@ func toSDKMessages(m []api.Message, subscription *Subscriber) PullMessages { type pullOptions struct { maxMessages int32 - longPullDuration *int32 + longPullDuration int32 } type PullOption func(*pullOptions) @@ -124,33 +124,29 @@ func WithMaxMessages(maximum int32) PullOption { func WithLongPullDuration(milliseconds int32) PullOption { return func(opts *pullOptions) { - opts.longPullDuration = &milliseconds + opts.longPullDuration = milliseconds } } func (s *Subscriber) Pull(ctx context.Context, opts ...PullOption) (PullMessages, error) { cfg := &pullOptions{ - maxMessages: 32, + maxMessages: 32, + longPullDuration: 100, } for _, opt := range opts { opt(cfg) } - var longPullDuration *int32 - if cfg.longPullDuration != nil && *cfg.longPullDuration != 0 { - ms := *cfg.longPullDuration - if ms < 100 || ms > 5000 { - return nil, &ConfigurationError{ - Msg: fmt.Sprintf("long_pull_duration must be 0 (default) or between 100–5000, got %d", ms), - } + if cfg.longPullDuration < 100 || cfg.longPullDuration > 5000 { + return nil, &ConfigurationError{ + Msg: fmt.Sprintf("long_pull_duration must be between 100–5000, got %d", cfg.longPullDuration), } - longPullDuration = &ms } reqBody := api.PullMessagesParams{ PubSubMaxMessages: &cfg.maxMessages, - PubSubLongPullDuration: longPullDuration, + PubSubLongPullDuration: &cfg.longPullDuration, } s.logger.V(4).Info("pulling messages", "max_messages", int(cfg.maxMessages)) diff --git a/pkg/pubsub/subscriber_test.go b/pkg/pubsub/subscriber_test.go index 4e08d14..422c56c 100644 --- a/pkg/pubsub/subscriber_test.go +++ b/pkg/pubsub/subscriber_test.go @@ -86,7 +86,6 @@ var _ = Describe("WithLongPullDuration validation", func() { Expect(errors.As(err, &cfgErr)).To(BeFalse(), "expected no ConfigurationError for ms=%d", ms) } }, - Entry("disabled (0)", int32(0)), Entry("minimum (100)", int32(100)), Entry("maximum (5000)", int32(5000)), ) From 798471009c2edf1d491805db13c3ee47b7264907 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Paul=20Gro=C3=9Fmann?= Date: Tue, 11 Aug 2026 15:17:44 +0200 Subject: [PATCH 4/4] fix: flakey pulljob tests --- pkg/pubsub/pulljob.go | 24 +++++++++++------------ pkg/pubsub/subscriber_test.go | 37 ++++++++--------------------------- 2 files changed, 20 insertions(+), 41 deletions(-) diff --git a/pkg/pubsub/pulljob.go b/pkg/pubsub/pulljob.go index 2a14dcb..ff72fba 100644 --- a/pkg/pubsub/pulljob.go +++ b/pkg/pubsub/pulljob.go @@ -55,7 +55,7 @@ func WithErrorHandler(handler func(err error) bool) PullJobOption { } func newPullJob(s *Subscriber, opts []PullJobOption) (*pullJob, error) { - b := &pullJob{ + job := &pullJob{ subscription: s, maxPullMessages: 10, interval: 1, @@ -68,20 +68,20 @@ func newPullJob(s *Subscriber, opts []PullJobOption) (*pullJob, error) { } for _, opt := range opts { - opt(b) + opt(job) } - if b.interval < 1 { + if job.interval < 1 { return nil, &ConfigurationError{ - Msg: fmt.Sprintf("interval must be at least set to 1, got %d", b.interval), + Msg: fmt.Sprintf("interval must be at least set to 1, got %d", job.interval), } } - return b, nil + return job, nil } -func (b *pullJob) runLoop(ctx context.Context, handler func(context.Context, PullMessages)) { - ticker := time.NewTicker(b.interval) +func (j *pullJob) runLoop(ctx context.Context, handler func(context.Context, PullMessages)) { + ticker := time.NewTicker(j.interval) defer ticker.Stop() for { @@ -89,24 +89,24 @@ func (b *pullJob) runLoop(ctx context.Context, handler func(context.Context, Pul case <-ctx.Done(): return case <-ticker.C: - pullOpts := []PullOption{WithMaxMessages(b.maxPullMessages), WithLongPullDuration(b.longPullDuration)} - messages, err := b.subscription.Pull(ctx, pullOpts...) + pullOpts := []PullOption{WithMaxMessages(j.maxPullMessages), WithLongPullDuration(j.longPullDuration)} + messages, err := j.subscription.Pull(ctx, pullOpts...) if err != nil { //nolint:nestif var sdkErr SDKError // Declare the target variable if errors.As(err, &sdkErr) { // Pass a pointer to sdkErr if !sdkErr.IsTransient() { // Only exit the loop if the users error handler returns false - if !b.errHandler(err) { + if !j.errHandler(err) { return } continue } - b.subscription.logger.Error(err, "transient error, retrying") + j.subscription.logger.Error(err, "transient error, retrying") continue } - b.subscription.logger.Error(err, "unknown error, retrying") + j.subscription.logger.Error(err, "unknown error, retrying") continue } diff --git a/pkg/pubsub/subscriber_test.go b/pkg/pubsub/subscriber_test.go index 422c56c..fb8a790 100644 --- a/pkg/pubsub/subscriber_test.go +++ b/pkg/pubsub/subscriber_test.go @@ -131,6 +131,8 @@ var _ = Describe("PullJob", func() { err := publisher.Purge(ctx) Expect(err).ToNot(HaveOccurred()) + time.Sleep(1 * time.Second) // wait for topic to be purged + // publishing test messages _, err = publisher.PublishStrings(ctx, "testMessage") Expect(err).ToNot(HaveOccurred()) @@ -139,36 +141,14 @@ var _ = Describe("PullJob", func() { Context("using PullJobChan", func() { It("should receive a message from channel", func(ctx context.Context) { subscriber := pubsub.NewSubscriber(topicId, subscriptionId, pubsub.WithHTTPRoundTripper(rt), pubsub.WithHost(environment)) + defer subscriber.Wait() - ctx, cancel := context.WithTimeout(ctx, 5*time.Second) - defer cancel() - - jobChan, err := subscriber.PullJobChan(ctx, pubsub.WithInterval(100*time.Millisecond)) - Expect(err).ToNot(HaveOccurred()) - - var receivedMessages pubsub.PullMessages - Eventually(jobChan, "5s").Should(Receive(&receivedMessages)) - Expect(receivedMessages).To(HaveLen(1)) - - decodedString, err := receivedMessages[0].DecodeString() - Expect(err).ToNot(HaveOccurred()) - Expect(decodedString).To(Equal("testMessage")) - err = subscriber.Ack(ctx, receivedMessages.AckIDs()) - Expect(err).ToNot(HaveOccurred()) - - cancel() - subscriber.Wait() - }) - - It("should receive a message from channel with long pull duration set", func(ctx context.Context) { - subscriber := pubsub.NewSubscriber(topicId, subscriptionId, pubsub.WithHTTPRoundTripper(rt), pubsub.WithHost(environment)) - - ctx, cancel := context.WithTimeout(ctx, 10*time.Second) + ctx, cancel := context.WithTimeout(ctx, 15*time.Second) defer cancel() jobChan, err := subscriber.PullJobChan(ctx, - pubsub.WithInterval(100*time.Millisecond), - pubsub.WithPullLongPullDuration(500), + pubsub.WithPullLongPullDuration(100), + pubsub.WithPullMaxMessages(1), ) Expect(err).ToNot(HaveOccurred()) @@ -183,7 +163,6 @@ var _ = Describe("PullJob", func() { Expect(err).ToNot(HaveOccurred()) cancel() - subscriber.Wait() }) }) @@ -191,7 +170,7 @@ var _ = Describe("PullJob", func() { It("should invoke the callback with messages", func(ctx context.Context) { subscriber := pubsub.NewSubscriber(topicId, subscriptionId, pubsub.WithHTTPRoundTripper(rt), pubsub.WithHost(environment)) defer subscriber.Wait() - ctx, cancel := context.WithTimeout(ctx, 5*time.Second) + ctx, cancel := context.WithTimeout(ctx, 15*time.Second) defer cancel() var callbackInvoked atomic.Bool // check if callback was invoked @@ -210,8 +189,8 @@ var _ = Describe("PullJob", func() { err := subscriber.PullJobCallback( ctx, callback, - pubsub.WithInterval(100*time.Millisecond), pubsub.WithPullMaxMessages(1), + pubsub.WithPullLongPullDuration(100), ) Expect(err).ToNot(HaveOccurred())