metadata_response.go 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266
  1. package sarama
  2. type PartitionMetadata struct {
  3. Err KError
  4. ID int32
  5. Leader int32
  6. Replicas []int32
  7. Isr []int32
  8. }
  9. func (pm *PartitionMetadata) decode(pd packetDecoder) (err error) {
  10. tmp, err := pd.getInt16()
  11. if err != nil {
  12. return err
  13. }
  14. pm.Err = KError(tmp)
  15. pm.ID, err = pd.getInt32()
  16. if err != nil {
  17. return err
  18. }
  19. pm.Leader, err = pd.getInt32()
  20. if err != nil {
  21. return err
  22. }
  23. pm.Replicas, err = pd.getInt32Array()
  24. if err != nil {
  25. return err
  26. }
  27. pm.Isr, err = pd.getInt32Array()
  28. if err != nil {
  29. return err
  30. }
  31. return nil
  32. }
  33. func (pm *PartitionMetadata) encode(pe packetEncoder) (err error) {
  34. pe.putInt16(int16(pm.Err))
  35. pe.putInt32(pm.ID)
  36. pe.putInt32(pm.Leader)
  37. err = pe.putInt32Array(pm.Replicas)
  38. if err != nil {
  39. return err
  40. }
  41. err = pe.putInt32Array(pm.Isr)
  42. if err != nil {
  43. return err
  44. }
  45. return nil
  46. }
  47. type TopicMetadata struct {
  48. Err KError
  49. Name string
  50. IsInternal bool // Only valid for Version >= 1
  51. Partitions []*PartitionMetadata
  52. }
  53. func (tm *TopicMetadata) decode(pd packetDecoder, version int16) (err error) {
  54. tmp, err := pd.getInt16()
  55. if err != nil {
  56. return err
  57. }
  58. tm.Err = KError(tmp)
  59. tm.Name, err = pd.getString()
  60. if err != nil {
  61. return err
  62. }
  63. if version >= 1 {
  64. tm.IsInternal, err = pd.getBool()
  65. if err != nil {
  66. return err
  67. }
  68. }
  69. n, err := pd.getArrayLength()
  70. if err != nil {
  71. return err
  72. }
  73. tm.Partitions = make([]*PartitionMetadata, n)
  74. for i := 0; i < n; i++ {
  75. tm.Partitions[i] = new(PartitionMetadata)
  76. err = tm.Partitions[i].decode(pd)
  77. if err != nil {
  78. return err
  79. }
  80. }
  81. return nil
  82. }
  83. func (tm *TopicMetadata) encode(pe packetEncoder, version int16) (err error) {
  84. pe.putInt16(int16(tm.Err))
  85. err = pe.putString(tm.Name)
  86. if err != nil {
  87. return err
  88. }
  89. if version >= 1 {
  90. pe.putBool(tm.IsInternal)
  91. }
  92. err = pe.putArrayLength(len(tm.Partitions))
  93. if err != nil {
  94. return err
  95. }
  96. for _, pm := range tm.Partitions {
  97. err = pm.encode(pe)
  98. if err != nil {
  99. return err
  100. }
  101. }
  102. return nil
  103. }
  104. type MetadataResponse struct {
  105. Version int16
  106. Brokers []*Broker
  107. ControllerID int32
  108. Topics []*TopicMetadata
  109. }
  110. func (r *MetadataResponse) decode(pd packetDecoder, version int16) (err error) {
  111. n, err := pd.getArrayLength()
  112. if err != nil {
  113. return err
  114. }
  115. r.Brokers = make([]*Broker, n)
  116. for i := 0; i < n; i++ {
  117. r.Brokers[i] = new(Broker)
  118. err = r.Brokers[i].decode(pd, version)
  119. if err != nil {
  120. return err
  121. }
  122. }
  123. if version >= 1 {
  124. r.ControllerID, err = pd.getInt32()
  125. if err != nil {
  126. return err
  127. }
  128. } else {
  129. r.ControllerID = -1
  130. }
  131. n, err = pd.getArrayLength()
  132. if err != nil {
  133. return err
  134. }
  135. r.Topics = make([]*TopicMetadata, n)
  136. for i := 0; i < n; i++ {
  137. r.Topics[i] = new(TopicMetadata)
  138. err = r.Topics[i].decode(pd, version)
  139. if err != nil {
  140. return err
  141. }
  142. }
  143. return nil
  144. }
  145. func (r *MetadataResponse) encode(pe packetEncoder) error {
  146. err := pe.putArrayLength(len(r.Brokers))
  147. if err != nil {
  148. return err
  149. }
  150. for _, broker := range r.Brokers {
  151. err = broker.encode(pe, r.Version)
  152. if err != nil {
  153. return err
  154. }
  155. }
  156. if r.Version >= 1 {
  157. pe.putInt32(r.ControllerID)
  158. }
  159. err = pe.putArrayLength(len(r.Topics))
  160. if err != nil {
  161. return err
  162. }
  163. for _, tm := range r.Topics {
  164. err = tm.encode(pe, r.Version)
  165. if err != nil {
  166. return err
  167. }
  168. }
  169. return nil
  170. }
  171. func (r *MetadataResponse) key() int16 {
  172. return 3
  173. }
  174. func (r *MetadataResponse) version() int16 {
  175. return 0
  176. }
  177. func (r *MetadataResponse) requiredVersion() KafkaVersion {
  178. return MinVersion
  179. }
  180. // testing API
  181. func (r *MetadataResponse) AddBroker(addr string, id int32) {
  182. r.Brokers = append(r.Brokers, &Broker{id: id, addr: addr})
  183. }
  184. func (r *MetadataResponse) AddTopic(topic string, err KError) *TopicMetadata {
  185. var tmatch *TopicMetadata
  186. for _, tm := range r.Topics {
  187. if tm.Name == topic {
  188. tmatch = tm
  189. goto foundTopic
  190. }
  191. }
  192. tmatch = new(TopicMetadata)
  193. tmatch.Name = topic
  194. r.Topics = append(r.Topics, tmatch)
  195. foundTopic:
  196. tmatch.Err = err
  197. return tmatch
  198. }
  199. func (r *MetadataResponse) AddTopicPartition(topic string, partition, brokerID int32, replicas, isr []int32, err KError) {
  200. tmatch := r.AddTopic(topic, ErrNoError)
  201. var pmatch *PartitionMetadata
  202. for _, pm := range tmatch.Partitions {
  203. if pm.ID == partition {
  204. pmatch = pm
  205. goto foundPartition
  206. }
  207. }
  208. pmatch = new(PartitionMetadata)
  209. pmatch.ID = partition
  210. tmatch.Partitions = append(tmatch.Partitions, pmatch)
  211. foundPartition:
  212. pmatch.Leader = brokerID
  213. pmatch.Replicas = replicas
  214. pmatch.Isr = isr
  215. pmatch.Err = err
  216. }