Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. 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.
Comment thread
koloss2001 marked this conversation as resolved.
Outdated

### Buffering: linger by destination agent, not by partition

Expand Down
7 changes: 4 additions & 3 deletions docs/internal/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`

Expand Down
72 changes: 60 additions & 12 deletions pkg/wgo/agentpool.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,22 +61,30 @@ 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
}

// 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))
Expand All @@ -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.
Expand All @@ -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
Expand All @@ -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 {
Comment thread
stephclay marked this conversation as resolved.
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.
Expand Down
107 changes: 101 additions & 6 deletions pkg/wgo/agentpool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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{{
Expand Down Expand Up @@ -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{},
Expand All @@ -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)
})
}
}
Expand Down
28 changes: 26 additions & 2 deletions pkg/wgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Expand All @@ -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).
Comment thread
koloss2001 marked this conversation as resolved.
c.noteLeaderDrops(dropped, true)
Comment thread
stephclay marked this conversation as resolved.
}

// 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
Expand Down
Loading
Loading