metadata_response.go 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239
  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. Partitions []*PartitionMetadata
  51. }
  52. func (tm *TopicMetadata) decode(pd packetDecoder) (err error) {
  53. tmp, err := pd.getInt16()
  54. if err != nil {
  55. return err
  56. }
  57. tm.Err = KError(tmp)
  58. tm.Name, err = pd.getString()
  59. if err != nil {
  60. return err
  61. }
  62. n, err := pd.getArrayLength()
  63. if err != nil {
  64. return err
  65. }
  66. tm.Partitions = make([]*PartitionMetadata, n)
  67. for i := 0; i < n; i++ {
  68. tm.Partitions[i] = new(PartitionMetadata)
  69. err = tm.Partitions[i].decode(pd)
  70. if err != nil {
  71. return err
  72. }
  73. }
  74. return nil
  75. }
  76. func (tm *TopicMetadata) encode(pe packetEncoder) (err error) {
  77. pe.putInt16(int16(tm.Err))
  78. err = pe.putString(tm.Name)
  79. if err != nil {
  80. return err
  81. }
  82. err = pe.putArrayLength(len(tm.Partitions))
  83. if err != nil {
  84. return err
  85. }
  86. for _, pm := range tm.Partitions {
  87. err = pm.encode(pe)
  88. if err != nil {
  89. return err
  90. }
  91. }
  92. return nil
  93. }
  94. type MetadataResponse struct {
  95. Brokers []*Broker
  96. Topics []*TopicMetadata
  97. }
  98. func (r *MetadataResponse) decode(pd packetDecoder, version int16) (err error) {
  99. n, err := pd.getArrayLength()
  100. if err != nil {
  101. return err
  102. }
  103. r.Brokers = make([]*Broker, n)
  104. for i := 0; i < n; i++ {
  105. r.Brokers[i] = new(Broker)
  106. err = r.Brokers[i].decode(pd)
  107. if err != nil {
  108. return err
  109. }
  110. }
  111. n, err = pd.getArrayLength()
  112. if err != nil {
  113. return err
  114. }
  115. r.Topics = make([]*TopicMetadata, n)
  116. for i := 0; i < n; i++ {
  117. r.Topics[i] = new(TopicMetadata)
  118. err = r.Topics[i].decode(pd)
  119. if err != nil {
  120. return err
  121. }
  122. }
  123. return nil
  124. }
  125. func (r *MetadataResponse) encode(pe packetEncoder) error {
  126. err := pe.putArrayLength(len(r.Brokers))
  127. if err != nil {
  128. return err
  129. }
  130. for _, broker := range r.Brokers {
  131. err = broker.encode(pe)
  132. if err != nil {
  133. return err
  134. }
  135. }
  136. err = pe.putArrayLength(len(r.Topics))
  137. if err != nil {
  138. return err
  139. }
  140. for _, tm := range r.Topics {
  141. err = tm.encode(pe)
  142. if err != nil {
  143. return err
  144. }
  145. }
  146. return nil
  147. }
  148. func (r *MetadataResponse) key() int16 {
  149. return 3
  150. }
  151. func (r *MetadataResponse) version() int16 {
  152. return 0
  153. }
  154. func (r *MetadataResponse) requiredVersion() KafkaVersion {
  155. return minVersion
  156. }
  157. // testing API
  158. func (r *MetadataResponse) AddBroker(addr string, id int32) {
  159. r.Brokers = append(r.Brokers, &Broker{id: id, addr: addr})
  160. }
  161. func (r *MetadataResponse) AddTopic(topic string, err KError) *TopicMetadata {
  162. var tmatch *TopicMetadata
  163. for _, tm := range r.Topics {
  164. if tm.Name == topic {
  165. tmatch = tm
  166. goto foundTopic
  167. }
  168. }
  169. tmatch = new(TopicMetadata)
  170. tmatch.Name = topic
  171. r.Topics = append(r.Topics, tmatch)
  172. foundTopic:
  173. tmatch.Err = err
  174. return tmatch
  175. }
  176. func (r *MetadataResponse) AddTopicPartition(topic string, partition, brokerID int32, replicas, isr []int32, err KError) {
  177. tmatch := r.AddTopic(topic, ErrNoError)
  178. var pmatch *PartitionMetadata
  179. for _, pm := range tmatch.Partitions {
  180. if pm.ID == partition {
  181. pmatch = pm
  182. goto foundPartition
  183. }
  184. }
  185. pmatch = new(PartitionMetadata)
  186. pmatch.ID = partition
  187. tmatch.Partitions = append(tmatch.Partitions, pmatch)
  188. foundPartition:
  189. pmatch.Leader = brokerID
  190. pmatch.Replicas = replicas
  191. pmatch.Isr = isr
  192. pmatch.Err = err
  193. }