task_lifecycle_consumer.go 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136
  1. package rule
  2. import (
  3. "encoding/json"
  4. "github.com/streadway/amqp"
  5. "sparrow/pkg/server"
  6. )
  7. type ExternalConsumer interface {
  8. AddMessageHandle(msg *TaskLifecycleMessage) error
  9. RemoveMessageHandle(msg *TaskLifecycleMessage) error
  10. UpdateMessageHandle(msg *TaskLifecycleMessage) error
  11. SnapMessageHandle(msg *TaskLifecycleMessage) error
  12. StartMessageHandle(msg *TaskLifecycleMessage) error
  13. StopMessageHandle(msg *TaskLifecycleMessage) error
  14. }
  15. const TaskLifecycleExchange = "task_lifecycle_exchange"
  16. const TaskExchange = "task_exchange"
  17. type TaskLifecycleConsumer struct {
  18. ch *amqp.Channel
  19. ec ExternalConsumer
  20. stopChan chan struct{}
  21. messageChan chan []byte
  22. }
  23. func NewTaskLifecycleConsumer(ch *amqp.Channel) *TaskLifecycleConsumer {
  24. return &TaskLifecycleConsumer{
  25. ch: ch,
  26. stopChan: make(chan struct{}),
  27. messageChan: make(chan []byte, 10),
  28. }
  29. }
  30. func (a *TaskLifecycleConsumer) SetExternalConsumer(ec ExternalConsumer) {
  31. a.ec = ec
  32. }
  33. func (a *TaskLifecycleConsumer) Stop() {
  34. close(a.stopChan)
  35. }
  36. func (a *TaskLifecycleConsumer) Start() error {
  37. go func() {
  38. for {
  39. select {
  40. case <-a.stopChan:
  41. return
  42. case msg := <-a.messageChan:
  43. a.handleMessage(msg)
  44. }
  45. }
  46. }()
  47. return nil
  48. }
  49. func (a *TaskLifecycleConsumer) handleMessage(msg []byte) {
  50. if a.ec == nil {
  51. server.Log.Errorf("TaskLifecycleConsumer is unset")
  52. return
  53. }
  54. var tm TaskLifecycleMessage
  55. err := json.Unmarshal(msg, &tm)
  56. if err != nil {
  57. server.Log.Errorf("handle lifecycle message error :%v", err)
  58. return
  59. }
  60. switch tm.Action {
  61. case "add":
  62. err = a.ec.AddMessageHandle(&tm)
  63. case "remove":
  64. err = a.ec.RemoveMessageHandle(&tm)
  65. case "update":
  66. err = a.ec.UpdateMessageHandle(&tm)
  67. case "snap":
  68. err = a.ec.SnapMessageHandle(&tm)
  69. case "start":
  70. err = a.ec.StartMessageHandle(&tm)
  71. case "stop":
  72. err = a.ec.StopMessageHandle(&tm)
  73. }
  74. }
  75. func (a *TaskLifecycleConsumer) Init() error {
  76. err := a.ch.ExchangeDeclare(
  77. TaskLifecycleExchange, // name
  78. "fanout", // type
  79. true, // durable
  80. false, // auto-deleted
  81. false, // internal
  82. false, // no-wait
  83. nil, // arguments
  84. )
  85. if err != nil {
  86. return err
  87. }
  88. //绑定queue到交换机
  89. q, err := a.ch.QueueDeclare(
  90. "", // name
  91. false, // durable
  92. false, // delete when unused
  93. true, // exclusive
  94. false, // no-wait
  95. nil, // arguments
  96. )
  97. if err != nil {
  98. return err
  99. }
  100. err = a.ch.QueueBind(
  101. q.Name, // queue name
  102. "", // routing key
  103. TaskLifecycleExchange,
  104. false,
  105. nil,
  106. )
  107. if err != nil {
  108. return err
  109. }
  110. // creat consumer
  111. msg, err := a.ch.Consume(q.Name, "", true, false, false, false, nil)
  112. if err != nil {
  113. return err
  114. }
  115. go func() {
  116. for {
  117. select {
  118. case <-a.stopChan:
  119. return
  120. case d := <-msg:
  121. a.messageChan <- d.Body
  122. }
  123. }
  124. }()
  125. return nil
  126. }