message_queue_test.go 5.5 KB

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