broker_test.go 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178
  1. package sarama
  2. import (
  3. "fmt"
  4. "testing"
  5. )
  6. func ExampleBroker() error {
  7. broker := NewBroker("localhost:9092")
  8. err := broker.Open(nil)
  9. if err != nil {
  10. return err
  11. }
  12. defer broker.Close()
  13. request := MetadataRequest{Topics: []string{"myTopic"}}
  14. response, err := broker.GetMetadata("myClient", &request)
  15. if err != nil {
  16. return err
  17. }
  18. fmt.Println("There are", len(response.Topics), "topics active in the cluster.")
  19. return nil
  20. }
  21. type mockEncoder struct {
  22. bytes []byte
  23. }
  24. func (m mockEncoder) encode(pe packetEncoder) error {
  25. pe.putRawBytes(m.bytes)
  26. return nil
  27. }
  28. func TestBrokerAccessors(t *testing.T) {
  29. broker := NewBroker("abc:123")
  30. if broker.ID() != -1 {
  31. t.Error("New broker didn't have an ID of -1.")
  32. }
  33. if broker.Addr() != "abc:123" {
  34. t.Error("New broker didn't have the correct address")
  35. }
  36. broker.id = 34
  37. if broker.ID() != 34 {
  38. t.Error("Manually setting broker ID did not take effect.")
  39. }
  40. }
  41. func TestSimpleBrokerCommunication(t *testing.T) {
  42. mb := NewMockBroker(t, 0)
  43. defer mb.Close()
  44. broker := NewBroker(mb.Addr())
  45. err := broker.Open(nil)
  46. if err != nil {
  47. t.Fatal(err)
  48. }
  49. for _, tt := range brokerTestTable {
  50. mb.Returns(&mockEncoder{tt.response})
  51. }
  52. for _, tt := range brokerTestTable {
  53. tt.runner(t, broker)
  54. }
  55. err = broker.Close()
  56. if err != nil {
  57. t.Error(err)
  58. }
  59. }
  60. // We're not testing encoding/decoding here, so most of the requests/responses will be empty for simplicity's sake
  61. var brokerTestTable = []struct {
  62. response []byte
  63. runner func(*testing.T, *Broker)
  64. }{
  65. {[]byte{0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00},
  66. func(t *testing.T, broker *Broker) {
  67. request := MetadataRequest{}
  68. response, err := broker.GetMetadata("clientID", &request)
  69. if err != nil {
  70. t.Error(err)
  71. }
  72. if response == nil {
  73. t.Error("Metadata request got no response!")
  74. }
  75. }},
  76. {[]byte{0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 't', 0x00, 0x00, 0x00, 0x00},
  77. func(t *testing.T, broker *Broker) {
  78. request := ConsumerMetadataRequest{}
  79. response, err := broker.GetConsumerMetadata("clientID", &request)
  80. if err != nil {
  81. t.Error(err)
  82. }
  83. if response == nil {
  84. t.Error("Consumer Metadata request got no response!")
  85. }
  86. }},
  87. {[]byte{},
  88. func(t *testing.T, broker *Broker) {
  89. request := ProduceRequest{}
  90. request.RequiredAcks = NoResponse
  91. response, err := broker.Produce("clientID", &request)
  92. if err != nil {
  93. t.Error(err)
  94. }
  95. if response != nil {
  96. t.Error("Produce request with NoResponse got a response!")
  97. }
  98. }},
  99. {[]byte{0x00, 0x00, 0x00, 0x00},
  100. func(t *testing.T, broker *Broker) {
  101. request := ProduceRequest{}
  102. request.RequiredAcks = WaitForLocal
  103. response, err := broker.Produce("clientID", &request)
  104. if err != nil {
  105. t.Error(err)
  106. }
  107. if response == nil {
  108. t.Error("Produce request without NoResponse got no response!")
  109. }
  110. }},
  111. {[]byte{0x00, 0x00, 0x00, 0x00},
  112. func(t *testing.T, broker *Broker) {
  113. request := FetchRequest{}
  114. response, err := broker.Fetch("clientID", &request)
  115. if err != nil {
  116. t.Error(err)
  117. }
  118. if response == nil {
  119. t.Error("Fetch request got no response!")
  120. }
  121. }},
  122. {[]byte{0x00, 0x00, 0x00, 0x00},
  123. func(t *testing.T, broker *Broker) {
  124. request := OffsetFetchRequest{}
  125. response, err := broker.FetchOffset("clientID", &request)
  126. if err != nil {
  127. t.Error(err)
  128. }
  129. if response == nil {
  130. t.Error("OffsetFetch request got no response!")
  131. }
  132. }},
  133. {[]byte{0x00, 0x00, 0x00, 0x00},
  134. func(t *testing.T, broker *Broker) {
  135. request := OffsetCommitRequest{}
  136. response, err := broker.CommitOffset("clientID", &request)
  137. if err != nil {
  138. t.Error(err)
  139. }
  140. if response == nil {
  141. t.Error("OffsetCommit request got no response!")
  142. }
  143. }},
  144. {[]byte{0x00, 0x00, 0x00, 0x00},
  145. func(t *testing.T, broker *Broker) {
  146. request := OffsetRequest{}
  147. response, err := broker.GetAvailableOffsets("clientID", &request)
  148. if err != nil {
  149. t.Error(err)
  150. }
  151. if response == nil {
  152. t.Error("Offset request got no response!")
  153. }
  154. }},
  155. }