Skip to content
Open
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
172 changes: 172 additions & 0 deletions pkg/kfake/issues_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1888,6 +1888,178 @@ func TestRequestCachedMetadata(t *testing.T) {
})
}

// TestRequestCachedMetadataBrokersMatchTopics reproduces a
// MetadataResponse whose Brokers and Topics come from different
// responses.
//
// The cluster answers brokers {0, 1, 2} with every partition led by 2,
// then {0, 1} led by 0. Ping applies the second answer as brokers only,
// so the topic cache still says leader 2. A cache hit must not return
// that leader next to the new broker list. A later request that spans a
// topic from each answer must not either.
func TestRequestCachedMetadataBrokersMatchTopics(t *testing.T) {
t.Parallel()
const topic1, topic2 = "topic1", "topic2"
c := newCluster(t, NumBrokers(3))

// kfake node i listens on ListenAddrs()[i]. The crafted brokers have
// to keep that mapping, or KIP-1242 rejects the connection.
addrs := c.ListenAddrs()

// gone: broker 2 has left and partitions moved to 0.
var gone atomic.Bool
var topicRequests atomic.Int32
c.ControlKey(int16(kmsg.Metadata), func(kreq kmsg.Request) (kmsg.Response, error, bool) {
c.KeepControl()
req := kreq.(*kmsg.MetadataRequest)
resp := req.ResponseKind().(*kmsg.MetadataResponse)

nodes := []int32{0, 1, 2}
leader := int32(2)
if gone.Load() {
nodes = []int32{0, 1}
leader = 0
}
for _, n := range nodes {
host, portStr, _ := net.SplitHostPort(addrs[n])
port, _ := strconv.Atoi(portStr)
b := kmsg.NewMetadataResponseBroker()
b.NodeID = n
b.Host = host
b.Port = int32(port)
resp.Brokers = append(resp.Brokers, b)
}
resp.ControllerID = 0
resp.ClusterID = kmsg.StringPtr("kfake")

brokersOnly := req.Topics != nil && len(req.Topics) == 0
if brokersOnly {
return resp, nil, true
}
topicRequests.Add(1)

var names []string
if req.Topics == nil {
names = []string{topic1, topic2}
}
for _, rt := range req.Topics {
if rt.Topic != nil {
names = append(names, *rt.Topic)
}
}
for _, name := range names {
st := kmsg.NewMetadataResponseTopic()
st.Topic = kmsg.StringPtr(name)
st.TopicID = [16]byte{name[0]}
sp := kmsg.NewMetadataResponseTopicPartition()
sp.Partition = 0
sp.Leader = leader
sp.Replicas = []int32{leader}
sp.ISR = []int32{leader}
st.Partitions = append(st.Partitions, sp)
resp.Topics = append(resp.Topics, st)
}
return resp, nil, true
})

cl := newPlainClient(t, c)
ctx := context.Background()

nodeIDs := func(resp *kmsg.MetadataResponse) []int32 {
ids := make([]int32, 0, len(resp.Brokers))
for _, b := range resp.Brokers {
ids = append(ids, b.NodeID)
}
slices.Sort(ids)
return ids
}
leadersInBrokers := func(t *testing.T, resp *kmsg.MetadataResponse) {
t.Helper()
ids := nodeIDs(resp)
for _, rt := range resp.Topics {
for _, p := range rt.Partitions {
if !slices.Contains(ids, p.Leader) {
t.Errorf("topic %s partition %d: leader %d not in brokers %v", *rt.Topic, p.Partition, p.Leader, ids)
}
}
}
}
metaReq := func(topics ...string) *kmsg.MetadataRequest {
req := kmsg.NewPtrMetadataRequest()
req.Topics = []kmsg.MetadataRequestTopic{}
for _, topic := range topics {
rt := kmsg.NewMetadataRequestTopic()
rt.Topic = kmsg.StringPtr(topic)
req.Topics = append(req.Topics, rt)
}
return req
}
all := kmsg.NewPtrMetadataRequest()

// Fill the cache from the {0,1,2} view.
resp, err := cl.RequestCachedMetadata(ctx, all, time.Hour)
if err != nil {
t.Fatal(err)
}
if got := nodeIDs(resp); !slices.Equal(got, []int32{0, 1, 2}) {
t.Fatalf("initial brokers: got %v, expected [0 1 2]", got)
}
leadersInBrokers(t, resp)

// Broker 2 leaves. Ping is a brokers-only Metadata: it rewrites the
// connection table to {0,1} and leaves the topic cache alone.
gone.Store(true)
if err := cl.Ping(ctx); err != nil {
t.Fatal(err)
}
if got := len(cl.DiscoveredBrokers()); got != 2 {
t.Fatalf("discovered brokers after ping: got %d, expected 2", got)
}

// A cache hit must return the broker list the cached leaders came
// with, not the connection table.
before := topicRequests.Load()
resp, err = cl.RequestCachedMetadata(ctx, all, time.Hour)
if err != nil {
t.Fatal(err)
}
if topicRequests.Load() != before {
t.Fatal("expected a cache hit")
}
if got := nodeIDs(resp); !slices.Equal(got, []int32{0, 1, 2}) {
t.Errorf("cache hit brokers: got %v, expected [0 1 2]", got)
}
leadersInBrokers(t, resp)

// topic1 is still from {0, 1, 2}; refresh only topic2, which is now
// led by 0. One request for both must not keep leader 2 beside the
// broker list that no longer has 2, and it must take one fetch to
// get there. The request after that shares one response, so it hits.
if _, err := cl.RequestCachedMetadata(ctx, metaReq(topic2), time.Nanosecond); err != nil {
t.Fatal(err)
}
before = topicRequests.Load()
resp, err = cl.RequestCachedMetadata(ctx, metaReq(topic1, topic2), time.Hour)
if err != nil {
t.Fatal(err)
}
if got := topicRequests.Load() - before; got != 1 {
t.Errorf("mixed cache hit issued %d topic metadata requests, expected 1", got)
}
if len(resp.Topics) != 2 {
t.Fatalf("got %d topics, expected 2", len(resp.Topics))
}
leadersInBrokers(t, resp)

before = topicRequests.Load()
if _, err := cl.RequestCachedMetadata(ctx, metaReq(topic1, topic2), time.Hour); err != nil {
t.Fatal(err)
}
if topicRequests.Load() != before {
t.Error("expected a cache hit once both topics share one response")
}
}

func TestKadmCachedMetadata(t *testing.T) {
t.Parallel()
c := newCluster(t,
Expand Down
154 changes: 132 additions & 22 deletions pkg/kgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -1644,11 +1644,36 @@ func (cl *Client) RequestCachedMetadata(ctx context.Context, req *kmsg.MetadataR
}

// Phase 3: fetch all resolved topic names, using the cache.
// resolveTopicMeta reuses this slice for the names it still has to
// fetch, so keep the caller's set aside for the refetch below.
var refetch []string
if len(topics) > 0 {
refetch = slices.Clone(topics)
}
cached, err := cl.resolveTopicMeta(ctx, topics, true, limit)
if err != nil {
return nil, err
}

// A cache hit can join topics from two responses: a partial hit
// refetched only the missing names, or a targeted fetch (the metadata
// loop asking for the topics it tracks) overwrote some entries of a
// still fresh all-topics fetch. Neither response's broker list is
// right for the other half, so fetch the caller's set once more. That
// fetch's results are written by one storeCachedMeta call and cannot
// mix again.
if cachedMetaMixed(cached) {
cached, err = cl.resolveTopicMeta(ctx, refetch, false, limit)
if err != nil {
return nil, err
}
}
var brokers *cachedMetaBrokers
for _, t := range cached {
brokers = t.brokers
break
}

// Phase 4: build the response. We deeply clone all cached data so
// that the end user cannot modify internal data.
dups := func(s *string) *string {
Expand Down Expand Up @@ -1685,21 +1710,39 @@ func (cl *Client) RequestCachedMetadata(ctx context.Context, req *kmsg.MetadataR

resp := kmsg.NewPtrMetadataResponse()

cl.brokersMu.RLock()
for _, b := range cl.brokers {
resp.Brokers = append(resp.Brokers, kmsg.MetadataResponseBroker{
NodeID: b.meta.NodeID,
Host: b.meta.Host,
Port: b.meta.Port,
Rack: dups(b.meta.Rack),
})
}
cl.brokersMu.RUnlock()
if brokers != nil {
// The topics came from a metadata response, so the brokers,
// controller, and cluster ID must come from that same response.
// On a cache hit they are as old as the topics, which is the age
// the caller already accepted via limit.
resp.Brokers = make([]kmsg.MetadataResponseBroker, 0, len(brokers.brokers))
for _, b := range brokers.brokers {
b.Rack = dups(b.Rack)
resp.Brokers = append(resp.Brokers, b)
}
resp.ClusterID = dups(brokers.clusterID)
resp.ControllerID = brokers.controllerID
} else {
// No topics: a brokers-only request, or nothing cached. There is
// no topic half to agree with, so the live connection table is
// the answer. This is what kadm.BrokerMetadata relies on, and it
// does not fetch when we already know a broker.
cl.brokersMu.RLock()
for _, b := range cl.brokers {
resp.Brokers = append(resp.Brokers, kmsg.MetadataResponseBroker{
NodeID: b.meta.NodeID,
Host: b.meta.Host,
Port: b.meta.Port,
Rack: dups(b.meta.Rack),
})
}
cl.brokersMu.RUnlock()

resp.ClusterID = dups(cl.clusterID.Load())
cl.controllerIDMu.Lock()
resp.ControllerID = cl.controllerID
cl.controllerIDMu.Unlock()
resp.ClusterID = dups(cl.clusterID.Load())
cl.controllerIDMu.Lock()
resp.ControllerID = cl.controllerID
cl.controllerIDMu.Unlock()
}

for _, t := range cached {
resp.Topics = append(resp.Topics, dupt(t.t))
Expand Down Expand Up @@ -3119,11 +3162,49 @@ func firstErrMerger(sresps []ResponseShard, merge func(kresp kmsg.Response)) err
return firstErr
}

// cachedMetaBrokers is the broker half of one metadata response, saved with
// every topic that response contained. RequestCachedMetadata returns those
// topics with this list, not with cl.brokers. cl.brokers is the connection
// table, and every metadata response rewrites it: the metadata loop, Ping,
// and fetchBrokerMetadata included. Topics stored by one storeCachedMeta
// call share one pointer, and that pointer is the generation.
//
// A caller can otherwise observe a pair no broker sent. RequestCachedMetadata
// fetches a response listing brokers {1, 2, 3} and a partition led by 3, and
// storeCachedMeta caches the topic. Broker 3 then leaves. Ping gets {1, 2}
// back and updateBrokers replaces cl.brokers. A brokers-only response does
// not touch the topic cache, so the topic still says leader 3. The next
// RequestCachedMetadata within limit copies cl.brokers and the cached topic,
// and the caller reads that struct as one response.
type cachedMetaBrokers struct {
brokers []kmsg.MetadataResponseBroker
controllerID int32
clusterID *string
}

type cachedMetaTopic struct {
id [16]byte
t kmsg.MetadataResponseTopic
ps map[int32]kmsg.MetadataResponseTopicPartition
when time.Time
id [16]byte
t kmsg.MetadataResponseTopic
ps map[int32]kmsg.MetadataResponseTopicPartition
when time.Time
brokers *cachedMetaBrokers // nil when a test fills the entry without storeCachedMeta
}

// cachedMetaMixed reports whether these topics were stored from more than
// one metadata response. Pointer identity is the generation.
func cachedMetaMixed(topics map[string]cachedMetaTopic) bool {
var first *cachedMetaBrokers
var seen bool
for _, t := range topics {
if !seen {
first, seen = t.brokers, true
continue
}
if t.brokers != first {
return true
}
}
return false
}

// For NOT_LEADER_FOR_PARTITION:
Expand Down Expand Up @@ -3272,6 +3353,34 @@ func (cl *Client) storeCachedMeta(req *kmsg.MetadataRequest, meta *kmsg.Metadata
cl.metaCache.byID = make(map[[16]byte]string)
}
when := time.Now()

// One broker snapshot for every topic in this response. Rack and
// ClusterID are cloned for the same reason the topics below are:
// cl.Request hands this response back to the user, who may write
// through those pointers. ControllerID is stored as the response
// sent it, including -1. updateMetadataBrokers ignores -1 so the
// client can still dial the last controller it knew, but that id
// may not be in this response's broker list.
var brokers *cachedMetaBrokers
if len(meta.Topics) > 0 {
brokers = &cachedMetaBrokers{
brokers: slices.Clone(meta.Brokers),
controllerID: meta.ControllerID,
}
for i := range brokers.brokers {
b := &brokers.brokers[i]
if b.Rack != nil {
rack := *b.Rack
b.Rack = &rack
}
b.UnknownTags = kmsg.Tags{}
}
if meta.ClusterID != nil {
clusterID := *meta.ClusterID
brokers.clusterID = &clusterID
}
}

var stored int
for _, topic := range meta.Topics {
if topic.Topic == nil {
Expand Down Expand Up @@ -3302,10 +3411,11 @@ func (cl *Client) storeCachedMeta(req *kmsg.MetadataRequest, meta *kmsg.Metadata
p.OfflineReplicas = slices.Clone(p.OfflineReplicas)
}
t := cachedMetaTopic{
id: topic.TopicID,
t: topic,
ps: make(map[int32]kmsg.MetadataResponseTopicPartition),
when: when,
id: topic.TopicID,
t: topic,
ps: make(map[int32]kmsg.MetadataResponseTopicPartition),
when: when,
brokers: brokers,
}
// A recreated topic comes back under a new ID. Delete the old
// ID's mapping when overwriting the entry, else byID accumulates
Expand Down
Loading