From 19c806b49bb7b4f023a28765c132438ae245519a Mon Sep 17 00:00:00 2001 From: Maxim Kolosov Date: Fri, 2 Oct 2026 13:05:22 -0600 Subject: [PATCH 1/3] Back off on-demand Metadata refreshes while a leader stays excluded --- README.md | 4 +- pkg/wgo/client.go | 71 ++++++++--- pkg/wgo/client_test.go | 267 ++++++++++++++++++++++++++++++++++++++++- pkg/wgo/config.go | 17 +-- 4 files changed, 331 insertions(+), 28 deletions(-) diff --git a/README.md b/README.md index 8e1832a..f1f6fea 100644 --- a/README.md +++ b/README.md @@ -26,7 +26,7 @@ This client has been designed around the following non-negotiable assumptions: 1. **Warpstream-specific.** Hedging the same batch across agents only works because any agent can serve any partition. Pointed at vanilla Kafka, the secondary leg would fail with `NotLeaderForPartition`. 2. **At-least-once delivery only.** Duplicates are tolerable. Any code that assumes exactly-once or in-partition record ordering must stay on franz-go. 3. **No transactional or idempotent producer support.** `DisableIdempotentWrite()` semantics are baked in — no `producerId`/`producerEpoch`/`baseSequence` handshake. -4. **Produce never blocks on Metadata.** The agent pool is refreshed on a timer and also on demand. A routing miss fails that Produce immediately and asks for a refresh; a later Produce can use the updated pool. A refresh that excluded a leader asks for another fetch even when that Produce already succeeded on a stand-in. Neither waits on the fetch. On-demand refreshes are coalesced and paced by `OnDemandMetadataRefreshInterval` (default 1s). +4. **Produce never blocks on Metadata.** The agent pool is refreshed on a timer and also on demand. A routing miss fails that Produce immediately and asks for a refresh; a later Produce can use the updated pool. A refresh that excluded a leader asks for another fetch on its own, with no Produce in flight. Those follow-ups start at `OnDemandMetadataRefreshInterval` (default 1s) and double up to `MetadataRefreshInterval` (default 10s) while the exclusion continues. The chain stops when a refresh excludes nobody: that fetch does not ask for another, the periodic refresh runs, and the gap returns to the floor. Neither the Produce nor the follow-up waits on the fetch. On-demand refreshes are coalesced. ## How it works @@ -48,7 +48,7 @@ For every record the client asks a `PartitionAssignmentStrategy` for an ordered - **Deterministic.** Given the same Metadata view, every client instance picks the same primary and the same secondary for a given partition. Hedge load is predictable and analysable instead of randomly smeared across agents. - **State-aware.** A wrapper around the base strategy (the **Demoter**, see below) can mark an agent as demoted so it's elided from the candidate list or surfaced as a probe. -A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason. +A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own, even when no Produce is in flight. Further exclusions double the wait from `OnDemandMetadataRefreshInterval` up to `MetadataRefreshInterval`. A refresh that excludes nobody stops the follow-ups, and the next periodic refresh puts the wait back. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason. ### Buffering: linger by destination agent, not by partition diff --git a/pkg/wgo/client.go b/pkg/wgo/client.go index 4fed0c3..971febc 100644 --- a/pkg/wgo/client.go +++ b/pkg/wgo/client.go @@ -383,26 +383,69 @@ func (c *WarpstreamClient) Close() { // interrupts an in-flight Refresh, and aborts the cooldown. We rely on kgo's // per-attempt RequestTimeoutOverhead and overall retry policy rather than // adding a separate per-refresh deadline. +// refreshBackoff paces on-demand Metadata fetches. The delay starts at floor, +// grows after each on-demand fetch up to ceiling, and returns to floor when a +// periodic fetch runs. +type refreshBackoff struct { + floor time.Duration + ceiling time.Duration + delay time.Duration +} + +func newRefreshBackoff(floor, ceiling time.Duration) refreshBackoff { + return refreshBackoff{floor: floor, ceiling: ceiling, delay: floor} +} + +func (b *refreshBackoff) current() time.Duration { return b.delay } + +func (b *refreshBackoff) advance() { + // time.Duration is an int64. Doubling past MaxInt64 wraps negative. + // At or below half the ceiling, delay*2 still fits and stays within it. + if b.delay > b.ceiling/2 { + b.delay = b.ceiling + return + } + b.delay *= 2 +} + +func (b *refreshBackoff) reset() { + b.delay = b.floor +} + func (c *WarpstreamClient) startBackgroundRefresh() { c.refreshWG.Add(1) go func() { defer c.refreshWG.Done() ticker := time.NewTicker(c.cfg.MetadataRefreshInterval) defer ticker.Stop() + backoff := newRefreshBackoff(c.cfg.OnDemandMetadataRefreshInterval, c.cfg.MetadataRefreshInterval) for { + periodic := false select { case <-c.refreshCtx.Done(): return case <-c.refreshNowCh: - ticker.Reset(c.cfg.MetadataRefreshInterval) - startedAt := time.Now() - c.refreshPool(metadataRefreshTriggerOnDemand) - if !c.waitRefreshCooldown(time.Since(startedAt)) { - return - } case <-ticker.C: + // A nudge queued at the same instant wins. At the ceiling the + // cooldown ends as the tick lands, and taking the tick would + // reset the delay. + select { + case <-c.refreshNowCh: + default: + periodic = true + } + } + + if periodic { c.refreshPool(metadataRefreshTriggerPeriodic) + backoff.reset() + continue } + ticker.Reset(c.cfg.MetadataRefreshInterval) + startedAt := time.Now() + c.refreshPool(metadataRefreshTriggerOnDemand) + c.waitRefreshCooldown(backoff.current(), time.Since(startedAt)) + backoff.advance() } }() } @@ -431,8 +474,7 @@ func (c *WarpstreamClient) refreshPool(trigger metadataRefreshTrigger) { } c.demoter.Refresh(c.pool.Agents()) // Nudge so a stand-in does not wait for the next periodic tick. Produce does - // not block. After a periodic tick the follow-up starts at once; later - // repeats wait out OnDemandMetadataRefreshInterval (default 1s). + // not block. The refresh loop paces the follow-up. c.noteLeaderDrops(dropped, true) } @@ -452,21 +494,18 @@ func (c *WarpstreamClient) noteLeaderDrops(dropped leaderDrops, nudge bool) { } } -// waitRefreshCooldown enforces the configured minimum between on-demand -// refresh start times. Time spent fetching already counts toward the interval. -// Returns false during close so the refresh loop exits without spinning. -func (c *WarpstreamClient) waitRefreshCooldown(elapsed time.Duration) bool { - remaining := c.cfg.OnDemandMetadataRefreshInterval - elapsed +// waitRefreshCooldown sleeps out the rest of delay. Time already spent +// fetching counts. Close cancels refreshCtx and the wait returns. +func (c *WarpstreamClient) waitRefreshCooldown(delay, elapsed time.Duration) { + remaining := delay - elapsed if remaining <= 0 { - return c.refreshCtx.Err() == nil + return } t := time.NewTimer(remaining) defer t.Stop() select { case <-c.refreshCtx.Done(): - return false case <-t.C: - return true } } diff --git a/pkg/wgo/client_test.go b/pkg/wgo/client_test.go index 3d786c9..c844ce5 100644 --- a/pkg/wgo/client_test.go +++ b/pkg/wgo/client_test.go @@ -3,6 +3,7 @@ package wgo import ( "bytes" "context" + "math" "slices" "sync" "sync/atomic" @@ -1639,7 +1640,7 @@ func TestWarpstreamClient_WaitRefreshCooldown(t *testing.T) { } startedAt := time.Now() - assert.True(t, c.waitRefreshCooldown(250*time.Millisecond)) + c.waitRefreshCooldown(time.Second, 250*time.Millisecond) assert.Equal(t, 750*time.Millisecond, time.Since(startedAt)) }) }) @@ -1654,10 +1655,272 @@ func TestWarpstreamClient_WaitRefreshCooldown(t *testing.T) { } startedAt := time.Now() - assert.True(t, c.waitRefreshCooldown(2*time.Second)) + c.waitRefreshCooldown(time.Second, 2*time.Second) assert.Zero(t, time.Since(startedAt)) }) }) + + t.Run("close interrupts the wait", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + c := &WarpstreamClient{ + cfg: Config{OnDemandMetadataRefreshInterval: time.Second}, + refreshCtx: ctx, + } + + startedAt := time.Now() + go func() { + time.Sleep(100 * time.Millisecond) + cancel() + }() + c.waitRefreshCooldown(time.Second, 0) + assert.Equal(t, 100*time.Millisecond, time.Since(startedAt)) + }) + }) +} + +func TestRefreshBackoff_Advance(t *testing.T) { + b := newRefreshBackoff(time.Second, 10*time.Second) + assert.Equal(t, time.Second, b.current()) + for _, want := range []time.Duration{2 * time.Second, 4 * time.Second, 8 * time.Second, 10 * time.Second, 10 * time.Second} { + b.advance() + assert.Equal(t, want, b.current()) + } + b.reset() + assert.Equal(t, time.Second, b.current()) + + // Doubling this delay would wrap time.Duration. advance saturates instead. + ceiling := time.Duration(math.MaxInt64) + b = newRefreshBackoff(ceiling/2+1, ceiling) + b.advance() + assert.Equal(t, ceiling, b.current()) +} + +func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) { + const topic = "test-topic" + + t.Run("leader drop backs off until a clean snapshot restores the leader", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + vnet := &kfake.VirtualNetwork{} + cluster, clusterAddr := testkafka.CreateCluster(t, 1, topic, + testkafka.WithVirtualNetwork(vnet), testkafka.WithNumBrokers(3)) + c, err := NewWarpstreamClient(nil, prometheus.NewPedanticRegistry(), append( + testWarpstreamOpts(clusterAddr, topic), WithDialer(vnet.DialContext))...) + require.NoError(t, err) + t.Cleanup(c.Close) + + before := c.demoter.Candidates(topic, 0, 1) + require.Len(t, before, 1) + leader := before[0].NodeID + + raw, err := c.Request(t.Context(), kmsg.NewPtrMetadataRequest()) + require.NoError(t, err) + template := raw.(*kmsg.MetadataResponse) + + var poison atomic.Bool + cluster.ControlKey(int16(kmsg.Metadata), func(req kmsg.Request) (kmsg.Response, error, bool) { + cluster.KeepControl() + if !poison.Load() { + return nil, nil, false + } + return metadataWithMissingLeaders(template, 99, req.GetVersion()), nil, true + }) + + poison.Store(true) + time.Sleep(time.Nanosecond) + c.triggerRefresh() + synctest.Wait() + assertOnDemandRefresh(t, c, 1, 0) + assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal)) + assertMissingLeader(t, c, topic) + + results := c.ProduceSync(t.Context(), []*kgo.Record{ + {Topic: topic, Partition: 0, Value: []byte("stand-in"), Timestamp: time.Now()}, + }) + require.Len(t, results, 1) + require.NoError(t, results[0].Err) + + // Gaps after the first fetch: 1s, 2s, 4s, 8s, then 10s. The 10s + // gap ends as the periodic tick lands; the queued nudge keeps it + // on-demand, so the delay does not fall back to 1s. + for i, gap := range []time.Duration{time.Second, 2 * time.Second, 4 * time.Second, 8 * time.Second, 10 * time.Second} { + time.Sleep(gap) + synctest.Wait() + assertOnDemandRefresh(t, c, float64(i+2), 0) + } + assert.Equal(t, float64(6), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal)) + + time.Sleep(time.Second) + synctest.Wait() + assertOnDemandRefresh(t, c, 6, 0) + + poison.Store(false) + time.Sleep(9 * time.Second) + synctest.Wait() + assertOnDemandRefresh(t, c, 7, 0) + assert.Equal(t, float64(6), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal)) + assertLeaderRestored(t, c, topic, leader) + + restored := c.ProduceSync(t.Context(), []*kgo.Record{ + {Topic: topic, Partition: 0, Value: []byte("leader"), Timestamp: time.Now()}, + }) + require.Len(t, restored, 1) + require.NoError(t, restored[0].Err) + + // The clean fetch did not ask for another. The periodic tick + // resets the gap to the floor. + time.Sleep(9 * time.Second) + synctest.Wait() + assertOnDemandRefresh(t, c, 7, 0) + time.Sleep(time.Second) + synctest.Wait() + assertOnDemandRefresh(t, c, 7, 1) + + c.triggerRefresh() + synctest.Wait() + assertOnDemandRefresh(t, c, 8, 1) + // Queue the next nudge during the floor cooldown. A delay still + // at the ceiling would not fetch again inside this second. + c.triggerRefresh() + time.Sleep(500 * time.Millisecond) + synctest.Wait() + assertOnDemandRefresh(t, c, 8, 1) + time.Sleep(500 * time.Millisecond) + synctest.Wait() + assertOnDemandRefresh(t, c, 9, 1) + }) + }) + + t.Run("routing miss uses the same backoff", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + c, _, _, _ := newTestWarpstreamClient(t, topic, 1) + + time.Sleep(time.Nanosecond) + miss := func() { + results := c.ProduceSync(t.Context(), []*kgo.Record{ + {Topic: "does-not-exist", Partition: 0, Value: []byte("v"), Timestamp: time.Now()}, + }) + require.ErrorContains(t, results[0].Err, "no agent assigned") + } + + miss() + synctest.Wait() + assertOnDemandRefresh(t, c, 1, 0) + + miss() + time.Sleep(500 * time.Millisecond) + synctest.Wait() + assertOnDemandRefresh(t, c, 1, 0) + time.Sleep(500 * time.Millisecond) + synctest.Wait() + assertOnDemandRefresh(t, c, 2, 0) + + miss() + time.Sleep(time.Second) + synctest.Wait() + assertOnDemandRefresh(t, c, 2, 0) + time.Sleep(time.Second) + synctest.Wait() + assertOnDemandRefresh(t, c, 3, 0) + }) + }) + + t.Run("failed fetch advances backoff and close interrupts the cooldown", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + c, cluster, _, _ := newTestWarpstreamClient(t, topic, 1) + raw, err := c.Request(t.Context(), kmsg.NewPtrMetadataRequest()) + require.NoError(t, err) + template := raw.(*kmsg.MetadataResponse) + // A response error fails the fetch in one round trip. Closing the + // connection instead makes kgo retry, and that retry time would + // swallow the cooldown this test is measuring. + cluster.ControlKey(int16(kmsg.Metadata), func(req kmsg.Request) (kmsg.Response, error, bool) { + cluster.KeepControl() + resp := cloneMetadataResponse(template, req.GetVersion()) + resp.ErrorCode = kerr.UnknownServerError.Code + return resp, nil, true + }) + + time.Sleep(time.Nanosecond) + c.triggerRefresh() + synctest.Wait() + assert.Equal(t, float64(1), failedOnDemandRefreshes(c)) + + c.triggerRefresh() + time.Sleep(500 * time.Millisecond) + synctest.Wait() + assert.Equal(t, float64(1), failedOnDemandRefreshes(c)) + time.Sleep(500 * time.Millisecond) + synctest.Wait() + assert.Equal(t, float64(2), failedOnDemandRefreshes(c)) + + c.triggerRefresh() + time.Sleep(time.Second) + synctest.Wait() + assert.Equal(t, float64(2), failedOnDemandRefreshes(c)) + time.Sleep(time.Second) + synctest.Wait() + assert.Equal(t, float64(3), failedOnDemandRefreshes(c)) + + startedAt := time.Now() + c.Close() + assert.Zero(t, time.Since(startedAt)) + }) + }) +} + +func assertOnDemandRefresh(t *testing.T, c *WarpstreamClient, onDemand, periodic float64) { + t.Helper() + assert.Equal(t, onDemand, testutil.ToFloat64(c.metrics.metadataRefreshResultsTotal.WithLabelValues( + string(metadataRefreshTriggerOnDemand), metadataRefreshResultUnchanged))) + assert.Equal(t, periodic, testutil.ToFloat64(c.metrics.metadataRefreshResultsTotal.WithLabelValues( + string(metadataRefreshTriggerPeriodic), metadataRefreshResultUnchanged))) +} + +func failedOnDemandRefreshes(c *WarpstreamClient) float64 { + return testutil.ToFloat64(c.metrics.metadataRefreshResultsTotal.WithLabelValues( + string(metadataRefreshTriggerOnDemand), metadataRefreshResultFailed)) +} + +func assertMissingLeader(t *testing.T, c *WarpstreamClient, topic string) { + t.Helper() + _, ok := c.pool.Strategy().(*DefaultPartitionAssignmentStrategy).leaders[topicPartition{topic: topic, partition: 0}] + assert.False(t, ok) + cands := c.demoter.Candidates(topic, 0, 1) + require.Len(t, cands, 1) + assert.NotEqual(t, int32(99), cands[0].NodeID) + assert.Contains(t, c.pool.Agents(), cands[0].NodeID) +} + +func assertLeaderRestored(t *testing.T, c *WarpstreamClient, topic string, leader int32) { + t.Helper() + got, ok := c.pool.Strategy().(*DefaultPartitionAssignmentStrategy).leaders[topicPartition{topic: topic, partition: 0}] + require.True(t, ok) + assert.Equal(t, leader, got) + cands := c.demoter.Candidates(topic, 0, 1) + require.Len(t, cands, 1) + assert.Equal(t, leader, cands[0].NodeID) +} + +func cloneMetadataResponse(src *kmsg.MetadataResponse, version int16) *kmsg.MetadataResponse { + out := *src + out.Version = version + out.Brokers = slices.Clone(src.Brokers) + out.Topics = slices.Clone(src.Topics) + for i := range out.Topics { + out.Topics[i].Partitions = slices.Clone(src.Topics[i].Partitions) + } + return &out +} + +func metadataWithMissingLeaders(src *kmsg.MetadataResponse, missingLeader int32, version int16) *kmsg.MetadataResponse { + out := cloneMetadataResponse(src, version) + for i := range out.Topics { + for j := range out.Topics[i].Partitions { + out.Topics[i].Partitions[j].Leader = missingLeader + } + } + return out } // lockedBuffer is a concurrency-safe sink for logger output, which franz-go diff --git a/pkg/wgo/config.go b/pkg/wgo/config.go index 7e6862a..2a6574d 100644 --- a/pkg/wgo/config.go +++ b/pkg/wgo/config.go @@ -71,10 +71,11 @@ type Config struct { // purges agent-stats entries for agents that have left the cluster. MetadataRefreshInterval time.Duration - // OnDemandMetadataRefreshInterval is the minimum time between the start of - // on-demand Metadata refreshes. That includes a routing miss, and repeated - // follow-ups while a refresh keeps excluding a leader. Must be > 0 and no - // greater than MetadataRefreshInterval. + // OnDemandMetadataRefreshInterval is the shortest gap between the start of + // on-demand Metadata refreshes. Repeated on-demand fetches double that gap + // up to MetadataRefreshInterval, and a periodic fetch puts it back. That + // includes a routing miss, and follow-ups while a refresh keeps excluding + // a leader. Must be > 0 and no greater than MetadataRefreshInterval. OnDemandMetadataRefreshInterval time.Duration } @@ -336,10 +337,10 @@ func WithMetadataRefreshInterval(d time.Duration) Opt { return opt{func(c *Config) { c.MetadataRefreshInterval = d }} } -// WithOnDemandMetadataRefreshInterval sets the minimum time between the start -// of on-demand Metadata refreshes. That includes a routing miss, and repeated -// follow-ups while a refresh keeps excluding a leader. It must be no greater -// than the background MetadataRefreshInterval. +// WithOnDemandMetadataRefreshInterval sets the shortest gap between the start +// of on-demand Metadata refreshes. Repeated on-demand fetches double that gap +// up to the background MetadataRefreshInterval, which puts the gap back when +// it runs. The gap must be no greater than that interval. func WithOnDemandMetadataRefreshInterval(d time.Duration) Opt { return opt{func(c *Config) { c.OnDemandMetadataRefreshInterval = d }} } From 7cb273be94fef861c15aa54d05e43c92f264799c Mon Sep 17 00:00:00 2001 From: Maxim Kolosov Date: Fri, 2 Oct 2026 13:32:23 -0600 Subject: [PATCH 2/3] Return from the refresh loop when Close cancels an on-demand wait. --- pkg/wgo/client.go | 4 ++++ pkg/wgo/client_test.go | 37 +++++++++++++++++++++++++++++++++++++ 2 files changed, 41 insertions(+) diff --git a/pkg/wgo/client.go b/pkg/wgo/client.go index 971febc..92c7204 100644 --- a/pkg/wgo/client.go +++ b/pkg/wgo/client.go @@ -445,6 +445,10 @@ func (c *WarpstreamClient) startBackgroundRefresh() { startedAt := time.Now() c.refreshPool(metadataRefreshTriggerOnDemand) c.waitRefreshCooldown(backoff.current(), time.Since(startedAt)) + // Return before select can take the nudge the fetch just queued. + if c.refreshCtx.Err() != nil { + return + } backoff.advance() } }() diff --git a/pkg/wgo/client_test.go b/pkg/wgo/client_test.go index c844ce5..bb72e1f 100644 --- a/pkg/wgo/client_test.go +++ b/pkg/wgo/client_test.go @@ -1867,6 +1867,43 @@ func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) { assert.Zero(t, time.Since(startedAt)) }) }) + + t.Run("close during a leader-drop cooldown does not fetch again", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + vnet := &kfake.VirtualNetwork{} + cluster, clusterAddr := testkafka.CreateCluster(t, 1, topic, + testkafka.WithVirtualNetwork(vnet), testkafka.WithNumBrokers(3)) + c, err := NewWarpstreamClient(nil, prometheus.NewPedanticRegistry(), append( + testWarpstreamOpts(clusterAddr, topic), WithDialer(vnet.DialContext))...) + require.NoError(t, err) + t.Cleanup(c.Close) + + raw, err := c.Request(t.Context(), kmsg.NewPtrMetadataRequest()) + require.NoError(t, err) + template := raw.(*kmsg.MetadataResponse) + + var poison atomic.Bool + cluster.ControlKey(int16(kmsg.Metadata), func(req kmsg.Request) (kmsg.Response, error, bool) { + cluster.KeepControl() + if !poison.Load() { + return nil, nil, false + } + return metadataWithMissingLeaders(template, 99, req.GetVersion()), nil, true + }) + + poison.Store(true) + time.Sleep(time.Nanosecond) + c.triggerRefresh() + synctest.Wait() + assertOnDemandRefresh(t, c, 1, 0) + + // noteLeaderDrops already queued the next nudge. Close must not run it. + startedAt := time.Now() + c.Close() + assert.Zero(t, time.Since(startedAt)) + assertOnDemandRefresh(t, c, 1, 0) + }) + }) } func assertOnDemandRefresh(t *testing.T, c *WarpstreamClient, onDemand, periodic float64) { From 641fee0e1933d07c520dcba6beacd385646a0a3b Mon Sep 17 00:00:00 2001 From: Maxim Kolosov Date: Fri, 2 Oct 2026 14:03:33 -0600 Subject: [PATCH 3/3] Return from the refresh loop when Close cancels a periodic fetch --- README.md | 2 +- pkg/wgo/client.go | 16 ++++++++++------ pkg/wgo/client_test.go | 40 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 51 insertions(+), 7 deletions(-) diff --git a/README.md b/README.md index f1f6fea..a0faad3 100644 --- a/README.md +++ b/README.md @@ -48,7 +48,7 @@ For every record the client asks a `PartitionAssignmentStrategy` for an ordered - **Deterministic.** Given the same Metadata view, every client instance picks the same primary and the same secondary for a given partition. Hedge load is predictable and analysable instead of randomly smeared across agents. - **State-aware.** A wrapper around the base strategy (the **Demoter**, see below) can mark an agent as demoted so it's elided from the candidate list or surfaced as a probe. -A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own, even when no Produce is in flight. Further exclusions double the wait from `OnDemandMetadataRefreshInterval` up to `MetadataRefreshInterval`. A refresh that excludes nobody stops the follow-ups, and the next periodic refresh puts the wait back. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason. +A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own, even when no Produce is in flight. Further exclusions double the wait from `OnDemandMetadataRefreshInterval` up to `MetadataRefreshInterval`. Repeated routing misses climb the same way. While on-demand refreshes keep being requested, the periodic refresh does not run, so the gap stays at the ceiling until they stop. A refresh that excludes nobody stops the follow-ups, and the next periodic refresh puts the wait back. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason. ### Buffering: linger by destination agent, not by partition diff --git a/pkg/wgo/client.go b/pkg/wgo/client.go index 92c7204..b169bff 100644 --- a/pkg/wgo/client.go +++ b/pkg/wgo/client.go @@ -377,12 +377,6 @@ func (c *WarpstreamClient) Close() { }) } -// startBackgroundRefresh owns every post-startup AgentPool.Refresh: the -// periodic ticker and on-demand nudges from triggerRefresh. One owner means -// Refresh is never concurrent. refreshCtx (cancelled by Close) stops the loop, -// interrupts an in-flight Refresh, and aborts the cooldown. We rely on kgo's -// per-attempt RequestTimeoutOverhead and overall retry policy rather than -// adding a separate per-refresh deadline. // refreshBackoff paces on-demand Metadata fetches. The delay starts at floor, // grows after each on-demand fetch up to ceiling, and returns to floor when a // periodic fetch runs. @@ -412,6 +406,12 @@ func (b *refreshBackoff) reset() { b.delay = b.floor } +// startBackgroundRefresh owns every post-startup AgentPool.Refresh: the +// periodic ticker and on-demand nudges from triggerRefresh. One owner means +// Refresh is never concurrent. refreshCtx (cancelled by Close) stops the loop, +// interrupts an in-flight Refresh, and aborts the cooldown. We rely on kgo's +// per-attempt RequestTimeoutOverhead and overall retry policy rather than +// adding a separate per-refresh deadline. func (c *WarpstreamClient) startBackgroundRefresh() { c.refreshWG.Add(1) go func() { @@ -438,6 +438,10 @@ func (c *WarpstreamClient) startBackgroundRefresh() { if periodic { c.refreshPool(metadataRefreshTriggerPeriodic) + // A nudge queued during this fetch must not start another after Close. + if c.refreshCtx.Err() != nil { + return + } backoff.reset() continue } diff --git a/pkg/wgo/client_test.go b/pkg/wgo/client_test.go index bb72e1f..b31f7ec 100644 --- a/pkg/wgo/client_test.go +++ b/pkg/wgo/client_test.go @@ -1904,6 +1904,40 @@ func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) { assertOnDemandRefresh(t, c, 1, 0) }) }) + + t.Run("close during a periodic fetch does not run a queued nudge", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + c, cluster, _, _ := newTestWarpstreamClient(t, topic, 1) + release := make(chan struct{}) + var blockNext atomic.Bool + blockNext.Store(true) + cluster.ControlKey(int16(kmsg.Metadata), func(req kmsg.Request) (kmsg.Response, error, bool) { + cluster.KeepControl() + // AgentPool.Refresh asks for every topic. kgo's own loop asks + // for brokers only, and blocking that request stalls the cluster. + if req.(*kmsg.MetadataRequest).Topics == nil && blockNext.CompareAndSwap(true, false) { + cluster.SleepControl(func() { <-release }) + } + return nil, nil, false + }) + + time.Sleep(10 * time.Second) + synctest.Wait() + // Queue a nudge while a Metadata request may be in flight. The count + // below is taken after that nudge has either run or is still waiting + // behind the blocked request. + c.triggerRefresh() + synctest.Wait() + before := onDemandRefreshes(c) + + go c.Close() + synctest.Wait() + close(release) + synctest.Wait() + + assert.Equal(t, before, onDemandRefreshes(c)) + }) + }) } func assertOnDemandRefresh(t *testing.T, c *WarpstreamClient, onDemand, periodic float64) { @@ -1914,6 +1948,12 @@ func assertOnDemandRefresh(t *testing.T, c *WarpstreamClient, onDemand, periodic string(metadataRefreshTriggerPeriodic), metadataRefreshResultUnchanged))) } +func onDemandRefreshes(c *WarpstreamClient) float64 { + return testutil.ToFloat64(c.metrics.metadataRefreshResultsTotal.WithLabelValues( + string(metadataRefreshTriggerOnDemand), metadataRefreshResultUnchanged)) + + failedOnDemandRefreshes(c) +} + func failedOnDemandRefreshes(c *WarpstreamClient) float64 { return testutil.ToFloat64(c.metrics.metadataRefreshResultsTotal.WithLabelValues( string(metadataRefreshTriggerOnDemand), metadataRefreshResultFailed))