diff --git a/README.md b/README.md index 545cc6f..8e1832a 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 when routing finds no candidate. The current Produce call still fails immediately rather than waiting for the fetch; a later Produce can use the updated pool. On-demand refreshes are coalesced and paced by `OnDemandMetadataRefreshInterval` (default 1s) to avoid request storms. +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). ## 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 that happens while another leader entry still keeps the topic known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. A topic with no leader entries is left alone, so it still gets an on-demand refresh. +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. ### Buffering: linger by destination agent, not by partition diff --git a/docs/internal/metrics.md b/docs/internal/metrics.md index fa5e90b..4783390 100644 --- a/docs/internal/metrics.md +++ b/docs/internal/metrics.md @@ -47,9 +47,10 @@ identical without a duplicate registration. ## 3. Warpstream-specific metrics — `warpstream_` prefix Metrics with no franz-go counterpart describe behaviour unique to this client: -hedging, agent demotion, direct-request and attempt accounting, and -client-boundary record counters. They carry a `warpstream_` prefix so they never -collide with franz-go/kprom names and are unambiguously backend-specific. +hedging, agent demotion, leaders excluded from the broker list, direct-request +and attempt accounting, and client-boundary record counters. They carry a +`warpstream_` prefix so they never collide with franz-go/kprom names and are +unambiguously backend-specific. ## 4. Build info — `warpstream_client_build_info` diff --git a/pkg/wgo/agentpool.go b/pkg/wgo/agentpool.go index 9e105bc..be1148a 100644 --- a/pkg/wgo/agentpool.go +++ b/pkg/wgo/agentpool.go @@ -61,7 +61,7 @@ func NewAgentPool(client *kgo.Client) *AgentPool { p := &AgentPool{client: client} p.state.Store(&poolState{ topicIDs: map[string][16]byte{}, - strategy: newDefaultPartitionAssignmentStrategy(nil, nil), + strategy: newDefaultPartitionAssignmentStrategy(nil, nil, nil, nil), }) return p } @@ -69,14 +69,22 @@ func NewAgentPool(client *kgo.Client) *AgentPool { // Refresh atomically replaces the snapshot. Returns the NodeIDs that have left // the cluster since the last Refresh — callers must purge per-agent state // (stats, etc.) for those IDs. Not safe for concurrent calls. -func (p *AgentPool) Refresh(ctx context.Context) (removed []int32, err error) { +func (p *AgentPool) Refresh(ctx context.Context) ([]int32, error) { + removed, _, err := p.refresh(ctx) + return removed, err +} + +// refresh atomically replaces the snapshot. removed is the NodeIDs that have +// left since the last refresh. dropped counts leaders excluded from the broker +// set and names one when the count is above zero. Not safe for concurrent calls. +func (p *AgentPool) refresh(ctx context.Context) (removed []int32, dropped leaderDrops, err error) { // Topics=nil requests metadata for every topic in the cluster. // A tiny positive cache age keeps production refreshes fresh while using // kgo's bounded internal Metadata retry policy. req := kmsg.NewPtrMetadataRequest() meta, err := p.client.RequestCachedMetadata(ctx, req, time.Nanosecond) if err != nil { - return nil, fmt.Errorf("fetching metadata: %w", err) + return nil, leaderDrops{}, fmt.Errorf("fetching metadata: %w", err) } newAgents := make([]int32, 0, len(meta.Brokers)) @@ -91,15 +99,15 @@ func (p *AgentPool) Refresh(ctx context.Context) (removed []int32, err error) { } prev := p.state.Load() - newLeaders, newTopicIDs := buildLeadersAndTopicIDs(meta.Topics, agentSet, prev.topicIDs) + newLeaders, newTopicIDs, noLiveLeader, noLeader, dropped := buildLeadersAndTopicIDs(meta.Topics, agentSet, prev.topicIDs) removed = diffRemovedAgents(prev.agents, agentSet) p.state.Store(&poolState{ agents: newAgents, topicIDs: newTopicIDs, - strategy: newDefaultPartitionAssignmentStrategy(newAgents, newLeaders), + strategy: newDefaultPartitionAssignmentStrategy(newAgents, newLeaders, noLiveLeader, noLeader), }) - return removed, nil + return removed, dropped, nil } // Strategy returns the strategy from the last Refresh. @@ -120,18 +128,34 @@ func (p *AgentPool) TopicID(topic string) ([16]byte, bool) { return id, ok } +// leaderDrops counts leaders whose NodeID was absent from the broker set. +// Topic, Partition, and NodeID are one sample when Count is above zero. +type leaderDrops struct { + Count int + Topic string + Partition int32 + NodeID int32 +} + // buildLeadersAndTopicIDs extracts the leader map and topic UUIDs from a -// Metadata response. Drops leaders pointing to NodeIDs absent from agentSet -// (transient mid-update window). Carries previous UUIDs forward for topics -// returned with a non-zero ErrorCode, so transient errors don't blank the -// topic from the producer's view. +// Metadata response, plus topics whose named leaders were all excluded (nil +// when there are none) and the excluded-leader count with one sample. +// Drops leaders pointing to NodeIDs absent from agentSet (transient +// mid-update window). Carries previous UUIDs forward for topics returned +// with a non-zero ErrorCode, so transient errors don't blank the topic +// from the producer's view. A topic-level ErrorCode is not a drop. A +// partition Leader below 0 names no node, so it is not a drop and not a +// fallback. noLeader is nil when there are none. func buildLeadersAndTopicIDs( respTopics []kmsg.MetadataResponseTopic, agentSet map[int32]struct{}, prevTopicIDs map[string][16]byte, -) (map[topicPartition]int32, map[string][16]byte) { +) (map[topicPartition]int32, map[string][16]byte, map[string]struct{}, map[topicPartition]struct{}, leaderDrops) { topicIDs := make(map[string][16]byte, len(respTopics)) leaders := make(map[topicPartition]int32, len(respTopics)*8) + var topicsWithNoLiveLeader map[string]struct{} + var noLeader map[topicPartition]struct{} + var dropped leaderDrops for _, t := range respTopics { if t.Topic == nil { continue @@ -145,14 +169,38 @@ func buildLeadersAndTopicIDs( continue } topicIDs[name] = t.TopicID + kept, excluded := 0, 0 for _, part := range t.Partitions { + // Leader below 0 names nobody. It is not a node id missing from + // the broker list, so it must not count as a drop or fall back. + if part.Leader < 0 { + if noLeader == nil { + noLeader = make(map[topicPartition]struct{}) + } + noLeader[topicPartition{topic: name, partition: part.Partition}] = struct{}{} + continue + } if _, known := agentSet[part.Leader]; !known { + if dropped.Count == 0 { + dropped.Topic = name + dropped.Partition = part.Partition + dropped.NodeID = part.Leader + } + dropped.Count++ + excluded++ continue } + kept++ leaders[topicPartition{topic: name, partition: part.Partition}] = part.Leader } + if excluded > 0 && kept == 0 { + if topicsWithNoLiveLeader == nil { + topicsWithNoLiveLeader = make(map[string]struct{}) + } + topicsWithNoLiveLeader[name] = struct{}{} + } } - return leaders, topicIDs + return leaders, topicIDs, topicsWithNoLiveLeader, noLeader, dropped } // diffRemovedAgents returns NodeIDs in old that are absent from newSet. diff --git a/pkg/wgo/agentpool_test.go b/pkg/wgo/agentpool_test.go index 0afae74..db17c52 100644 --- a/pkg/wgo/agentpool_test.go +++ b/pkg/wgo/agentpool_test.go @@ -30,9 +30,10 @@ func TestAgentPool_Refresh(t *testing.T) { pool := NewAgentPool(client) - removed, err := pool.Refresh(t.Context()) + removed, dropped, err := pool.refresh(t.Context()) require.NoError(t, err) assert.Empty(t, removed) + assert.Zero(t, dropped.Count) assert.NotNil(t, pool.Strategy()) id, ok := pool.TopicID(topicName) @@ -154,10 +155,13 @@ func TestBuildLeadersAndTopicIDs(t *testing.T) { idB := [16]byte{0x43} tests := map[string]struct { - respTopics []kmsg.MetadataResponseTopic - prevTopicIDs map[string][16]byte - wantLeaders map[topicPartition]int32 - wantTopicIDs map[string][16]byte + respTopics []kmsg.MetadataResponseTopic + prevTopicIDs map[string][16]byte + wantLeaders map[topicPartition]int32 + wantTopicIDs map[string][16]byte + wantNoLiveLeader map[string]struct{} + wantNoLeader map[topicPartition]struct{} + wantDropped leaderDrops }{ "happy path: single topic, all leaders known": { respTopics: []kmsg.MetadataResponseTopic{{ @@ -200,6 +204,7 @@ func TestBuildLeadersAndTopicIDs(t *testing.T) { {topic: "a", partition: 2}: 3, }, wantTopicIDs: map[string][16]byte{"a": idA}, + wantDropped: leaderDrops{Count: 1, Topic: "a", Partition: 1, NodeID: 99}, }, "topic absent from response is evicted (deletion is authoritative)": { respTopics: []kmsg.MetadataResponseTopic{}, @@ -224,15 +229,105 @@ func TestBuildLeadersAndTopicIDs(t *testing.T) { TopicID: idA, Partitions: []kmsg.MetadataResponseTopicPartition{{Partition: 0, Leader: 99}}, }}, + wantLeaders: map[topicPartition]int32{}, + wantTopicIDs: map[string][16]byte{"a": idA}, + wantNoLiveLeader: map[string]struct{}{"a": {}}, + wantDropped: leaderDrops{Count: 1, Topic: "a", Partition: 0, NodeID: 99}, + }, + "every partition leader unknown: first excluded leader is the sample": { + respTopics: []kmsg.MetadataResponseTopic{{ + Topic: stringPtr("a"), + TopicID: idA, + Partitions: []kmsg.MetadataResponseTopicPartition{ + {Partition: 0, Leader: 99}, + {Partition: 1, Leader: 98}, + }, + }}, + wantLeaders: map[topicPartition]int32{}, + wantTopicIDs: map[string][16]byte{"a": idA}, + wantNoLiveLeader: map[string]struct{}{"a": {}}, + wantDropped: leaderDrops{Count: 2, Topic: "a", Partition: 0, NodeID: 99}, + }, + "one topic fully excluded, another kept": { + respTopics: []kmsg.MetadataResponseTopic{ + {Topic: stringPtr("a"), TopicID: idA, Partitions: []kmsg.MetadataResponseTopicPartition{{Partition: 0, Leader: 99}}}, + {Topic: stringPtr("b"), TopicID: idB, Partitions: []kmsg.MetadataResponseTopicPartition{{Partition: 0, Leader: 1}}}, + }, + wantLeaders: map[topicPartition]int32{{topic: "b", partition: 0}: 1}, + wantTopicIDs: map[string][16]byte{"a": idA, "b": idB}, + wantNoLiveLeader: map[string]struct{}{"a": {}}, + wantDropped: leaderDrops{Count: 1, Topic: "a", Partition: 0, NodeID: 99}, + }, + "two topics both lose every leader: count sums and the sample stays on the first": { + respTopics: []kmsg.MetadataResponseTopic{ + {Topic: stringPtr("a"), TopicID: idA, Partitions: []kmsg.MetadataResponseTopicPartition{{Partition: 0, Leader: 99}}}, + {Topic: stringPtr("b"), TopicID: idB, Partitions: []kmsg.MetadataResponseTopicPartition{{Partition: 3, Leader: 98}}}, + }, + wantLeaders: map[topicPartition]int32{}, + wantTopicIDs: map[string][16]byte{"a": idA, "b": idB}, + wantNoLiveLeader: map[string]struct{}{"a": {}, "b": {}}, + wantDropped: leaderDrops{Count: 2, Topic: "a", Partition: 0, NodeID: 99}, + }, + "partition leader below zero is not a drop and not a fallback topic": { + respTopics: []kmsg.MetadataResponseTopic{{ + Topic: stringPtr("a"), + TopicID: idA, + Partitions: []kmsg.MetadataResponseTopicPartition{{ + Partition: 0, + Leader: -1, + ErrorCode: 5, // LEADER_NOT_AVAILABLE + }}, + }}, + wantLeaders: map[topicPartition]int32{}, + wantTopicIDs: map[string][16]byte{"a": idA}, + wantNoLeader: map[topicPartition]struct{}{{topic: "a", partition: 0}: {}}, + }, + "leader below zero does not count when a sibling names a missing node": { + respTopics: []kmsg.MetadataResponseTopic{{ + Topic: stringPtr("a"), + TopicID: idA, + Partitions: []kmsg.MetadataResponseTopicPartition{ + {Partition: 0, Leader: -1, ErrorCode: 5}, + {Partition: 1, Leader: 1}, + {Partition: 2, Leader: 99}, + }, + }}, + wantLeaders: map[topicPartition]int32{{topic: "a", partition: 1}: 1}, + wantTopicIDs: map[string][16]byte{"a": idA}, + wantNoLeader: map[topicPartition]struct{}{{topic: "a", partition: 0}: {}}, + wantDropped: leaderDrops{Count: 1, Topic: "a", Partition: 2, NodeID: 99}, + }, + "partition error that still names a live leader is kept": { + respTopics: []kmsg.MetadataResponseTopic{{ + Topic: stringPtr("a"), + TopicID: idA, + Partitions: []kmsg.MetadataResponseTopicPartition{{ + Partition: 0, + Leader: 1, + ErrorCode: 9, // REPLICA_NOT_AVAILABLE + }}, + }}, + wantLeaders: map[topicPartition]int32{{topic: "a", partition: 0}: 1}, + wantTopicIDs: map[string][16]byte{"a": idA}, + }, + "empty partition list with a zero error code stays out of the no-live-leader set": { + respTopics: []kmsg.MetadataResponseTopic{{ + Topic: stringPtr("a"), + TopicID: idA, + Partitions: nil, + }}, wantLeaders: map[topicPartition]int32{}, wantTopicIDs: map[string][16]byte{"a": idA}, }, } for name, tc := range tests { t.Run(name, func(t *testing.T) { - leaders, topicIDs := buildLeadersAndTopicIDs(tc.respTopics, knownAgents, tc.prevTopicIDs) + leaders, topicIDs, noLiveLeader, noLeader, dropped := buildLeadersAndTopicIDs(tc.respTopics, knownAgents, tc.prevTopicIDs) assert.Equal(t, tc.wantLeaders, leaders) assert.Equal(t, tc.wantTopicIDs, topicIDs) + assert.Equal(t, tc.wantNoLiveLeader, noLiveLeader) + assert.Equal(t, tc.wantNoLeader, noLeader) + assert.Equal(t, tc.wantDropped, dropped) }) } } diff --git a/pkg/wgo/client.go b/pkg/wgo/client.go index 7f2289d..4fed0c3 100644 --- a/pkg/wgo/client.go +++ b/pkg/wgo/client.go @@ -121,7 +121,8 @@ func NewWarpstreamClient(logger kgo.Logger, reg prometheus.Registerer, opts ...O m := newMetrics(reg) pool := NewAgentPool(kgoClient) - if _, err := pool.Refresh(context.Background()); err != nil { + _, dropped, err := pool.refresh(context.Background()) + if err != nil { kgoClient.Close() return nil, fmt.Errorf("initial agent pool refresh: %w", err) } @@ -150,6 +151,9 @@ func NewWarpstreamClient(logger kgo.Logger, reg prometheus.Registerer, opts ...O refreshCancel: refreshCancel, refreshNowCh: make(chan struct{}, 1), } + // Count a dropped leader and do not nudge: the refresh goroutine starts below. + // The periodic tick fetches the next snapshot. + c.noteLeaderDrops(dropped, false) // Demoter sits on top of the lazy pool strategy so refresh-driven // agent-pool changes flow through transparently while the Demoter's // per-agent probe-timing state persists across refreshes. @@ -416,7 +420,7 @@ func (c *WarpstreamClient) triggerRefresh() { // and leave the previous snapshot in place. func (c *WarpstreamClient) refreshPool(trigger metadataRefreshTrigger) { before := c.pool.Agents() - removed, err := c.pool.Refresh(c.refreshCtx) + removed, dropped, err := c.pool.refresh(c.refreshCtx) c.metrics.observeMetadataRefresh(trigger, before, c.pool.Agents(), err) if err != nil { log(c.logger, kgo.LogLevelWarn, "warpstream client metadata refresh failed", "err", err) @@ -426,6 +430,26 @@ func (c *WarpstreamClient) refreshPool(trigger metadataRefreshTrigger) { c.tracker.PurgeAgents(removed) } 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). + c.noteLeaderDrops(dropped, true) +} + +// noteLeaderDrops counts and logs excluded leaders. nudge asks for another fetch. +func (c *WarpstreamClient) noteLeaderDrops(dropped leaderDrops, nudge bool) { + if dropped.Count == 0 { + return + } + c.metrics.agentPoolLeaderDroppedTotal.Add(float64(dropped.Count)) + log(c.logger, kgo.LogLevelWarn, "warpstream agentpool: leaders excluded from map", + "count", dropped.Count, + "first_topic", dropped.Topic, + "first_partition", dropped.Partition, + "first_node_id", dropped.NodeID) + if nudge { + c.triggerRefresh() + } } // waitRefreshCooldown enforces the configured minimum between on-demand diff --git a/pkg/wgo/client_test.go b/pkg/wgo/client_test.go index ad8cd44..3d786c9 100644 --- a/pkg/wgo/client_test.go +++ b/pkg/wgo/client_test.go @@ -352,7 +352,7 @@ func TestWarpstreamClient_ProduceSync(t *testing.T) { topicIDs: map[string][16]byte{topic: topicID}, strategy: newDefaultPartitionAssignmentStrategy([]int32{leader}, map[topicPartition]int32{ {topic: topic, partition: 0}: leader, - }), + }, nil, nil), }) // Partition 0 (healthy sibling) and partition 1 (dropped leader) @@ -390,6 +390,74 @@ func TestWarpstreamClient_ProduceSync(t *testing.T) { }) }) + t.Run("topic whose every leader was dropped falls back; a sibling topic still routes", func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + const other = "other-topic" + c, _, clusterAddr, vnet := newTestWarpstreamClient(t, topic, 1) + + createReq := kmsg.NewPtrCreateTopicsRequest() + createReq.Topics = []kmsg.CreateTopicsRequestTopic{{ + Topic: other, + NumPartitions: 1, + ReplicationFactor: 1, + }} + _, err := c.Request(t.Context(), createReq) + require.NoError(t, err) + + // Past Refresh's one-nanosecond metadata cache so the new topic is visible. + time.Sleep(time.Nanosecond) + _, err = c.pool.Refresh(t.Context()) + require.NoError(t, err) + + leaderCands := c.demoter.Candidates(other, 0, 1) + require.Len(t, leaderCands, 1) + leader := leaderCands[0].NodeID + wipedID, ok := c.pool.TopicID(topic) + require.True(t, ok) + otherID, ok := c.pool.TopicID(other) + require.True(t, ok) + + // topic lost every leader, so it must hash onto the one live broker. + // other keeps its leader. A nil candidate for topic would fail both records. + c.pool.state.Store(&poolState{ + agents: []int32{leader}, + topicIDs: map[string][16]byte{topic: wipedID, other: otherID}, + strategy: newDefaultPartitionAssignmentStrategy([]int32{leader}, map[topicPartition]int32{ + {topic: other, partition: 0}: leader, + }, map[string]struct{}{topic: {}}, nil), + }) + + results := c.ProduceSync(t.Context(), []*kgo.Record{ + {Topic: topic, Partition: 0, Value: []byte("wiped"), Timestamp: time.Now()}, + {Topic: other, Partition: 0, Value: []byte("kept"), Timestamp: time.Now()}, + }) + require.Len(t, results, 2) + require.NoError(t, results[0].Err) + require.NoError(t, results[1].Err) + assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned))) + + consumer, err := kgo.NewClient( + kgo.SeedBrokers(clusterAddr), + kgo.Dialer(vnet.DialContext), + kgo.ConsumePartitions(map[string]map[int32]kgo.Offset{ + topic: {0: kgo.NewOffset().AtStart()}, + other: {0: kgo.NewOffset().AtStart()}, + }), + ) + require.NoError(t, err) + t.Cleanup(consumer.Close) + + fetches := consumer.PollFetches(t.Context()) + require.NoError(t, fetches.Err()) + require.Len(t, fetches.Records(), 2) + got := map[string]string{} + for _, r := range fetches.Records() { + got[r.Topic] = string(r.Value) + } + assert.Equal(t, map[string]string{topic: "wiped", other: "kept"}, got) + }) + }) + t.Run("partition split across flushes fails when a later flush fails it", func(t *testing.T) { // ProduceSync splits partition-A across flushes; one chunk shares a flush // with a Produce of partition-B. That flush lands A and fails B, and the @@ -1526,6 +1594,40 @@ func TestWarpstreamClient_IdleClusterStats(t *testing.T) { }) } +func TestWarpstreamClient_NoteLeaderDrops(t *testing.T) { + reg := prometheus.NewRegistry() + c := &WarpstreamClient{ + logger: nopLogger{}, + metrics: newMetrics(reg), + refreshNowCh: make(chan struct{}, 1), + } + + c.noteLeaderDrops(leaderDrops{}, true) + assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal)) + select { + case <-c.refreshNowCh: + t.Fatal("clean refresh nudged") + default: + } + + c.noteLeaderDrops(leaderDrops{Count: 2, Topic: "ingest", Partition: 35, NodeID: 99}, true) + assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal)) + select { + case <-c.refreshNowCh: + default: + t.Fatal("refresh that excluded leaders did not nudge") + } + + // Constructor path counts the drop and does not fill the nudge channel. + c.noteLeaderDrops(leaderDrops{Count: 1, Topic: "ingest", Partition: 0, NodeID: 7}, false) + assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal)) + select { + case <-c.refreshNowCh: + t.Fatal("constructor refresh nudged") + default: + } +} + func TestWarpstreamClient_WaitRefreshCooldown(t *testing.T) { t.Run("fetch time counts toward the interval", func(t *testing.T) { synctest.Test(t, func(t *testing.T) { diff --git a/pkg/wgo/config.go b/pkg/wgo/config.go index d16c4a8..7e6862a 100644 --- a/pkg/wgo/config.go +++ b/pkg/wgo/config.go @@ -72,7 +72,8 @@ type Config struct { MetadataRefreshInterval time.Duration // OnDemandMetadataRefreshInterval is the minimum time between the start of - // Metadata refreshes triggered by routing misses. Must be > 0 and no + // 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 time.Duration } @@ -336,7 +337,8 @@ func WithMetadataRefreshInterval(d time.Duration) Opt { } // WithOnDemandMetadataRefreshInterval sets the minimum time between the start -// of Metadata refreshes triggered by routing misses. It must be no greater +// 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. func WithOnDemandMetadataRefreshInterval(d time.Duration) Opt { return opt{func(c *Config) { c.OnDemandMetadataRefreshInterval = d }} diff --git a/pkg/wgo/demoter_test.go b/pkg/wgo/demoter_test.go index 294fe23..32c3271 100644 --- a/pkg/wgo/demoter_test.go +++ b/pkg/wgo/demoter_test.go @@ -862,7 +862,7 @@ func BenchmarkDemoter_Candidates(b *testing.B) { for p := int32(0); p < sc.numPartitions; p++ { leaders[topicPartition{topic, p}] = p % sc.numAgents } - inner := newDefaultPartitionAssignmentStrategy(agents, leaders) + inner := newDefaultPartitionAssignmentStrategy(agents, leaders, nil, nil) d, _ := newTestDemoter(inner, tracker, health, cfg) diff --git a/pkg/wgo/metrics.go b/pkg/wgo/metrics.go index 798d89d..55e74d8 100644 --- a/pkg/wgo/metrics.go +++ b/pkg/wgo/metrics.go @@ -50,6 +50,7 @@ type metrics struct { produceRecordsFailedTotal prometheus.Counter produceRecordsRejectedTotal *prometheus.CounterVec + agentPoolLeaderDroppedTotal prometheus.Counter metadataRefreshResultsTotal *prometheus.CounterVec clusterStatsAvailable prometheus.Gauge @@ -236,6 +237,10 @@ func newMetrics(reg prometheus.Registerer) *metrics { Name: "warpstream_produce_records_rejected_total", Help: "Total number of records rejected by the client before any wire dispatch, by reason (record_too_large, no_agent_assigned).", }, []string{"reason"}), + agentPoolLeaderDroppedTotal: promauto.With(reg).NewCounter(prometheus.CounterOpts{ + Name: "warpstream_agentpool_leader_dropped_total", + Help: "Partition leaders excluded from the assignment map because their NodeID was absent from that Metadata response's broker list. One increment per excluded leader, including the constructor refresh. Topic-level Metadata errors are not counted. A partition Leader below 0 is not counted.", + }), metadataRefreshResultsTotal: promauto.With(reg).NewCounterVec(prometheus.CounterOpts{ Name: "warpstream_metadata_refresh_results_total", Help: "Total number of live AgentPool Metadata refreshes, by trigger (periodic, on_demand) and result (membership_changed, unchanged, failed). membership_changed is the sorted Agent NodeID set only; leader-only or topic-only updates are unchanged. The constructor Refresh is not counted.", diff --git a/pkg/wgo/partition_assignment.go b/pkg/wgo/partition_assignment.go index c7a2fc1..c324997 100644 --- a/pkg/wgo/partition_assignment.go +++ b/pkg/wgo/partition_assignment.go @@ -108,23 +108,27 @@ type DefaultPartitionAssignmentStrategy struct { agents []int32 // sorted ascending, snapshot at construction leaders map[topicPartition]int32 - // knownTopics is every topic with at least one entry in leaders. + // knownTopics is every topic with at least one entry in leaders, plus every + // topic whose partitions this refresh listed and then excluded entirely. // Candidates uses it to tell "this topic is known, one partition's // leader is just missing" (safe to fall back to another agent) apart // from "this topic is unknown to Metadata" (falling back would hide a // topic that still needs an on-demand refresh). // - // Edge case: if every partition of a small topic loses its leader in - // the same refresh, that topic looks unknown here too for that one - // refresh, and briefly loses the fallback. Rare, and not seen in the - // incidents this fix is based on, so left as-is for now. + // A topic Metadata has never returned, and a topic-level error, stay out. + // A partition whose Leader was below 0 stays out of the fallback too, + // even when this topic is known through some other partition. knownTopics map[string]struct{} + noLeader map[topicPartition]struct{} } -func newDefaultPartitionAssignmentStrategy(agents []int32, leaders map[topicPartition]int32) *DefaultPartitionAssignmentStrategy { +func newDefaultPartitionAssignmentStrategy(agents []int32, leaders map[topicPartition]int32, topicsWithNoLiveLeader map[string]struct{}, noLeader map[topicPartition]struct{}) *DefaultPartitionAssignmentStrategy { // Built from empty, not sized off leaders: there are far fewer // distinct topics than partitions. knownTopics := make(map[string]struct{}) + for topic := range topicsWithNoLiveLeader { + knownTopics[topic] = struct{}{} + } for tp := range leaders { knownTopics[tp.topic] = struct{}{} } @@ -132,6 +136,7 @@ func newDefaultPartitionAssignmentStrategy(agents []int32, leaders map[topicPart agents: agents, leaders: leaders, knownTopics: knownTopics, + noLeader: noLeader, } } @@ -144,8 +149,10 @@ func newDefaultPartitionAssignmentStrategy(agents []int32, leaders map[topicPart // this falls back to a deterministic pick from the live agent set instead // of returning no candidates. Any live agent can serve any partition, so // this is a safe guess while the real leader is still unclear. A topic -// Metadata has never returned is left alone, so it still gets an -// on-demand refresh. +// whose partitions this refresh listed and then excluded entirely counts +// as known. A topic Metadata has never returned is left alone, so it still +// gets an on-demand refresh. A partition whose Leader was below 0 returns +// nil: WarpStream named no agent, so this does not pick one. // // Caveat: two clients that refreshed at different times can pick // different fallback agents for the same partition — a real leader @@ -157,8 +164,12 @@ func (s *DefaultPartitionAssignmentStrategy) Candidates(topic string, partition } var h uint64 - leader, ok := s.leaders[topicPartition{topic: topic, partition: partition}] + tp := topicPartition{topic: topic, partition: partition} + leader, ok := s.leaders[tp] if !ok { + if _, unnamed := s.noLeader[tp]; unnamed { + return nil + } if len(s.agents) == 0 { return nil } diff --git a/pkg/wgo/partition_assignment_test.go b/pkg/wgo/partition_assignment_test.go index 5413a35..202e510 100644 --- a/pkg/wgo/partition_assignment_test.go +++ b/pkg/wgo/partition_assignment_test.go @@ -51,7 +51,7 @@ func TestDefaultPartitionAssignmentStrategy_SecondaryAvailability(t *testing.T) t.Run(name, func(t *testing.T) { s := newDefaultPartitionAssignmentStrategy(tc.all, map[topicPartition]int32{ {topic: tc.topic, partition: tc.partition}: tc.primary, - }) + }, nil, nil) nodeID, ok := secondaryOf(s, tc.topic, tc.partition) assert.Equal(t, tc.wantOk, ok) if ok { @@ -70,7 +70,7 @@ func TestDefaultPartitionAssignmentStrategy_SecondaryNeverReturnsPrimary(t *test t.Run(fmt.Sprintf("primary=%d partition=%d", primary, part), func(t *testing.T) { s := newDefaultPartitionAssignmentStrategy(agents, map[topicPartition]int32{ {topic: "t", partition: part}: primary, - }) + }, nil, nil) nodeID, ok := secondaryOf(s, "t", part) require.True(t, ok) assert.NotEqual(t, primary, nodeID) @@ -83,7 +83,7 @@ func TestDefaultPartitionAssignmentStrategy_NodeIDZeroIsValid(t *testing.T) { t.Run("NodeID 0 can be the secondary", func(t *testing.T) { s := newDefaultPartitionAssignmentStrategy([]int32{0, 1}, map[topicPartition]int32{ {topic: "t", partition: 0}: 1, - }) + }, nil, nil) nodeID, ok := secondaryOf(s, "t", 0) require.True(t, ok) assert.Equal(t, int32(0), nodeID) @@ -92,7 +92,7 @@ func TestDefaultPartitionAssignmentStrategy_NodeIDZeroIsValid(t *testing.T) { t.Run("NodeID 0 can be the primary without aliasing no-secondary sentinel", func(t *testing.T) { s := newDefaultPartitionAssignmentStrategy([]int32{0, 1}, map[topicPartition]int32{ {topic: "t", partition: 0}: 0, - }) + }, nil, nil) nodeID, ok := secondaryOf(s, "t", 0) require.True(t, ok) assert.Equal(t, int32(1), nodeID) @@ -104,7 +104,7 @@ func TestDefaultPartitionAssignmentStrategy_Candidates(t *testing.T) { leaders := map[topicPartition]int32{ {topic: "t", partition: 0}: 1, } - s := newDefaultPartitionAssignmentStrategy(agents, leaders) + s := newDefaultPartitionAssignmentStrategy(agents, leaders, nil, nil) t.Run("returns up to maxCandidates entries", func(t *testing.T) { c := s.Candidates("t", 0, 3) @@ -121,7 +121,7 @@ func TestDefaultPartitionAssignmentStrategy_Candidates(t *testing.T) { }) t.Run("returns nil for unknown partition when no agents are live", func(t *testing.T) { - empty := newDefaultPartitionAssignmentStrategy(nil, leaders) + empty := newDefaultPartitionAssignmentStrategy(nil, leaders, nil, nil) assert.Nil(t, empty.Candidates("t", 99, 3)) }) @@ -155,7 +155,7 @@ func TestDefaultPartitionAssignmentStrategy_Candidates(t *testing.T) { t.Run("single-agent cluster has no secondary candidate", func(t *testing.T) { s1 := newDefaultPartitionAssignmentStrategy([]int32{1}, map[topicPartition]int32{ {topic: "t", partition: 0}: 1, - }) + }, nil, nil) c := s1.Candidates("t", 0, 5) assert.Len(t, c, 1) assert.Equal(t, int32(1), c[0].NodeID) @@ -168,7 +168,7 @@ func TestDefaultPartitionAssignmentStrategy_FallbackOnUnknownLeader(t *testing.T // looked up below takes the fallback path, not "topic unknown". s := newDefaultPartitionAssignmentStrategy(agents, map[topicPartition]int32{ {topic: "t", partition: -1}: agents[0], - }) + }, nil, nil) t.Run("fallback pick is deterministic across calls and instances", func(t *testing.T) { c1 := s.Candidates("t", 7, 1) @@ -178,7 +178,7 @@ func TestDefaultPartitionAssignmentStrategy_FallbackOnUnknownLeader(t *testing.T s2 := newDefaultPartitionAssignmentStrategy(agents, map[topicPartition]int32{ {topic: "t", partition: -1}: agents[0], - }) + }, nil, nil) c3 := s2.Candidates("t", 7, 1) // Two independently constructed strategies over the same agent set must agree. assert.Equal(t, c1, c3) @@ -227,7 +227,7 @@ func TestDefaultPartitionAssignmentStrategy_FallbackOnUnknownLeader(t *testing.T t.Run("single live agent: fallback returns it with no secondary", func(t *testing.T) { single := newDefaultPartitionAssignmentStrategy([]int32{1}, map[topicPartition]int32{ {topic: "t", partition: -1}: 1, - }) + }, nil, nil) c := single.Candidates("t", 0, 5) require.Len(t, c, 1) assert.Equal(t, int32(1), c[0].NodeID) @@ -236,7 +236,7 @@ func TestDefaultPartitionAssignmentStrategy_FallbackOnUnknownLeader(t *testing.T t.Run("no live agents at all: still returns nil, not a zero-value agent", func(t *testing.T) { empty := newDefaultPartitionAssignmentStrategy(nil, map[topicPartition]int32{ {topic: "t", partition: -1}: 1, - }) + }, nil, nil) assert.Nil(t, empty.Candidates("t", 0, 3)) }) @@ -245,14 +245,14 @@ func TestDefaultPartitionAssignmentStrategy_FallbackOnUnknownLeader(t *testing.T // unknown, not just missing one leader. No fallback here. unknownTopic := newDefaultPartitionAssignmentStrategy(agents, map[topicPartition]int32{ {topic: "other-topic", partition: 0}: agents[0], - }) + }, nil, nil) assert.Nil(t, unknownTopic.Candidates("t", 0, 3)) }) t.Run("mixed: some partitions have a known leader, others fall back, in the same strategy", func(t *testing.T) { mixed := newDefaultPartitionAssignmentStrategy(agents, map[topicPartition]int32{ {topic: "t", partition: 0}: 3, - }) + }, nil, nil) known := mixed.Candidates("t", 0, 1) require.Len(t, known, 1) assert.Equal(t, int32(3), known[0].NodeID) @@ -261,6 +261,29 @@ func TestDefaultPartitionAssignmentStrategy_FallbackOnUnknownLeader(t *testing.T require.Len(t, fallback, 1) assert.Contains(t, agents, fallback[0].NodeID) }) + + t.Run("topic with no remaining leaders still falls back when named in the set", func(t *testing.T) { + wiped := newDefaultPartitionAssignmentStrategy(agents, nil, map[string]struct{}{"t": {}}, nil) + for p := int32(0); p < 20; p++ { + c := wiped.Candidates("t", p, 1) + require.Len(t, c, 1) + want := agents[hashTopicPartition("t", p)%uint64(len(agents))] + assert.Equal(t, want, c[0].NodeID) + } + assert.Nil(t, wiped.Candidates("absent", 0, 1)) + }) + + t.Run("partition with no leader returns nil even when the topic is known", func(t *testing.T) { + s := newDefaultPartitionAssignmentStrategy(agents, map[topicPartition]int32{ + {topic: "t", partition: 0}: 3, + }, nil, map[topicPartition]struct{}{ + {topic: "t", partition: 1}: {}, + }) + known := s.Candidates("t", 0, 1) + require.Len(t, known, 1) + assert.Equal(t, int32(3), known[0].NodeID) + assert.Nil(t, s.Candidates("t", 1, 1)) + }) } func TestDefaultPartitionAssignmentStrategy_PrimaryAndSecondaryViaCandidates(t *testing.T) { @@ -269,7 +292,7 @@ func TestDefaultPartitionAssignmentStrategy_PrimaryAndSecondaryViaCandidates(t * {topic: "t", partition: 0}: 1, {topic: "t", partition: 1}: 2, } - s := newDefaultPartitionAssignmentStrategy(agents, leaders) + s := newDefaultPartitionAssignmentStrategy(agents, leaders, nil, nil) t.Run("Candidates[0] is the leader", func(t *testing.T) { c := s.Candidates("t", 0, 1) @@ -299,8 +322,8 @@ func TestDefaultPartitionAssignmentStrategy_PrimaryAndSecondaryViaCandidates(t * }) t.Run("each Refresh creates a fresh strategy with its own cache", func(t *testing.T) { - s1 := newDefaultPartitionAssignmentStrategy(agents, leaders) - s2 := newDefaultPartitionAssignmentStrategy(agents, leaders) + s1 := newDefaultPartitionAssignmentStrategy(agents, leaders, nil, nil) + s2 := newDefaultPartitionAssignmentStrategy(agents, leaders, nil, nil) // Both strategies should produce the same deterministic result independently. c1 := s1.Candidates("t", 0, 2) c2 := s2.Candidates("t", 0, 2) @@ -313,7 +336,7 @@ func BenchmarkDefaultPartitionAssignmentStrategy_Candidates(b *testing.B) { all := makeNodeIDs(n) s := newDefaultPartitionAssignmentStrategy(all, map[topicPartition]int32{ {topic: "ingest", partition: 7}: all[0], - }) + }, nil, nil) b.Run(fmt.Sprintf("agents=%d", n), func(b *testing.B) { b.ReportAllocs() for range b.N { @@ -362,13 +385,10 @@ func BenchmarkNewDefaultPartitionAssignmentStrategy(b *testing.B) { } b.ResetTimer() b.ReportAllocs() - // Assigning to a sink outside the loop, not to _, matters here: - // discarding the result lets the compiler prove it never - // escapes and skip the allocations entirely, understating the - // real cost — AgentPool.Refresh always stores this behind an - // atomic.Pointer, which does escape. + // A package-level sink keeps these allocations in the reported count. + // Discarding the result lets the compiler elide them. for range b.N { - sink = newDefaultPartitionAssignmentStrategy(agents, leaders) + sink = newDefaultPartitionAssignmentStrategy(agents, leaders, nil, nil) } }) }