Skip to content
Merged
Show file tree
Hide file tree
Changes from 12 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: 3 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@ read it when planning changes to the areas it covers, and keep it current:

- [`docs/internal/metrics.md`](docs/internal/metrics.md) — how metrics are split
(transport vs producer-state vs warpstream-specific) and the franz-go/`kprom`
drop-in-compatibility contract. Update it when changing metrics.
drop-in-compatibility contract. Update it when that categorization or parity
contract changes. Do not catalog individual custom metrics here; their
Prometheus help strings are the source of truth.
- [`docs/internal/tracing.md`](docs/internal/tracing.md) — how the client drives
franz-go's produce-record hooks on its own produce path to support tracing (e.g.
`kotel`) as a drop-in. Update it when changing hook or tracing behaviour.
Expand Down
25 changes: 25 additions & 0 deletions pkg/wgo/agentpool.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package wgo
import (
"context"
"fmt"
"slices"
"sort"
"sync/atomic"
"time"
Expand Down Expand Up @@ -92,6 +93,8 @@ func (p *AgentPool) refresh(ctx context.Context) (removed []int32, dropped leade
newAgents = append(newAgents, b.NodeID)
}
sort.Slice(newAgents, func(i, j int) bool { return newAgents[i] < newAgents[j] })
// A NodeID listed twice is still one agent.
newAgents = slices.Compact(newAgents)
Comment thread
koloss2001 marked this conversation as resolved.
Outdated

agentSet := make(map[int32]struct{}, len(newAgents))
for _, id := range newAgents {
Expand Down Expand Up @@ -224,3 +227,25 @@ func diffRemovedAgents(old []int32, newSet map[int32]struct{}) []int32 {
}
return removed
}

// diffAgentMembership counts NodeIDs that appeared or disappeared between two
// sorted, unique agent lists.
func diffAgentMembership(old, new []int32) (added, removed int) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

diffAgentMembership assumes sorted, unique input. refresh sorts newAgents but doesn't dedupe it. If a Metadata response ever lists a NodeID twice, [5,5] followed by [5] counts one removed with no real membership change, which inflates agents_changed_total. diffRemovedAgents is unaffected because it goes through agentSet. Probably rare in practice, but a slices.Compact after the sort would make the documented precondition true. (Or derive the counts from the same set-based diff refresh already does.)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. refresh now deduplicates NodeIDs after sorting. It was worse than described: the pool itself held the duplicate and membership_changed was also counted. Added a test that duplicates a broker in the Metadata response.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm reverting this. A correct Metadata response contains one entry per NodeID, and a duplicate is treated as malformed, so the sorted list already satisfies diffAgentMembership. Compacting also changes fallback routing for a malformed response (hash % len(agents) and the secondary walk), which this PR otherwise doesn't touch. I've removed the slices.Compact call and its test.

i, j := 0, 0
for i < len(old) && j < len(new) {
switch {
case old[i] == new[j]:
i++
j++
case old[i] < new[j]:
removed++
i++
default:
added++
j++
}
}
removed += len(old) - i
added += len(new) - j
return added, removed
}
35 changes: 35 additions & 0 deletions pkg/wgo/agentpool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -366,6 +366,41 @@ func TestBuildLeadersAndTopicIDs(t *testing.T) {

func stringPtr(s string) *string { return &s }

func TestAgentPool_DiffAgentMembership(t *testing.T) {
tests := map[string]struct {
old []int32
new []int32
wantAdded int
wantRemoved int
}{
"unchanged": {
old: []int32{1, 2, 3}, new: []int32{1, 2, 3},
},
"one added": {
old: []int32{1, 2}, new: []int32{1, 2, 3}, wantAdded: 1,
},
"one removed": {
old: []int32{1, 2, 3}, new: []int32{1, 3}, wantRemoved: 1,
},
"replace one": {
old: []int32{1, 2, 3}, new: []int32{1, 4}, wantAdded: 1, wantRemoved: 2,
},
"empty to some": {
old: nil, new: []int32{1, 2}, wantAdded: 2,
},
"all removed": {
old: []int32{1, 2}, new: nil, wantRemoved: 2,
},
}
for name, tc := range tests {
t.Run(name, func(t *testing.T) {
added, removed := diffAgentMembership(tc.old, tc.new)
assert.Equal(t, tc.wantAdded, added)
assert.Equal(t, tc.wantRemoved, removed)
})
}
}

func TestDiffRemovedAgents(t *testing.T) {
tests := map[string]struct {
old []int32
Expand Down
8 changes: 5 additions & 3 deletions pkg/wgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,7 @@ 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.
// Publish the excluded-leader count 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
Expand Down Expand Up @@ -503,12 +503,14 @@ func (c *WarpstreamClient) refreshPool(trigger metadataRefreshTrigger) {
c.noteLeaderDrops(dropped, true)
}

// noteLeaderDrops counts and logs excluded leaders. nudge asks for another fetch.
// noteLeaderDrops publishes the excluded-leader count and logs it when it is
// above zero. nudge asks for another fetch. The gauge is set even at zero so it
// clears when the exclusion does.
func (c *WarpstreamClient) noteLeaderDrops(dropped leaderDrops, nudge bool) {
c.metrics.agentPoolExcludedLeaders.Set(float64(dropped.Count))
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,
Expand Down
Loading
Loading