Skip to content
Merged
Show file tree
Hide file tree
Changes from 11 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
22 changes: 22 additions & 0 deletions pkg/wgo/agentpool.go
Original file line number Diff line number Diff line change
Expand Up @@ -224,3 +224,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
201 changes: 194 additions & 7 deletions pkg/wgo/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,30 @@ func newTestWarpstreamClient(t *testing.T, topic string, numPartitions int32, op
func TestWarpstreamClient_ProduceSync(t *testing.T) {
const topic = "test-topic"

t.Run("attempted payload equals the wire counters when every produce succeeds", func(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
c, _, _, _ := newTestWarpstreamClient(t, topic, 1)

// Three records for one partition go out as one merged batch.
results := c.ProduceSync(t.Context(), []*kgo.Record{
{Topic: topic, Partition: 0, Value: []byte("a"), Timestamp: time.Now()},
{Topic: topic, Partition: 0, Value: []byte("b"), Timestamp: time.Now()},
{Topic: topic, Partition: 0, Value: []byte("c"), Timestamp: time.Now()},
})
require.Len(t, results, 3)
for _, res := range results {
require.NoError(t, res.Err)
}

assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.produceAttemptRecords[attemptPrimary]))
assert.Equal(t, testutil.ToFloat64(c.metrics.produceWireRecordsTotal), testutil.ToFloat64(c.metrics.produceAttemptRecords[attemptPrimary]))
assert.Equal(t, testutil.ToFloat64(c.metrics.produceWireCompressedBytesTotal), testutil.ToFloat64(c.metrics.produceAttemptBytes[attemptPrimary]))
assert.Positive(t, testutil.ToFloat64(c.metrics.produceAttemptBytes[attemptPrimary]))
assert.Zero(t, testutil.ToFloat64(c.metrics.produceAttemptRecords[attemptHedge]))
assert.Zero(t, testutil.ToFloat64(c.metrics.produceAttemptBytes[attemptHedge]))
})
})

t.Run("single record produces and is consumable", func(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
c, _, clusterAddr, vnet := newTestWarpstreamClient(t, topic, 1)
Expand Down Expand Up @@ -275,6 +299,11 @@ func TestWarpstreamClient_ProduceSync(t *testing.T) {
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceRecordsFailedTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned)))
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedRecordTooLarge)))

// The caller cancel detaches this wait. The flush keeps its own
// deadline, and a successful background produce is not a Hedger failure.
synctest.Wait()
assertFinalOutcomes(t, c.metrics, nil)
})
})

Expand Down Expand Up @@ -1229,6 +1258,13 @@ func TestWarpstreamClient_WriteTimeoutUnblocksStuckProduce(t *testing.T) {
// Returned at the write-timeout ceiling, not hanging.
assert.GreaterOrEqual(t, elapsed, writeTimeout)
assert.Less(t, elapsed, writeTimeout+time.Second)
assertFinalOutcomes(t, c.metrics, map[produceFinalOutcome]float64{
produceFinalOutcomeWriteTimeout: 1,
})
// The stuck attempt was dispatched, so it counts as attempted,
// but nothing was acked on the wire.
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceAttemptRecords[attemptPrimary]))
assert.Zero(t, testutil.ToFloat64(c.metrics.produceWireRecordsTotal))
})
})

Expand All @@ -1246,6 +1282,9 @@ func TestWarpstreamClient_WriteTimeoutUnblocksStuckProduce(t *testing.T) {
case <-time.After(5 * time.Second):
t.Fatal("Produce promise did not fire within the write-timeout bound")
}
assertFinalOutcomes(t, c.metrics, map[produceFinalOutcome]float64{
produceFinalOutcomeWriteTimeout: 1,
})
})
})
})
Expand Down Expand Up @@ -2012,6 +2051,76 @@ func TestWarpstreamClient_OnDemandMetadataRefresh(t *testing.T) {
})
}

// TestWarpstreamClient_AgentPoolChurnOnLiveRefresh drives real membership
// changes through the periodic refresh and checks the churn counters: the
// constructor refresh and an unchanged refresh add nothing, an added Agent
// counts once as added, and a removed Agent counts once as removed.
func TestWarpstreamClient_AgentPoolChurnOnLiveRefresh(t *testing.T) {
const topic = "test-topic"

synctest.Test(t, func(t *testing.T) {
reg := prometheus.NewPedanticRegistry()
vnet := &kfake.VirtualNetwork{}
cluster, clusterAddr := testkafka.CreateCluster(t, 1, topic, testkafka.WithVirtualNetwork(vnet))
c, err := NewWarpstreamClient(nil, reg, append(testWarpstreamOpts(clusterAddr, topic), WithDialer(vnet.DialContext))...)
require.NoError(t, err)
t.Cleanup(c.Close)

churn := func() (added, removed float64) {
return testutil.ToFloat64(c.metrics.agentpoolAgentsChanged[agentpoolChurnAdded]),
testutil.ToFloat64(c.metrics.agentpoolAgentsChanged[agentpoolChurnRemoved])
}
refreshOnce := func() {
time.Sleep(10 * time.Second)
synctest.Wait()
}

// The constructor refresh is not churn.
added, removed := churn()
assert.Zero(t, added)
assert.Zero(t, removed)

refreshOnce()
added, removed = churn()
assert.Zero(t, added, "an unchanged refresh adds nothing")
assert.Zero(t, removed)

const newNode = int32(10)
_, _, err = cluster.AddNode(newNode, 9092+5)
require.NoError(t, err)
refreshOnce()
added, removed = churn()
assert.Equal(t, float64(1), added)
assert.Zero(t, removed)

// The first refresh after RemoveNode may fail while the client still
// holds a connection to the departed node. A failed refresh keeps the
// previous snapshot, so it must not count as churn. Refresh until the
// pool drops the node, then check the removal is counted exactly once.
require.NoError(t, cluster.RemoveNode(newNode))
for range 5 {
if !slices.Contains(c.pool.Agents(), newNode) {
break
}
refreshOnce()
if slices.Contains(c.pool.Agents(), newNode) {
_, removed = churn()
assert.Zero(t, removed, "no churn while the pool still holds the node")
}
}
require.NotContains(t, c.pool.Agents(), newNode)
added, removed = churn()
assert.Equal(t, float64(1), added)
assert.Equal(t, float64(1), removed)

// Further unchanged refreshes add nothing.
refreshOnce()
added, removed = churn()
assert.Equal(t, float64(1), added)
assert.Equal(t, float64(1), removed)
})
}

func TestWarpstreamClient_IdleClusterStats(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
reg := prometheus.NewPedanticRegistry()
Expand Down Expand Up @@ -2048,29 +2157,39 @@ func TestWarpstreamClient_NoteLeaderDrops(t *testing.T) {
}

c.noteLeaderDrops(leaderDrops{}, true)
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))
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))
assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))
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.
// Constructor path publishes the count and does not fill the nudge channel.
// The gauge holds the latest count; it does not accumulate across refreshes.
c.noteLeaderDrops(leaderDrops{Count: 1, Topic: "ingest", Partition: 0, NodeID: 7}, false)
assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))
select {
case <-c.refreshNowCh:
t.Fatal("constructor refresh nudged")
default:
}

// A clean refresh clears the gauge and does not nudge.
c.noteLeaderDrops(leaderDrops{}, true)
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))
select {
case <-c.refreshNowCh:
t.Fatal("clean refresh nudged")
default:
}
}

func TestWarpstreamClient_WaitRefreshCooldown(t *testing.T) {
Expand Down Expand Up @@ -2175,7 +2294,7 @@ func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) {
c.triggerRefresh()
synctest.Wait()
assertOnDemandRefresh(t, c, 1, 0)
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))
assertMissingLeader(t, c, topic)

results := c.ProduceSync(t.Context(), []*kgo.Record{
Expand All @@ -2192,7 +2311,9 @@ func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) {
synctest.Wait()
assertOnDemandRefresh(t, c, float64(i+2), 0)
}
assert.Equal(t, float64(6), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal))
// Six refreshes later the gauge still reads one excluded leader: it is
// state, not a running total, so the backoff does not move it.
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))

time.Sleep(time.Second)
synctest.Wait()
Expand All @@ -2202,7 +2323,8 @@ func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) {
time.Sleep(9 * time.Second)
synctest.Wait()
assertOnDemandRefresh(t, c, 7, 0)
assert.Equal(t, float64(6), testutil.ToFloat64(c.metrics.agentPoolLeaderDroppedTotal))
// The clean fetch clears the gauge.
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders))
assertLeaderRestored(t, c, topic, leader)

restored := c.ProduceSync(t.Context(), []*kgo.Record{
Expand Down Expand Up @@ -2388,6 +2510,71 @@ func TestWarpstreamClient_OnDemandRefreshBackoff(t *testing.T) {
})
}

// TestWarpstreamClient_ExcludedLeadersGaugeAcrossRefreshes checks the gauge
// follows the latest successful snapshot: set while a leader is excluded,
// unchanged by a failed refresh, and back to zero once a clean refresh lands.
func TestWarpstreamClient_ExcludedLeadersGaugeAcrossRefreshes(t *testing.T) {
const topic = "test-topic"

synctest.Test(t, func(t *testing.T) {
vnet := &kfake.VirtualNetwork{}
cluster, clusterAddr := testkafka.CreateCluster(t, 1, topic,
testkafka.WithVirtualNetwork(vnet), testkafka.WithNumBrokers(3))
c, err := NewWarpstreamClient(nil, prometheus.NewPedanticRegistry(), append(
testWarpstreamOpts(clusterAddr, topic), WithDialer(vnet.DialContext))...)
require.NoError(t, err)
t.Cleanup(c.Close)

gauge := func() float64 { return testutil.ToFloat64(c.metrics.agentPoolExcludedLeaders) }
// The constructor refresh saw a healthy cluster.
assert.Zero(t, gauge())

raw, err := c.Request(t.Context(), kmsg.NewPtrMetadataRequest())
require.NoError(t, err)
template := raw.(*kmsg.MetadataResponse)

var poison, fail atomic.Bool
cluster.ControlKey(int16(kmsg.Metadata), func(req kmsg.Request) (kmsg.Response, error, bool) {
cluster.KeepControl()
switch {
case fail.Load():
// A response error fails the fetch in one round trip.
resp := cloneMetadataResponse(template, req.GetVersion())
resp.ErrorCode = kerr.UnknownServerError.Code
return resp, nil, true
case poison.Load():
return metadataWithMissingLeaders(template, 99, req.GetVersion()), nil, true
}
return nil, nil, false
})

// A refresh that excludes a leader publishes it.
poison.Store(true)
time.Sleep(time.Nanosecond)
c.triggerRefresh()
synctest.Wait()
assert.Equal(t, float64(1), gauge())
assertMissingLeader(t, c, topic)

// A failed refresh keeps the previous snapshot and the previous value.
fail.Store(true)
c.triggerRefresh()
time.Sleep(3 * time.Second)
synctest.Wait()
assert.GreaterOrEqual(t, failedOnDemandRefreshes(c), float64(1))
assert.Equal(t, float64(1), gauge(), "a failed refresh must not move the gauge")
assertMissingLeader(t, c, topic)

// A clean refresh clears it.
fail.Store(false)
poison.Store(false)
c.triggerRefresh()
time.Sleep(10 * time.Second)
synctest.Wait()
assert.Zero(t, gauge())
})
}

func assertOnDemandRefresh(t *testing.T, c *WarpstreamClient, onDemand, periodic float64) {
t.Helper()
assert.Equal(t, onDemand, testutil.ToFloat64(c.metrics.metadataRefreshResultsTotal.WithLabelValues(
Expand Down
Loading
Loading