Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
5. **No custom partitioner.** wgo has no default partitioning logic and does not accept a custom `kgo.Partitioner`. The caller must set `record.Partition` on every record before calling Produce. An unset `Partition` field silently routes to partition 0.

## How it works
Expand All @@ -49,7 +49,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

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 {
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 @@ -439,7 +443,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 @@ -449,6 +453,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
Expand Down
Loading
Loading