metadata_response.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227
  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 (m *MetadataResponse) decode(pd packetDecoder) (err error) {
  99. n, err := pd.getArrayLength()
  100. if err != nil {
  101. return err
  102. }
  103. m.Brokers = make([]*Broker, n)
  104. for i := 0; i < n; i++ {
  105. m.Brokers[i] = new(Broker)
  106. err = m.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. m.Topics = make([]*TopicMetadata, n)
  116. for i := 0; i < n; i++ {
  117. m.Topics[i] = new(TopicMetadata)
  118. err = m.Topics[i].decode(pd)
  119. if err != nil {
  120. return err
  121. }
  122. }
  123. return nil
  124. }
  125. func (m *MetadataResponse) encode(pe packetEncoder) error {
  126. err := pe.putArrayLength(len(m.Brokers))
  127. if err != nil {
  128. return err
  129. }
  130. for _, broker := range m.Brokers {
  131. err = broker.encode(pe)
  132. if err != nil {
  133. return err
  134. }
  135. }
  136. err = pe.putArrayLength(len(m.Topics))
  137. if err != nil {
  138. return err
  139. }
  140. for _, tm := range m.Topics {
  141. err = tm.encode(pe)
  142. if err != nil {
  143. return err
  144. }
  145. }
  146. return nil
  147. }
  148. // testing API
  149. func (m *MetadataResponse) AddBroker(addr string, id int32) {
  150. m.Brokers = append(m.Brokers, &Broker{id: id, addr: addr})
  151. }
  152. func (m *MetadataResponse) AddTopic(topic string, err KError) *TopicMetadata {
  153. var tmatch *TopicMetadata
  154. for _, tm := range m.Topics {
  155. if tm.Name == topic {
  156. tmatch = tm
  157. goto foundTopic
  158. }
  159. }
  160. tmatch = new(TopicMetadata)
  161. tmatch.Name = topic
  162. m.Topics = append(m.Topics, tmatch)
  163. foundTopic:
  164. tmatch.Err = err
  165. return tmatch
  166. }
  167. func (m *MetadataResponse) AddTopicPartition(topic string, partition, brokerID int32, replicas, isr []int32, err KError) {
  168. tmatch := m.AddTopic(topic, ErrNoError)
  169. var pmatch *PartitionMetadata
  170. for _, pm := range tmatch.Partitions {
  171. if pm.ID == partition {
  172. pmatch = pm
  173. goto foundPartition
  174. }
  175. }
  176. pmatch = new(PartitionMetadata)
  177. pmatch.ID = partition
  178. tmatch.Partitions = append(tmatch.Partitions, pmatch)
  179. foundPartition:
  180. pmatch.Leader = brokerID
  181. pmatch.Replicas = replicas
  182. pmatch.Isr = isr
  183. pmatch.Err = err
  184. }