@@ -2,6 +2,8 @@ package kfake
22
33import (
44 "hash/crc32"
5+ "net"
6+ "strconv"
57 "time"
68
79 "github.com/twmb/franz-go/pkg/kerr"
@@ -14,7 +16,7 @@ import (
1416// * Multiple batches in one produce
1517// * Compact
1618
17- func init () { regKey (0 , 3 , 9 ) }
19+ func init () { regKey (0 , 3 , 10 ) }
1820
1921func (c * Cluster ) handleProduce (b * broker , kreq kmsg.Request ) (kmsg.Response , error ) {
2022 var (
@@ -46,13 +48,25 @@ func (c *Cluster) handleProduce(b *broker, kreq kmsg.Request) (kmsg.Response, er
4648 donet (t , errCode )
4749 }
4850 }
51+ var includeBrokers bool
4952 toresp := func () kmsg.Response {
5053 for topic , partitions := range tdone {
5154 st := kmsg .NewProduceResponseTopic ()
5255 st .Topic = topic
5356 st .Partitions = partitions
5457 resp .Topics = append (resp .Topics , st )
5558 }
59+ if includeBrokers {
60+ for _ , b := range c .bs {
61+ sb := kmsg .NewProduceResponseBroker ()
62+ h , p , _ := net .SplitHostPort (b .ln .Addr ().String ())
63+ p32 , _ := strconv .Atoi (p )
64+ sb .NodeID = b .node
65+ sb .Host = h
66+ sb .Port = int32 (p32 )
67+ resp .Brokers = append (resp .Brokers , sb )
68+ }
69+ }
5670 return resp
5771 }
5872
@@ -76,7 +90,10 @@ func (c *Cluster) handleProduce(b *broker, kreq kmsg.Request) (kmsg.Response, er
7690 continue
7791 }
7892 if pd .leader != b {
79- donep (rt .Topic , rp , kerr .NotLeaderForPartition .Code )
93+ p := donep (rt .Topic , rp , kerr .NotLeaderForPartition .Code )
94+ p .CurrentLeader .LeaderID = pd .leader .node
95+ p .CurrentLeader .LeaderEpoch = pd .epoch
96+ includeBrokers = true
8097 continue
8198 }
8299
0 commit comments