diff --git a/pkg/wgo/client.go b/pkg/wgo/client.go index 63387fb..e4b195a 100644 --- a/pkg/wgo/client.go +++ b/pkg/wgo/client.go @@ -539,12 +539,20 @@ func (c *WarpstreamClient) waitRefreshCooldown(delay, elapsed time.Duration) { type rejectedTopicPartitionRecords struct { topicPartitionRecords err error + // miss is why the lookup found no agent, for the routing-miss counter. + miss routingMissReason +} + +// routedGroup is a routed partition group plus how its agent was chosen. +type routedGroup struct { + promised[routedTopicPartitionRecords] + outcome routeOutcome } // routeRecords routes each partition once. A partition with no candidate is // returned unsent. The first miss requests a metadata refresh. func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(groupRecords []*kgo.Record) func(ProduceResult)) ([]promised[routedTopicPartitionRecords], []rejectedTopicPartitionRecords) { - groups := make(map[topicPartition]*promised[routedTopicPartitionRecords]) + groups := make(map[topicPartition]*routedGroup) order := make([]topicPartition, 0) var rejectedByKey map[topicPartition]int var rejected []rejectedTopicPartitionRecords @@ -562,7 +570,7 @@ func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(grou } } - cands := c.demoter.Candidates(r.Topic, r.Partition, 1) + cands, outcome := c.demoter.candidatesWithRoute(r.Topic, r.Partition, 1) if len(cands) == 0 { if rejectedByKey == nil { rejectedByKey = make(map[topicPartition]int) @@ -575,30 +583,54 @@ func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(grou partition: r.Partition, records: []*kgo.Record{r}, }, - err: fmt.Errorf("no agent assigned for topic %q partition %d", r.Topic, r.Partition), + err: fmt.Errorf("no agent assigned for topic %q partition %d", r.Topic, r.Partition), + miss: outcome.missReason(), }) continue } - groups[key] = &promised[routedTopicPartitionRecords]{ - item: routedTopicPartitionRecords{ - topicPartitionRecords: topicPartitionRecords{ - topic: r.Topic, - partition: r.Partition, - records: []*kgo.Record{r}, + groups[key] = &routedGroup{ + promised: promised[routedTopicPartitionRecords]{ + item: routedTopicPartitionRecords{ + topicPartitionRecords: topicPartitionRecords{ + topic: r.Topic, + partition: r.Partition, + records: []*kgo.Record{r}, + }, + nodeID: cands[0].NodeID, + nodeState: cands[0].State, }, - nodeID: cands[0].NodeID, - nodeState: cands[0].State, }, + outcome: outcome, } order = append(order, key) } + var ( + routes [routeSourceCount]int + misses [routingMissCount]int + ) out := make([]promised[routedTopicPartitionRecords], 0, len(order)) for _, key := range order { g := groups[key] g.done = doneFor(g.item.records) - out = append(out, *g) + out = append(out, g.promised) + if src, ok := g.outcome.routeSource(); ok { + routes[src] += len(g.item.records) + } + } + for i := range rejected { + misses[rejected[i].miss] += len(rejected[i].records) + } + for src, n := range routes { + if n > 0 { + c.metrics.partitionRoutes[src].Add(float64(n)) + } + } + for reason, n := range misses { + if n > 0 { + c.metrics.routingMisses[reason].Add(float64(n)) + } } return out, rejected } @@ -607,11 +639,15 @@ func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(grou // skips the per-partition map and avoids heap-allocating an intermediate // group. Returns an error if the record's partition has no known candidate. func (c *WarpstreamClient) routeRecord(record *kgo.Record, done func(ProduceResult)) (promised[routedTopicPartitionRecords], error) { - cands := c.demoter.Candidates(record.Topic, record.Partition, 1) + cands, outcome := c.demoter.candidatesWithRoute(record.Topic, record.Partition, 1) if len(cands) == 0 { + c.metrics.routingMisses[outcome.missReason()].Inc() c.triggerRefresh() return promised[routedTopicPartitionRecords]{}, fmt.Errorf("no agent assigned for topic %q partition %d", record.Topic, record.Partition) } + if src, ok := outcome.routeSource(); ok { + c.metrics.partitionRoutes[src].Inc() + } return promised[routedTopicPartitionRecords]{ item: routedTopicPartitionRecords{ topicPartitionRecords: topicPartitionRecords{ diff --git a/pkg/wgo/client_test.go b/pkg/wgo/client_test.go index 71931ad..0c0a1a0 100644 --- a/pkg/wgo/client_test.go +++ b/pkg/wgo/client_test.go @@ -2653,3 +2653,69 @@ func (b *lockedBuffer) String() string { defer b.mu.Unlock() return b.buf.String() } + +func TestWarpstreamClient_RoutingMetricsReconcile(t *testing.T) { + const topic = "test-topic" + synctest.Test(t, func(t *testing.T) { + c, _, _, _ := newTestWarpstreamClient(t, topic, 4) + // The single fake broker is node 0. Partition 0 names it as leader, + // partition 1 has no entry (stand-in), partition 2 is unnamed and + // partitions past 3 do not exist. + live := c.pool.state.Load() + c.pool.state.Store(&poolState{ + agents: live.agents, + topicIDs: live.topicIDs, + strategy: newDefaultPartitionAssignmentStrategy([]int32{0}, + map[topicPartition]int32{{topic: topic, partition: 0}: 0}, + nil, + map[topicPartition]struct{}{{topic: topic, partition: 2}: {}}, + map[string]int32{topic: 4}), + }) + + rec := func(topic string, partition int32, size int) *kgo.Record { + return &kgo.Record{Topic: topic, Partition: partition, Value: make([]byte, size), Timestamp: time.Now()} + } + results := c.ProduceSync(t.Context(), []*kgo.Record{ + rec(topic, 0, 1), rec(topic, 0, 1), + rec(topic, 1, 1), rec(topic, 1, 1), rec(topic, 1, 1), + rec(topic, 2, 1), + rec(topic, 9, 1), rec(topic, 9, 1), + rec("other", 0, 1), + rec(topic, 0, 2<<20), // larger than the test client's 1<<20 BatchMaxBytes + }) + require.Len(t, results, 10) + for _, i := range []int{0, 1, 2, 3, 4} { + assert.NoError(t, results[i].Err) + } + + var ( + routes = testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceLeader]) + testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceStandIn]) + misses float64 + ) + for _, m := range c.metrics.routingMisses { + misses += testutil.ToFloat64(m) + } + tooLarge := testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedRecordTooLarge)) + noAgent := testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned)) + + assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceLeader])) + assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceStandIn])) + assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.routingMisses[routingMissNoLeader])) + assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.routingMisses[routingMissPartitionOutOfRange])) + assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.routingMisses[routingMissUnknownTopic])) + // Every submitted record is exactly one of: too large, routed, or a miss. + assert.Equal(t, testutil.ToFloat64(c.metrics.produceRecordsTotal), tooLarge+routes+misses) + assert.Equal(t, noAgent, misses) + + // Produce counts one record per call. The rejected batch above nudged a + // refresh that may have replaced the injected strategy, so only the + // total is stable here. + var promised sync.WaitGroup + promised.Add(2) + c.Produce(t.Context(), rec(topic, 1, 1), func(*kgo.Record, error) { promised.Done() }) + c.Produce(t.Context(), rec(topic, 9, 1), func(*kgo.Record, error) { promised.Done() }) + promised.Wait() + assert.Equal(t, float64(6), testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceLeader])+testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceStandIn])) + assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.routingMisses[routingMissPartitionOutOfRange])) + }) +} diff --git a/pkg/wgo/demoter.go b/pkg/wgo/demoter.go index 7cb64cc..6e2c710 100644 --- a/pkg/wgo/demoter.go +++ b/pkg/wgo/demoter.go @@ -144,8 +144,16 @@ func NewDemoter(inner PartitionAssignmentStrategy, tracker AgentStatsReader, hea // demoted and none is due for a probe, the natural primary is surfaced as a // forced probe so we never refuse to route. func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) []Agent { + agents, _ := d.candidatesWithRoute(topic, partition, maxCandidates) + return agents +} + +// candidatesWithRoute is Candidates plus the inner strategy's classification of +// the lookup that produced the returned agents. Demotion only reorders or elides +// agents, so a demoted leader replaced here keeps that classification. +func (d *Demoter) candidatesWithRoute(topic string, partition int32, maxCandidates int) ([]Agent, routeOutcome) { if maxCandidates <= 0 { - return nil + return nil, routeUnclassified } // Use the shared HealthCheckConfig parameters so the Demoter and the // Hedger hit the same ClusterStats cache entry. cluster.SlowThreshold @@ -155,7 +163,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [ now := d.now() clusterStats, hasClusterStats := d.tracker.ClusterStats(now, d.healthCfg.SlowMultiplier, d.healthCfg.FaultyThreshold) if suppressed, _ := d.isDemotionSuppressed(clusterStats, hasClusterStats); suppressed { - return d.inner.Candidates(topic, partition, maxCandidates) + return candidatesOf(d.inner, topic, partition, maxCandidates) } // The returned list holds up to maxCandidates entries total @@ -170,6 +178,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [ const maxRetries = 6 var ( agents []Agent + route routeOutcome noDemotedAgents bool ) for retry, extra := 0, 2; retry < maxRetries; retry, extra = retry+1, extra*2 { @@ -178,7 +187,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [ nonDemoted = 0 ) - agents = d.inner.Candidates(topic, partition, asked) + agents, route = candidatesOf(d.inner, topic, partition, asked) noDemotedAgents = true for _, c := range agents { @@ -199,10 +208,10 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [ } } if len(agents) == 0 { - return nil + return nil, route } if noDemotedAgents { - return agents[:min(len(agents), maxCandidates)] + return agents[:min(len(agents), maxCandidates)], route } candidates := make([]Agent, 0, maxCandidates) @@ -239,7 +248,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [ candidates = append(candidates, forced) } - return candidates + return candidates, route } // isDemoted reports whether agent nodeID currently meets the demotion criteria. diff --git a/pkg/wgo/demoter_test.go b/pkg/wgo/demoter_test.go index 9b35408..8f98614 100644 --- a/pkg/wgo/demoter_test.go +++ b/pkg/wgo/demoter_test.go @@ -1042,3 +1042,158 @@ func TestDemoter_Refresh(t *testing.T) { assert.True(t, kept3) assert.Len(t, d.lastDemotedProbe, 2) } + +// routeStub is an inner strategy that classifies its answers. outcome receives +// the 1-based call number so a test can change the classification between the +// Demoter's lookups. +type routeStub struct { + agents []Agent + outcome func(call int) routeOutcome + calls int +} + +func (s *routeStub) Candidates(topic string, partition int32, maxCandidates int) []Agent { + agents, _ := s.candidatesWithRoute(topic, partition, maxCandidates) + return agents +} + +func (s *routeStub) candidatesWithRoute(_ string, _ int32, maxCandidates int) ([]Agent, routeOutcome) { + s.calls++ + return s.agents[:min(len(s.agents), maxCandidates)], s.outcome(s.calls) +} + +func TestDemoter_CandidatesWithRoute(t *testing.T) { + const topic = "t" + health := HealthCheckConfig{ + SlowMultiplier: 2.0, + MaxSlowFraction: 0.3, + FaultyThreshold: 0.05, + MaxFaultyFraction: 0.3, + } + cfg := DemoterConfig{ProbeInterval: time.Second} + + t.Run("keeps the inner classification", func(t *testing.T) { + inner := newDefaultPartitionAssignmentStrategy([]int32{1, 2}, map[topicPartition]int32{ + {topic: topic, partition: 0}: 1, + }, nil, nil, map[string]int32{topic: 2}) + d, _ := newTestDemoter(inner, noopAgentStatsTracker{}, HealthCheckConfig{}, DemoterConfig{}) + + got, outcome := d.candidatesWithRoute(topic, 0, 1) + require.Len(t, got, 1) + assert.Equal(t, routeLeader, outcome) + + _, outcome = d.candidatesWithRoute(topic, 1, 1) + assert.Equal(t, routeStandIn, outcome) + + got, outcome = d.candidatesWithRoute(topic, 5, 1) + assert.Empty(t, got) + assert.Equal(t, routeMissOutOfRange, outcome) + }) + + t.Run("a custom inner strategy is unclassified", func(t *testing.T) { + custom := &mockPartitionAssignmentStrategy{candidates: map[partitionKey][]Agent{{topic, 0}: healthyAgents(9)}} + d, _ := newTestDemoter(custom, noopAgentStatsTracker{}, HealthCheckConfig{}, DemoterConfig{}) + got, outcome := d.candidatesWithRoute(topic, 0, 1) + assert.Equal(t, healthyAgents(9), got) + assert.Equal(t, routeUnclassified, outcome) + }) + + t.Run("Candidates is unchanged", func(t *testing.T) { + inner := newDefaultPartitionAssignmentStrategy([]int32{1, 2}, map[topicPartition]int32{ + {topic: topic, partition: 0}: 1, + }, nil, nil, map[string]int32{topic: 2}) + d, _ := newTestDemoter(inner, noopAgentStatsTracker{}, HealthCheckConfig{}, DemoterConfig{}) + want, _ := d.candidatesWithRoute(topic, 0, 2) + assert.Equal(t, want, d.Candidates(topic, 0, 2)) + }) + + t.Run("a demoted leader replaced by an alternate keeps the classification", func(t *testing.T) { + const leader, healthy = int32(2), int32(1) + inner := newDefaultPartitionAssignmentStrategy([]int32{healthy, leader, 3}, map[topicPartition]int32{ + {topic: topic, partition: 0}: leader, + }, nil, nil, map[string]int32{topic: 1}) + tr := NewAverageAgentStatsTracker() + nowNs := time.Now().UnixNano() + for _, id := range []int32{healthy, 3} { + seedFullWindow(tr, id, nowNs, 20, 10, 0) + } + seedFullWindow(tr, leader, nowNs, 10, 10, 10) + d, _ := newTestDemoter(inner, tr, health, cfg) + now := time.Now() + d.now = func() time.Time { return now } + + // The first call spends the leader's probe slot. + got, outcome := d.candidatesWithRoute(topic, 0, 1) + require.Len(t, got, 1) + assert.Equal(t, leader, got[0].NodeID) + assert.Equal(t, AgentStateDemoted, got[0].State) + assert.Equal(t, routeLeader, outcome) + + // Within the probe interval the leader is skipped for an alternate. + got, outcome = d.candidatesWithRoute(topic, 0, 1) + require.Len(t, got, 1) + assert.NotEqual(t, leader, got[0].NodeID) + assert.Equal(t, routeLeader, outcome) + }) + + t.Run("a forced probe keeps the classification", func(t *testing.T) { + const slow = int32(2) + inner := &routeStub{ + agents: healthyAgents(slow), + outcome: func(int) routeOutcome { return routeStandIn }, + } + tr := NewAverageAgentStatsTracker() + nowNs := time.Now().UnixNano() + seedFullWindow(tr, 1, nowNs, 20, 10, 0) + seedFullWindow(tr, 3, nowNs, 20, 10, 0) + seedFullWindow(tr, slow, nowNs, 10, 10, 10) + d, _ := newTestDemoter(inner, tr, health, cfg) + now := time.Now() + d.now = func() time.Time { return now } + + _, _ = d.candidatesWithRoute(topic, 0, 1) // spends the probe slot + + got, outcome := d.candidatesWithRoute(topic, 0, 1) + require.Len(t, got, 1) + assert.Equal(t, slow, got[0].NodeID) + assert.Equal(t, AgentStateDemoted, got[0].State) + assert.Equal(t, routeStandIn, outcome) + }) + + t.Run("the classification comes from the lookup that produced the agents", func(t *testing.T) { + // Ten demoted agents ahead of one healthy agent make the Demoter widen + // its ask. A refresh between the lookups changes the answer, so the + // outcome of the last lookup is the one that matches the returned agents. + // A quarter of the agents are faulty, below the cluster-wide guard. + faulty := []int32{10, 11, 12, 13, 14, 15, 16, 17, 18, 19} + var healthy []int32 + for id := int32(100); id < 130; id++ { + healthy = append(healthy, id) + } + inner := &routeStub{ + agents: healthyAgents(append(append([]int32{}, faulty...), healthy...)...), + outcome: func(call int) routeOutcome { + if call == 1 { + return routeLeader + } + return routeStandIn + }, + } + tr := NewAverageAgentStatsTracker() + nowNs := time.Now().UnixNano() + for _, id := range healthy { + seedFullWindow(tr, id, nowNs, 20, 10, 0) + } + for _, id := range faulty { + seedFullWindow(tr, id, nowNs, 20, 10, 10) + } + d, _ := newTestDemoter(inner, tr, health, cfg) + now := time.Now() + d.now = func() time.Time { return now } + + got, outcome := d.candidatesWithRoute(topic, 0, 1) + require.Len(t, got, 1) + require.Greater(t, inner.calls, 1) + assert.Equal(t, routeStandIn, outcome) + }) +} diff --git a/pkg/wgo/metrics.go b/pkg/wgo/metrics.go index 1e55e36..9411fbe 100644 --- a/pkg/wgo/metrics.go +++ b/pkg/wgo/metrics.go @@ -57,6 +57,9 @@ type metrics struct { produceRecordsFailedTotal prometheus.Counter produceRecordsRejectedTotal *prometheus.CounterVec + partitionRoutes [routeSourceCount]prometheus.Counter + routingMisses [routingMissCount]prometheus.Counter + agentPoolExcludedLeaders prometheus.Gauge metadataRefreshResultsTotal *prometheus.CounterVec @@ -89,6 +92,68 @@ const ( produceRejectedNoAgentAssigned = "no_agent_assigned" ) +// routeSource is how the strategy chose the primary agent for a routed record. +type routeSource int8 + +const ( + routeSourceLeader routeSource = iota + routeSourceStandIn + routeSourceCount +) + +const ( + routeSourceLabelLeader = "leader" + routeSourceLabelStandIn = "stand_in" +) + +// routingMissReason is why an initial route found no agent. +type routingMissReason int8 + +// The zero value is other, so an unset miss is not counted as an empty pool. +const ( + routingMissOther routingMissReason = iota + routingMissEmptyPool + routingMissUnknownTopic + routingMissNoLeader + routingMissPartitionOutOfRange + routingMissCount +) + +const ( + routingMissLabelEmptyPool = "empty_pool" + routingMissLabelUnknownTopic = "unknown_topic" + routingMissLabelNoLeader = "no_leader" + routingMissLabelPartitionOutOfRange = "partition_out_of_range" + routingMissLabelOther = "other" +) + +// routeSource reports the counter an accepted route belongs to. A strategy that +// cannot classify its answer is not counted. +func (o routeOutcome) routeSource() (routeSource, bool) { + switch o { + case routeLeader: + return routeSourceLeader, true + case routeStandIn: + return routeSourceStandIn, true + } + return 0, false +} + +// missReason is the reason a lookup that returned no agent is counted under. +func (o routeOutcome) missReason() routingMissReason { + switch o { + case routeMissEmptyPool: + return routingMissEmptyPool + case routeMissUnknownTopic: + return routingMissUnknownTopic + case routeMissNoLeader: + return routingMissNoLeader + case routeMissOutOfRange: + return routingMissPartitionOutOfRange + } + return routingMissOther +} + const ( agentStateHealthy = "healthy" agentStateDemoted = "demoted" @@ -231,6 +296,16 @@ func newMetrics(reg prometheus.Registerer) *metrics { NativeHistogramMinResetDuration: time.Hour, }, []string{"outcome"}) + partitionRoutes := promauto.With(reg).NewCounterVec(prometheus.CounterOpts{ + Name: "warpstream_partition_routes_total", + Help: "Input records routed to an agent, by how the strategy chose it: leader (the partition's named leader is in the snapshot) or stand_in (the leader entry was missing for a known topic, so a live agent was picked). Counted once per record at the initial routing decision, not per hedge, retry or flush. A demoted leader replaced by the Demoter keeps the classification of the lookup. stand_in covers every missing leader entry, not only an excluded leader.", + }, []string{"source"}) + + routingMisses := promauto.With(reg).NewCounterVec(prometheus.CounterOpts{ + Name: "warpstream_routing_misses_total", + Help: "Input records rejected because the initial lookup found no agent, by reason: empty_pool (no agents), unknown_topic (the topic is not in the snapshot, including a topic Metadata returned with an error), no_leader (WarpStream named no leader for the partition), partition_out_of_range (the partition does not exist), or other (no agent was found and no reason was set; not expected with the default strategy). Counted once per record, matching warpstream_produce_records_rejected_total{reason=\"no_agent_assigned\"}.", + }, []string{"reason"}) + hedgeTriggers := promauto.With(reg).NewCounterVec(prometheus.CounterOpts{ Name: "warpstream_produce_hedge_triggers_total", Help: "Why a logical fallback cascade started: latency (hedge timer, or a healthy primary whose computed delay is already zero), primary_failure (the primary failed before the race), or demoted_probe (the routing-time primary was demoted). One increment per cascade entry, including a cascade that dispatches no request. Not a wire request or a hedge wave.", @@ -269,6 +344,17 @@ func newMetrics(reg prometheus.Registerer) *metrics { }).Set(1) return &metrics{ + partitionRoutes: [routeSourceCount]prometheus.Counter{ + routeSourceLeader: partitionRoutes.WithLabelValues(routeSourceLabelLeader), + routeSourceStandIn: partitionRoutes.WithLabelValues(routeSourceLabelStandIn), + }, + routingMisses: [routingMissCount]prometheus.Counter{ + routingMissEmptyPool: routingMisses.WithLabelValues(routingMissLabelEmptyPool), + routingMissUnknownTopic: routingMisses.WithLabelValues(routingMissLabelUnknownTopic), + routingMissNoLeader: routingMisses.WithLabelValues(routingMissLabelNoLeader), + routingMissPartitionOutOfRange: routingMisses.WithLabelValues(routingMissLabelPartitionOutOfRange), + routingMissOther: routingMisses.WithLabelValues(routingMissLabelOther), + }, hedgeAttemptsTotal: promauto.With(reg).NewCounter(prometheus.CounterOpts{ Name: "warpstream_hedge_attempts_total", Help: "Total number of produce requests for which a fanout to per-partition secondaries was attempted. Includes both latency-triggered hedges (primary still in flight) and primary-failure retries.", diff --git a/pkg/wgo/metrics_test.go b/pkg/wgo/metrics_test.go index 6f03e71..327e2b3 100644 --- a/pkg/wgo/metrics_test.go +++ b/pkg/wgo/metrics_test.go @@ -426,3 +426,27 @@ func BenchmarkMetrics_ClusterStatsCollect(b *testing.B) { } } } + +func TestNewMetrics_RoutingCounters(t *testing.T) { + reg := prometheus.NewPedanticRegistry() + m := newMetrics(reg) + + assert.Equal(t, int(routeSourceCount), testutil.CollectAndCount(reg, "warpstream_partition_routes_total")) + assert.Equal(t, int(routingMissCount), testutil.CollectAndCount(reg, "warpstream_routing_misses_total")) + + m.partitionRoutes[routeSourceStandIn].Add(3) + m.routingMisses[routingMissPartitionOutOfRange].Inc() + require.NoError(t, testutil.GatherAndCompare(reg, strings.NewReader(` + # HELP warpstream_partition_routes_total Input records routed to an agent, by how the strategy chose it: leader (the partition's named leader is in the snapshot) or stand_in (the leader entry was missing for a known topic, so a live agent was picked). Counted once per record at the initial routing decision, not per hedge, retry or flush. A demoted leader replaced by the Demoter keeps the classification of the lookup. stand_in covers every missing leader entry, not only an excluded leader. + # TYPE warpstream_partition_routes_total counter + warpstream_partition_routes_total{source="leader"} 0 + warpstream_partition_routes_total{source="stand_in"} 3 + # HELP warpstream_routing_misses_total Input records rejected because the initial lookup found no agent, by reason: empty_pool (no agents), unknown_topic (the topic is not in the snapshot, including a topic Metadata returned with an error), no_leader (WarpStream named no leader for the partition), partition_out_of_range (the partition does not exist), or other (no agent was found and no reason was set; not expected with the default strategy). Counted once per record, matching warpstream_produce_records_rejected_total{reason="no_agent_assigned"}. + # TYPE warpstream_routing_misses_total counter + warpstream_routing_misses_total{reason="empty_pool"} 0 + warpstream_routing_misses_total{reason="no_leader"} 0 + warpstream_routing_misses_total{reason="other"} 0 + warpstream_routing_misses_total{reason="partition_out_of_range"} 1 + warpstream_routing_misses_total{reason="unknown_topic"} 0 + `), "warpstream_partition_routes_total", "warpstream_routing_misses_total")) +} diff --git a/pkg/wgo/partition_assignment.go b/pkg/wgo/partition_assignment.go index 76a5c95..dec5538 100644 --- a/pkg/wgo/partition_assignment.go +++ b/pkg/wgo/partition_assignment.go @@ -47,6 +47,27 @@ func (a Agent) cloneWithState(state AgentState) Agent { return a } +// routeOutcome is how a strategy resolved one initial lookup. The zero value +// means the lookup was not classified, as for a custom strategy: a successful +// route is not counted, and a miss is counted as "other". +type routeOutcome int8 + +const ( + routeUnclassified routeOutcome = iota + routeLeader + routeStandIn + routeMissEmptyPool + routeMissUnknownTopic + routeMissNoLeader + routeMissOutOfRange +) + +// routeClassifier is implemented by strategies that report the outcome from the +// same snapshot that produced the candidates. +type routeClassifier interface { + candidatesWithRoute(topic string, partition int32, maxCandidates int) ([]Agent, routeOutcome) +} + // PartitionAssignmentStrategy maps a partition to an ordered list of // candidate agents. The first candidate is the primary (used for normal // routing); the rest are deterministic alternates used for hedging. @@ -96,6 +117,19 @@ func (l *LazyPartitionAssignmentStrategy) Candidates(topic string, partition int return l.resolve().Candidates(topic, partition, maxCandidates) } +func (l *LazyPartitionAssignmentStrategy) candidatesWithRoute(topic string, partition int32, maxCandidates int) ([]Agent, routeOutcome) { + return candidatesOf(l.resolve(), topic, partition, maxCandidates) +} + +// candidatesOf looks up candidates with the strategy's classification when it +// reports one. +func candidatesOf(s PartitionAssignmentStrategy, topic string, partition int32, maxCandidates int) ([]Agent, routeOutcome) { + if rc, ok := s.(routeClassifier); ok { + return rc.candidatesWithRoute(topic, partition, maxCandidates) + } + return s.Candidates(topic, partition, maxCandidates), routeUnclassified +} + // DefaultPartitionAssignmentStrategy is an immutable snapshot of the agent // pool. The leader map is precomputed in the constructor so the produce hot // path reads it lock-free; Candidates is computed lazily over the same agent @@ -170,34 +204,47 @@ func newDefaultPartitionAssignmentStrategy(agents []int32, leaders map[topicPart // doesn't have that problem. The extra cost from that (more segment // streams per partition) is unmeasured; not assumed to be small. func (s *DefaultPartitionAssignmentStrategy) Candidates(topic string, partition int32, maxCandidates int) []Agent { + agents, _ := s.candidatesWithRoute(topic, partition, maxCandidates) + return agents +} + +func (s *DefaultPartitionAssignmentStrategy) candidatesWithRoute(topic string, partition int32, maxCandidates int) ([]Agent, routeOutcome) { if maxCandidates <= 0 { - return nil + return nil, routeUnclassified } var h uint64 tp := topicPartition{topic: topic, partition: partition} leader, ok := s.leaders[tp] + route := routeLeader if !ok { if _, unnamed := s.noLeader[tp]; unnamed { - return nil + return nil, routeMissNoLeader } if len(s.agents) == 0 { - return nil + return nil, routeMissEmptyPool } if _, topicKnown := s.knownTopics[topic]; !topicKnown { - return nil + // A topic with no named leaders is in partitionCounts but not in + // knownTopics. Use the count only to pick the label. Adding the + // topic to knownTopics would route its holes to a stand-in. + if n, listed := s.partitionCounts[topic]; listed && (partition < 0 || partition >= n) { + return nil, routeMissOutOfRange + } + return nil, routeMissUnknownTopic } if partition < 0 || partition >= s.partitionCounts[topic] { - return nil + return nil, routeMissOutOfRange } h = hashTopicPartition(topic, partition) leader = s.agents[h%uint64(len(s.agents))] + route = routeStandIn } out := make([]Agent, 0, maxCandidates) out = append(out, Agent{NodeID: leader, State: AgentStateHealthy}) if maxCandidates == 1 { - return out + return out, route } // Walk the non-leader agents in deterministic hash order: start at @@ -206,7 +253,7 @@ func (s *DefaultPartitionAssignmentStrategy) Candidates(topic string, partition // acceptable here because n (agent count) is small. nonLeaderCount := len(s.agents) - 1 if nonLeaderCount <= 0 { - return out + return out, route } // The one-candidate return above never reaches this hash. A fallback // pick already stored it. @@ -218,7 +265,7 @@ func (s *DefaultPartitionAssignmentStrategy) Candidates(topic string, partition idx := (start + offset) % nonLeaderCount out = append(out, Agent{NodeID: nthNonLeader(s.agents, leader, idx), State: AgentStateHealthy}) } - return out + return out, route } // nthNonLeader returns the idx-th element of agents skipping leader. idx is diff --git a/pkg/wgo/partition_assignment_test.go b/pkg/wgo/partition_assignment_test.go index 5a86eb7..e86e974 100644 --- a/pkg/wgo/partition_assignment_test.go +++ b/pkg/wgo/partition_assignment_test.go @@ -528,3 +528,91 @@ func TestHashTopicPartition_NoHeapAllocation(t *testing.T) { allocs := testing.AllocsPerRun(1000, func() { _ = hashTopicPartition("some-topic-name", 7) }) assert.Zero(t, allocs) } + +// A topic with no named leaders is not in knownTopics, so its holes are still +// rejected. The partition count only changes the label. +func TestDefaultPartitionAssignmentStrategy_CandidatesWithRoute_AllLeadersUnnamed(t *testing.T) { + const topic = "t" + // Partitions 0 and 2 have no leader. Partition 1 is a hole. + s := newDefaultPartitionAssignmentStrategy([]int32{1, 2, 3}, + map[topicPartition]int32{}, + nil, + map[topicPartition]struct{}{{topic: topic, partition: 0}: {}, {topic: topic, partition: 2}: {}}, + map[string]int32{topic: 3}) + + tests := []struct { + name string + topic string + partition int32 + want routeOutcome + }{ + {"listed partition", topic, 0, routeMissNoLeader}, + {"other listed partition", topic, 2, routeMissNoLeader}, + {"hole is rejected", topic, 1, routeMissUnknownTopic}, + {"past the count", topic, 3, routeMissOutOfRange}, + {"far past the count", topic, 9, routeMissOutOfRange}, + {"negative partition", topic, -1, routeMissOutOfRange}, + {"topic not in Metadata", "other", 0, routeMissUnknownTopic}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, outcome := s.candidatesWithRoute(tt.topic, tt.partition, 1) + assert.Equal(t, tt.want, outcome) + assert.Empty(t, got) + assert.Empty(t, s.Candidates(tt.topic, tt.partition, 1)) + }) + } +} + +func TestDefaultPartitionAssignmentStrategy_CandidatesWithRoute(t *testing.T) { + const topic = "t" + agents := []int32{1, 2, 3} + build := func(agents []int32) *DefaultPartitionAssignmentStrategy { + return newDefaultPartitionAssignmentStrategy(agents, + map[topicPartition]int32{{topic: topic, partition: 0}: 2}, + nil, + map[topicPartition]struct{}{{topic: topic, partition: 1}: {}}, + map[string]int32{topic: 4}) + } + + tests := []struct { + name string + strategy *DefaultPartitionAssignmentStrategy + topic string + partition int32 + max int + want routeOutcome + wantAgent bool + }{ + {"named leader", build(agents), topic, 0, 1, routeLeader, true}, + {"missing leader entry stands in", build(agents), topic, 2, 1, routeStandIn, true}, + {"unnamed leader", build(agents), topic, 1, 1, routeMissNoLeader, false}, + {"empty pool", build(nil), topic, 2, 1, routeMissEmptyPool, false}, + {"unknown topic", build(agents), "other", 0, 1, routeMissUnknownTopic, false}, + {"partition past the count", build(agents), topic, 4, 1, routeMissOutOfRange, false}, + {"negative partition", build(agents), topic, -1, 1, routeMissOutOfRange, false}, + {"no candidates requested", build(agents), topic, 0, 0, routeUnclassified, false}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, outcome := tt.strategy.candidatesWithRoute(tt.topic, tt.partition, tt.max) + assert.Equal(t, tt.want, outcome) + assert.Equal(t, tt.wantAgent, len(got) > 0) + assert.Equal(t, tt.strategy.Candidates(tt.topic, tt.partition, tt.max), got) + }) + } + + t.Run("the lazy strategy forwards the outcome from one snapshot", func(t *testing.T) { + lazy := NewLazyPartitionAssignmentStrategy(func() PartitionAssignmentStrategy { return build(agents) }) + _, outcome := lazy.candidatesWithRoute(topic, 2, 1) + assert.Equal(t, routeStandIn, outcome) + }) + + t.Run("a strategy that cannot classify reports unclassified", func(t *testing.T) { + custom := &mockPartitionAssignmentStrategy{candidates: map[partitionKey][]Agent{{topic, 0}: healthyAgents(7)}} + lazy := NewLazyPartitionAssignmentStrategy(func() PartitionAssignmentStrategy { return custom }) + got, outcome := lazy.candidatesWithRoute(topic, 0, 1) + assert.Equal(t, routeUnclassified, outcome) + assert.Equal(t, healthyAgents(7), got) + }) +} diff --git a/pkg/wgo/route_records_test.go b/pkg/wgo/route_records_test.go index 7231853..58b64d0 100644 --- a/pkg/wgo/route_records_test.go +++ b/pkg/wgo/route_records_test.go @@ -5,6 +5,7 @@ import ( "time" "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/testutil" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/twmb/franz-go/pkg/kgo" @@ -249,6 +250,7 @@ func TestWarpstreamClient_RouteRecords(t *testing.T) { func newRouteRecordsClient(strategy PartitionAssignmentStrategy, nudge chan struct{}) *WarpstreamClient { return &WarpstreamClient{ demoter: NewDemoter(strategy, noopAgentStatsTracker{}, HealthCheckConfig{}, DemoterConfig{}, nopLogger{}, prometheus.NewRegistry()), + metrics: newMetrics(prometheus.NewRegistry()), refreshNowCh: nudge, } } @@ -329,3 +331,73 @@ func benchRouteInputs(topic string, n, partitions, missEvery int) ([]*kgo.Record } return records, &mockPartitionAssignmentStrategy{candidates: candidates} } + +func TestWarpstreamClient_RouteRecordsCountsRoutesAndMisses(t *testing.T) { + const topic = "t" + strategy := newDefaultPartitionAssignmentStrategy([]int32{1, 2}, + map[topicPartition]int32{{topic: topic, partition: 0}: 1}, + nil, + map[topicPartition]struct{}{{topic: topic, partition: 2}: {}}, + map[string]int32{topic: 4}) + newClient := func() *WarpstreamClient { + return newRouteRecordsClient(strategy, make(chan struct{}, 4)) + } + rec := func(topic string, partition int32) *kgo.Record { + return &kgo.Record{Topic: topic, Partition: partition, Value: []byte("v")} + } + routes := func(c *WarpstreamClient) [routeSourceCount]float64 { + var out [routeSourceCount]float64 + for i := range out { + out[i] = testutil.ToFloat64(c.metrics.partitionRoutes[i]) + } + return out + } + misses := func(c *WarpstreamClient) [routingMissCount]float64 { + var out [routingMissCount]float64 + for i := range out { + out[i] = testutil.ToFloat64(c.metrics.routingMisses[i]) + } + return out + } + + t.Run("a batch counts input records, not groups", func(t *testing.T) { + c := newClient() + _, rejected := c.routeRecords([]*kgo.Record{ + rec(topic, 0), rec(topic, 0), rec(topic, 0), + rec(topic, 1), rec(topic, 1), + rec(topic, 2), + rec(topic, 9), rec(topic, 9), + rec("other", 0), + }, countAccepted(new(int))) + + require.Len(t, rejected, 3) + assert.Equal(t, [routeSourceCount]float64{routeSourceLeader: 3, routeSourceStandIn: 2}, routes(c)) + assert.Equal(t, [routingMissCount]float64{ + routingMissNoLeader: 1, + routingMissPartitionOutOfRange: 2, + routingMissUnknownTopic: 1, + }, misses(c)) + }) + + t.Run("a single record counts once", func(t *testing.T) { + c := newClient() + _, err := c.routeRecord(rec(topic, 1), func(ProduceResult) {}) + require.NoError(t, err) + _, err = c.routeRecord(rec(topic, 9), func(ProduceResult) {}) + require.Error(t, err) + + assert.Equal(t, [routeSourceCount]float64{routeSourceStandIn: 1}, routes(c)) + assert.Equal(t, [routingMissCount]float64{routingMissPartitionOutOfRange: 1}, misses(c)) + }) + + t.Run("a custom strategy is not counted as a route and misses are other", func(t *testing.T) { + custom := &mockPartitionAssignmentStrategy{candidates: map[partitionKey][]Agent{{topic, 0}: healthyAgents(5)}} + c := newRouteRecordsClient(custom, make(chan struct{}, 4)) + routed, rejected := c.routeRecords([]*kgo.Record{rec(topic, 0), rec(topic, 1)}, countAccepted(new(int))) + + require.Len(t, routed, 1) + require.Len(t, rejected, 1) + assert.Equal(t, [routeSourceCount]float64{}, routes(c)) + assert.Equal(t, [routingMissCount]float64{routingMissOther: 1}, misses(c)) + }) +}