message_queue_test.go 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191
  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. redis.WithInFlightCount(100),
  18. redis.WithVisibilityAgainSec(5))
  19. defer redis.Destroy(redisMessageQueue)
  20. testRedisMessageQueue(t, redisMessageQueue)
  21. }
  22. func TestMqttMessageQueue(t *testing.T) {
  23. mqttMessageQueue, err := mqtt.New("localhost:1883", "admin", "mtyzxhc")
  24. if err != nil {
  25. t.Fatalf("%+v\n", err)
  26. }
  27. defer mqtt.Destroy(mqttMessageQueue)
  28. testMqttMessageQueue(t, mqttMessageQueue)
  29. }
  30. func TestMessageQueueInfrastructure(t *testing.T) {
  31. i := infrastructure.NewInfrastructure(infrastructure.Config{
  32. MessageQueueConfig: infrastructure.MessageQueueConfig{
  33. Redis: &infrastructure.RedisConfig{
  34. Address: "localhost:30379",
  35. Password: "mtyzxhc",
  36. DB: 1,
  37. MaxLen: 1000,
  38. ConsumerNum: 1,
  39. VisibilityAgainSec: 5,
  40. InFlightCount: 100,
  41. },
  42. Mqtt: &infrastructure.MqttConfig{
  43. Address: "localhost:1883",
  44. UserName: "admin",
  45. Password: "mtyzxhc",
  46. },
  47. },
  48. })
  49. defer infrastructure.DestroyInfrastructure(i)
  50. testMessageQueue(t, i.RedisMessageQueue())
  51. testMessageQueue(t, i.MqttMessageQueue())
  52. }
  53. func testRedisMessageQueue(t *testing.T, redisMessageQueue *redis.MessageQueue) {
  54. wg := sync.WaitGroup{}
  55. wg.Add(2)
  56. err := redisMessageQueue.Subscribe("test1", "test-redis",
  57. func(queue common.MessageQueue, topic string, data string) error {
  58. if string(data) != "test-message" {
  59. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  60. }
  61. fmt.Println("redis test1 consumed")
  62. wg.Done()
  63. return nil
  64. })
  65. if err != nil {
  66. t.Fatalf("%+v\n", err)
  67. }
  68. err = redisMessageQueue.Subscribe("test2", "test-redis",
  69. func(queue common.MessageQueue, topic string, data string) error {
  70. if string(data) != "test-message" {
  71. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  72. }
  73. fmt.Println("redis test2 consumed")
  74. wg.Done()
  75. return nil
  76. })
  77. if err != nil {
  78. t.Fatalf("%+v\n", err)
  79. }
  80. err = redisMessageQueue.Publish("test-redis", "test-message")
  81. if err != nil {
  82. t.Fatalf("%+v\n", err)
  83. }
  84. wg.Wait()
  85. }
  86. func testMqttMessageQueue(t *testing.T, mqttMessageQueue *mqtt.MessageQueue) {
  87. wg := sync.WaitGroup{}
  88. wg.Add(2)
  89. err := mqttMessageQueue.Subscribe("test1", "test-mqtt",
  90. func(queue common.MessageQueue, topic string, data string) error {
  91. if string(data) != "test-message" {
  92. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  93. }
  94. fmt.Println("mqtt test1 consumed")
  95. wg.Done()
  96. return nil
  97. })
  98. if err != nil {
  99. t.Fatalf("%+v\n", err)
  100. }
  101. err = mqttMessageQueue.Subscribe("test2", "test-mqtt",
  102. func(queue common.MessageQueue, topic string, data string) error {
  103. if string(data) != "test-message" {
  104. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  105. }
  106. fmt.Println("mqtt test2 consumed")
  107. wg.Done()
  108. return nil
  109. })
  110. if err != nil {
  111. t.Fatalf("%+v\n", err)
  112. }
  113. err = mqttMessageQueue.Publish("test-mqtt", "test-message")
  114. if err != nil {
  115. t.Fatalf("%+v\n", err)
  116. }
  117. wg.Wait()
  118. }
  119. func testMessageQueue(t *testing.T, messageQueue common.MessageQueue) {
  120. wg := sync.WaitGroup{}
  121. wg.Add(2)
  122. err := message_queue.Subscribe(messageQueue, "test1", "test-message-queue",
  123. func(queue common.MessageQueue, topic string, data string) error {
  124. if string(data) != "test-message" {
  125. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  126. }
  127. fmt.Println("test1 consumed")
  128. wg.Done()
  129. return nil
  130. })
  131. if err != nil {
  132. t.Fatalf("%+v\n", err)
  133. }
  134. err = message_queue.Subscribe(messageQueue, "test2", "test-message-queue",
  135. func(queue common.MessageQueue, topic string, data string) error {
  136. if string(data) != "test-message" {
  137. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  138. }
  139. fmt.Println("test2 consumed")
  140. wg.Done()
  141. return nil
  142. })
  143. if err != nil {
  144. t.Fatalf("%+v\n", err)
  145. }
  146. err = message_queue.Publish(messageQueue, "test-message-queue", "test-message")
  147. if err != nil {
  148. t.Fatalf("%+v\n", err)
  149. }
  150. wg.Wait()
  151. }