diff --git a/pkg/kfake/issues_test.go b/pkg/kfake/issues_test.go index f7d4cf823..389b1b6ec 100644 --- a/pkg/kfake/issues_test.go +++ b/pkg/kfake/issues_test.go @@ -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, diff --git a/pkg/kgo/client.go b/pkg/kgo/client.go index ee0fdf0f5..8a3538b3a 100644 --- a/pkg/kgo/client.go +++ b/pkg/kgo/client.go @@ -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 { @@ -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)) @@ -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: @@ -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 { @@ -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 diff --git a/pkg/kgo/client_sweep_test.go b/pkg/kgo/client_sweep_test.go index 762beaea1..d1859354c 100644 --- a/pkg/kgo/client_sweep_test.go +++ b/pkg/kgo/client_sweep_test.go @@ -1,9 +1,11 @@ package kgo import ( + "context" "go/ast" "go/parser" "go/token" + "slices" "testing" "time" @@ -140,6 +142,106 @@ func mkreq(topics ...string) *kmsg.MetadataRequest { return req } +// A brokers-only metadata response replaces cl.brokers and leaves the topic +// cache alone. RequestCachedMetadata must not pair those newer brokers with +// the cached leaders: that response never existed, and the leader is missing +// from Brokers. +func TestRequestCachedMetadataBrokersFromSameResponse(t *testing.T) { + t.Parallel() + cl := &Client{cfg: defaultCfg()} + + rack := "r" + meta := kmsg.NewPtrMetadataResponse() + for _, id := range []int32{1, 2, 3} { + b := kmsg.NewMetadataResponseBroker() + b.NodeID = id + b.Rack = &rack + meta.Brokers = append(meta.Brokers, b) + } + meta.ControllerID = 3 + meta.ClusterID = kmsg.StringPtr("c") + rt := kmsg.NewMetadataResponseTopic() + rt.Topic = kmsg.StringPtr("foo") + rp := kmsg.NewMetadataResponseTopicPartition() + rp.Leader = 3 + rt.Partitions = append(rt.Partitions, rp) + meta.Topics = append(meta.Topics, rt) + + all := kmsg.NewPtrMetadataRequest() // nil Topics: all topics + cl.storeCachedMeta(all, meta, true, nil) + + // Broker 3 has since left. The connection table, the controller the + // client dials, and the cluster id it holds are all from a later + // response. The topic cache still says 3 leads foo. + cl.brokers = []*broker{ + {meta: BrokerMetadata{NodeID: 1}}, + {meta: BrokerMetadata{NodeID: 2}}, + } + cl.controllerID = 1 + live := "live" + cl.clusterID.Store(&live) + + ctx := context.Background() + resp, err := cl.RequestCachedMetadata(ctx, all, time.Hour) + if err != nil { + t.Fatal(err) + } + var ids []int32 + for _, b := range resp.Brokers { + ids = append(ids, b.NodeID) + } + if !slices.Equal(ids, []int32{1, 2, 3}) { + t.Errorf("cached brokers = %v, want [1 2 3]", ids) + } + if resp.ControllerID != 3 { + t.Errorf("cached controller = %d, want 3", resp.ControllerID) + } + if resp.ClusterID == nil || *resp.ClusterID != "c" { + t.Errorf("cached cluster id = %v, want %q", resp.ClusterID, "c") + } + if len(resp.Topics) != 1 || len(resp.Topics[0].Partitions) != 1 || resp.Topics[0].Partitions[0].Leader != 3 { + t.Fatalf("cached topics = %+v, want foo led by 3", resp.Topics) + } + // dups must not hand back the snapshot's Rack, nor the one on the + // response the snapshot was built from. A caller writing through it + // would change the cache. + if resp.Brokers[0].Rack == &rack || resp.Brokers[0].Rack == meta.Brokers[0].Rack { + t.Error("returned Rack aliases internal state") + } + + // A topic-bearing response that reports no controller returns -1, + // not the id the client still dials. A request with no topics keeps + // that dial id, since it has no snapshot to take one from. + unknown := kmsg.NewPtrMetadataResponse() + unknown.ControllerID = -1 + ub := kmsg.NewMetadataResponseBroker() + ub.NodeID = 1 + unknown.Brokers = append(unknown.Brokers, ub) + ut := kmsg.NewMetadataResponseTopic() + ut.Topic = kmsg.StringPtr("foo") + up := kmsg.NewMetadataResponseTopicPartition() + up.Leader = 1 + ut.Partitions = append(ut.Partitions, up) + unknown.Topics = append(unknown.Topics, ut) + cl.controllerID = 5 + cl.storeCachedMeta(all, unknown, true, nil) + + resp, err = cl.RequestCachedMetadata(ctx, all, time.Hour) + if err != nil { + t.Fatal(err) + } + if resp.ControllerID != -1 { + t.Errorf("cached controller = %d, want -1 from the response", resp.ControllerID) + } + resp, err = cl.RequestCachedMetadata(ctx, mkreq(), time.Hour) + if err != nil { + t.Fatal(err) + } + if resp.ControllerID != 5 { + t.Errorf("no-topics controller = %d, want the live 5", resp.ControllerID) + } +} + // We evict cached topics only from a request that says something about what // we want cached. A broker only request asks for no topics, and internal // targeted fetches (topic IDs, unknown produce topics, the group's topics)