message_queue_test.go 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118
  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/redis"
  8. "github.com/pkg/errors"
  9. "sync"
  10. "testing"
  11. )
  12. func TestRedisMessageQueue(t *testing.T) {
  13. redisMessageQueue := redis.New("localhost:30379", "", "mtyzxhc", 1,
  14. common.WithMaxLen(1000),
  15. common.WithConsumerNum(1))
  16. defer redis.Destroy(redisMessageQueue)
  17. testRedisMessageQueue(t, redisMessageQueue)
  18. }
  19. func TestMessageQueueInfrastructure(t *testing.T) {
  20. i := infrastructure.NewInfrastructure(infrastructure.Config{
  21. MessageQueueConfig: infrastructure.MessageQueueConfig{
  22. Redis: &infrastructure.RedisConfig{
  23. Address: "localhost:30379",
  24. Password: "mtyzxhc",
  25. DB: 1,
  26. },
  27. MaxLen: 1000,
  28. ConsumerNum: 1,
  29. },
  30. })
  31. defer infrastructure.DestroyInfrastructure(i)
  32. testMessageQueue(t, i.RedisMessageQueue())
  33. }
  34. func testRedisMessageQueue(t *testing.T, redisMessageQueue *redis.MessageQueue) {
  35. wg := sync.WaitGroup{}
  36. wg.Add(2)
  37. err := redisMessageQueue.Subscribe("test1", "test-redis",
  38. func(queue common.MessageQueue, topic string, data string) {
  39. if string(data) != "test-message" {
  40. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  41. }
  42. fmt.Println("test1 consumed")
  43. wg.Done()
  44. })
  45. if err != nil {
  46. t.Fatalf("%+v\n", err)
  47. }
  48. err = redisMessageQueue.Subscribe("test2", "test-redis",
  49. func(queue common.MessageQueue, topic string, data string) {
  50. if string(data) != "test-message" {
  51. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  52. }
  53. fmt.Println("test2 consumed")
  54. wg.Done()
  55. })
  56. if err != nil {
  57. t.Fatalf("%+v\n", err)
  58. }
  59. err = redisMessageQueue.Publish("test-redis", "test-message")
  60. if err != nil {
  61. t.Fatalf("%+v\n", err)
  62. }
  63. wg.Wait()
  64. }
  65. func testMessageQueue(t *testing.T, redisMessageQueue common.MessageQueue) {
  66. wg := sync.WaitGroup{}
  67. wg.Add(2)
  68. err := message_queue.Subscribe(redisMessageQueue, "test1", "test-message-queue",
  69. func(queue common.MessageQueue, topic string, data string) {
  70. if string(data) != "test-message" {
  71. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  72. }
  73. fmt.Println("test1 consumed")
  74. wg.Done()
  75. })
  76. if err != nil {
  77. t.Fatalf("%+v\n", err)
  78. }
  79. err = message_queue.Subscribe(redisMessageQueue, "test2", "test-message-queue",
  80. func(queue common.MessageQueue, topic string, data string) {
  81. if string(data) != "test-message" {
  82. t.Fatalf("%+v\n", errors.New("消息数据不一致"))
  83. }
  84. fmt.Println("test2 consumed")
  85. wg.Done()
  86. })
  87. if err != nil {
  88. t.Fatalf("%+v\n", err)
  89. }
  90. err = message_queue.Publish(redisMessageQueue, "test-message-queue", "test-message")
  91. if err != nil {
  92. t.Fatalf("%+v\n", err)
  93. }
  94. wg.Wait()
  95. }