message_queue_test.go 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239
  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.Subscribe("test1", "test-mqtt",
  100. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  101. if event.ID != "1" {
  102. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  103. }
  104. if event.Type != "test" {
  105. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  106. }
  107. if string(event.Data) != "test-message" {
  108. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  109. }
  110. fmt.Println("mqtt test1 consumed")
  111. wg.Done()
  112. return nil
  113. })
  114. if err != nil {
  115. t.Fatalf("%+v\n", err)
  116. }
  117. err = mqttMessageQueue.Subscribe("test2", "test-mqtt",
  118. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  119. if event.ID != "1" {
  120. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  121. }
  122. if event.Type != "test" {
  123. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  124. }
  125. if string(event.Data) != "test-message" {
  126. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  127. }
  128. fmt.Println("mqtt test2 consumed")
  129. wg.Done()
  130. return nil
  131. })
  132. if err != nil {
  133. t.Fatalf("%+v\n", err)
  134. }
  135. err = mqttMessageQueue.Publish("test-mqtt",
  136. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  137. if err != nil {
  138. t.Fatalf("%+v\n", err)
  139. }
  140. wg.Wait()
  141. }
  142. func testMessageQueue(t *testing.T, messageQueue common.MessageQueue) {
  143. wg := sync.WaitGroup{}
  144. wg.Add(2)
  145. err := message_queue.Subscribe(messageQueue, "test1", "test-message-queue",
  146. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  147. if event.ID != "1" {
  148. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  149. }
  150. if event.Type != "test" {
  151. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  152. }
  153. if string(event.Data) != "test-message" {
  154. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  155. }
  156. fmt.Println("test1 consumed")
  157. wg.Done()
  158. return nil
  159. })
  160. if err != nil {
  161. t.Fatalf("%+v\n", err)
  162. }
  163. err = message_queue.Subscribe(messageQueue, "test2", "test-message-queue",
  164. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) error {
  165. if event.ID != "1" {
  166. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  167. }
  168. if event.Type != "test" {
  169. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  170. }
  171. if string(event.Data) != "test-message" {
  172. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  173. }
  174. fmt.Println("test2 consumed")
  175. wg.Done()
  176. return nil
  177. })
  178. if err != nil {
  179. t.Fatalf("%+v\n", err)
  180. }
  181. err = message_queue.Publish(messageQueue, "test-message-queue",
  182. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  183. if err != nil {
  184. t.Fatalf("%+v\n", err)
  185. }
  186. wg.Wait()
  187. }