message_queue_test.go 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245
  1. package test
  2. import (
  3. "fmt"
  4. "git.sxidc.com/go-framework/baize/framework/core/data_protocol"
  5. "git.sxidc.com/go-framework/baize/framework/core/infrastructure"
  6. "git.sxidc.com/go-framework/baize/framework/core/infrastructure/message_queue"
  7. "git.sxidc.com/go-framework/baize/framework/core/infrastructure/message_queue/common"
  8. "git.sxidc.com/go-framework/baize/framework/core/infrastructure/message_queue/mqtt"
  9. "git.sxidc.com/go-framework/baize/framework/core/infrastructure/message_queue/redis"
  10. "github.com/pkg/errors"
  11. "sync"
  12. "testing"
  13. )
  14. func TestRedisMessageQueue(t *testing.T) {
  15. redisMessageQueue := redis.New("localhost:30379", "", "mtyzxhc", 1,
  16. redis.WithMaxLen(1000),
  17. redis.WithConsumerNum(1))
  18. defer redis.Destroy(redisMessageQueue)
  19. testRedisMessageQueue(t, redisMessageQueue)
  20. }
  21. func TestMqttMessageQueue(t *testing.T) {
  22. mqttMessageQueue, err := mqtt.New("localhost:1883", "admin", "mtyzxhc")
  23. if err != nil {
  24. t.Fatalf("%+v\n", err)
  25. }
  26. defer mqtt.Destroy(mqttMessageQueue)
  27. testMqttMessageQueue(t, mqttMessageQueue)
  28. }
  29. func TestMessageQueueInfrastructure(t *testing.T) {
  30. i := infrastructure.NewInfrastructure(infrastructure.Config{
  31. MessageQueueConfig: infrastructure.MessageQueueConfig{
  32. Redis: &infrastructure.RedisConfig{
  33. Address: "localhost:30379",
  34. Password: "mtyzxhc",
  35. DB: 1,
  36. MaxLen: 1000,
  37. ConsumerNum: 1,
  38. },
  39. Mqtt: &infrastructure.MqttConfig{
  40. Address: "localhost:1883",
  41. UserName: "admin",
  42. Password: "mtyzxhc",
  43. },
  44. },
  45. })
  46. defer infrastructure.DestroyInfrastructure(i)
  47. testMessageQueue(t, i.RedisMessageQueue())
  48. testMessageQueue(t, i.MqttMessageQueue())
  49. }
  50. func testRedisMessageQueue(t *testing.T, redisMessageQueue *redis.MessageQueue) {
  51. wg := sync.WaitGroup{}
  52. wg.Add(2)
  53. err := redisMessageQueue.Subscribe("test1", "test-redis",
  54. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  55. if event.ID != "1" {
  56. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  57. }
  58. if event.Type != "test" {
  59. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  60. }
  61. if string(event.Data) != "test-message" {
  62. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  63. }
  64. fmt.Println("redis test1 consumed")
  65. wg.Done()
  66. return nil
  67. })
  68. if err != nil {
  69. t.Fatalf("%+v\n", err)
  70. }
  71. err = redisMessageQueue.Subscribe("test2", "test-redis",
  72. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  73. if event.ID != "1" {
  74. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  75. }
  76. if event.Type != "test" {
  77. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  78. }
  79. if string(event.Data) != "test-message" {
  80. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  81. }
  82. fmt.Println("redis test2 consumed")
  83. wg.Done()
  84. return nil
  85. })
  86. if err != nil {
  87. t.Fatalf("%+v\n", err)
  88. }
  89. err = redisMessageQueue.Publish("test-redis",
  90. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  91. if err != nil {
  92. t.Fatalf("%+v\n", err)
  93. }
  94. wg.Wait()
  95. }
  96. func testMqttMessageQueue(t *testing.T, mqttMessageQueue *mqtt.MessageQueue) {
  97. wg := sync.WaitGroup{}
  98. wg.Add(2)
  99. //err := mqttMessageQueue.Publish("test-mqtt",
  100. // data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  101. //if err != nil {
  102. // t.Fatalf("%+v\n", err)
  103. //}
  104. err := mqttMessageQueue.Subscribe("test1", "test-mqtt",
  105. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  106. if event.ID != "1" {
  107. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  108. }
  109. if event.Type != "test" {
  110. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  111. }
  112. if string(event.Data) != "test-message" {
  113. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  114. }
  115. fmt.Println("mqtt test1 consumed")
  116. wg.Done()
  117. return nil
  118. })
  119. if err != nil {
  120. t.Fatalf("%+v\n", err)
  121. }
  122. err = mqttMessageQueue.Subscribe("test2", "test-mqtt",
  123. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  124. if event.ID != "1" {
  125. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  126. }
  127. if event.Type != "test" {
  128. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  129. }
  130. if string(event.Data) != "test-message" {
  131. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  132. }
  133. fmt.Println("mqtt test2 consumed")
  134. wg.Done()
  135. return nil
  136. })
  137. if err != nil {
  138. t.Fatalf("%+v\n", err)
  139. }
  140. err = mqttMessageQueue.Publish("test-mqtt",
  141. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  142. if err != nil {
  143. t.Fatalf("%+v\n", err)
  144. }
  145. wg.Wait()
  146. }
  147. func testMessageQueue(t *testing.T, messageQueue common.MessageQueue) {
  148. wg := sync.WaitGroup{}
  149. wg.Add(2)
  150. err := message_queue.Subscribe(messageQueue, "test1", "test-message-queue",
  151. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  152. if event.ID != "1" {
  153. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  154. }
  155. if event.Type != "test" {
  156. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  157. }
  158. if string(event.Data) != "test-message" {
  159. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  160. }
  161. fmt.Println("test1 consumed")
  162. wg.Done()
  163. return nil
  164. })
  165. if err != nil {
  166. t.Fatalf("%+v\n", err)
  167. }
  168. err = message_queue.Subscribe(messageQueue, "test2", "test-message-queue",
  169. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  170. if event.ID != "1" {
  171. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  172. }
  173. if event.Type != "test" {
  174. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  175. }
  176. if string(event.Data) != "test-message" {
  177. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  178. }
  179. fmt.Println("test2 consumed")
  180. wg.Done()
  181. return nil
  182. })
  183. if err != nil {
  184. t.Fatalf("%+v\n", err)
  185. }
  186. err = message_queue.Publish(messageQueue, "test-message-queue",
  187. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  188. if err != nil {
  189. t.Fatalf("%+v\n", err)
  190. }
  191. wg.Wait()
  192. }