message_queue_test.go 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227
  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) {
  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. })
  67. if err != nil {
  68. t.Fatalf("%+v\n", err)
  69. }
  70. err = redisMessageQueue.Subscribe("test2", "test-redis",
  71. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) {
  72. if event.ID != "1" {
  73. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  74. }
  75. if event.Type != "test" {
  76. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  77. }
  78. if string(event.Data) != "test-message" {
  79. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  80. }
  81. fmt.Println("redis test2 consumed")
  82. wg.Done()
  83. })
  84. if err != nil {
  85. t.Fatalf("%+v\n", err)
  86. }
  87. err = redisMessageQueue.Publish("test-redis",
  88. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  89. if err != nil {
  90. t.Fatalf("%+v\n", err)
  91. }
  92. wg.Wait()
  93. }
  94. func testMqttMessageQueue(t *testing.T, mqttMessageQueue *mqtt.MessageQueue) {
  95. wg := sync.WaitGroup{}
  96. wg.Add(2)
  97. err := mqttMessageQueue.Subscribe("test1", "test-mqtt",
  98. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) {
  99. if event.ID != "1" {
  100. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  101. }
  102. if event.Type != "test" {
  103. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  104. }
  105. if string(event.Data) != "test-message" {
  106. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  107. }
  108. fmt.Println("mqtt test1 consumed")
  109. wg.Done()
  110. })
  111. if err != nil {
  112. t.Fatalf("%+v\n", err)
  113. }
  114. err = mqttMessageQueue.Subscribe("test2", "test-mqtt",
  115. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) {
  116. if event.ID != "1" {
  117. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  118. }
  119. if event.Type != "test" {
  120. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  121. }
  122. if string(event.Data) != "test-message" {
  123. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  124. }
  125. fmt.Println("mqtt test2 consumed")
  126. wg.Done()
  127. })
  128. if err != nil {
  129. t.Fatalf("%+v\n", err)
  130. }
  131. err = mqttMessageQueue.Publish("test-mqtt",
  132. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  133. if err != nil {
  134. t.Fatalf("%+v\n", err)
  135. }
  136. wg.Wait()
  137. }
  138. func testMessageQueue(t *testing.T, messageQueue common.MessageQueue) {
  139. wg := sync.WaitGroup{}
  140. wg.Add(2)
  141. err := message_queue.Subscribe(messageQueue, "test1", "test-message-queue",
  142. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) {
  143. if event.ID != "1" {
  144. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  145. }
  146. if event.Type != "test" {
  147. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  148. }
  149. if string(event.Data) != "test-message" {
  150. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  151. }
  152. fmt.Println("test1 consumed")
  153. wg.Done()
  154. })
  155. if err != nil {
  156. t.Fatalf("%+v\n", err)
  157. }
  158. err = message_queue.Subscribe(messageQueue, "test2", "test-message-queue",
  159. func(queue common.MessageQueue, topic string, event *data_protocol.CloudEvent) {
  160. if event.ID != "1" {
  161. t.Fatalf("%+v\n", errors.New("消息ID不一致"))
  162. }
  163. if event.Type != "test" {
  164. t.Fatalf("%+v\n", errors.New("消息类型不一致"))
  165. }
  166. if string(event.Data) != "test-message" {
  167. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  168. }
  169. fmt.Println("test2 consumed")
  170. wg.Done()
  171. })
  172. if err != nil {
  173. t.Fatalf("%+v\n", err)
  174. }
  175. err = message_queue.Publish(messageQueue, "test-message-queue",
  176. data_protocol.NewCloudEvent("1", "test", "baize-test.com", "application/text", []byte("test-message")))
  177. if err != nil {
  178. t.Fatalf("%+v\n", err)
  179. }
  180. wg.Wait()
  181. }