mgr.go 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216
  1. package mqtt
  2. import (
  3. "runtime/debug"
  4. "time"
  5. "github.com/sirupsen/logrus"
  6. "lc/common/util"
  7. )
  8. // OptType 消息发布/订阅方向
  9. type OptType uint8
  10. const (
  11. ToAll OptType = 0 // 同时发布到云端和边缘端
  12. ToCloud OptType = 1 // 仅云端
  13. ToEdge OptType = 2 // 仅边缘端
  14. )
  15. // MQTTMessage 队列中的消息
  16. type MQTTMessage struct {
  17. topic string
  18. payload string
  19. qos QOS
  20. tp OptType
  21. }
  22. // RestartFunc 协程 panic 后重建回调
  23. type RestartFunc func(fn func(args ...interface{}) interface{}, args ...interface{})
  24. // MQTTMgr 管理 Cloud/Edge 双客户端,支持异步队列发布
  25. type MQTTMgr struct {
  26. Cloud *MqttClient
  27. Edge *MqttClient
  28. Queue *util.MlQueue
  29. restartFn RestartFunc
  30. }
  31. // MQTTMgrConfig 管理器构造参数
  32. type MQTTMgrConfig struct {
  33. CloudServer string
  34. CloudClientID string
  35. CloudUser string
  36. CloudPassword string
  37. CloudTimeout uint
  38. CloudOnline BaseMqttOnline // nil 则使用 EmptyMqttOnline
  39. EdgeServer string
  40. EdgeClientID string
  41. EdgeUser string
  42. EdgePassword string
  43. EdgeTimeout uint
  44. EdgeOnline BaseMqttOnline // nil 则使用 EmptyMqttOnline
  45. }
  46. // NewMQTTMgr 创建 MQTT 管理器
  47. func NewMQTTMgr(cfg MQTTMgrConfig) *MQTTMgr {
  48. mgr := &MQTTMgr{
  49. Queue: util.NewQueue(2000),
  50. }
  51. if cfg.CloudOnline == nil {
  52. cfg.CloudOnline = &EmptyMqttOnline{}
  53. }
  54. if cfg.EdgeOnline == nil {
  55. cfg.EdgeOnline = &EmptyMqttOnline{}
  56. }
  57. if cfg.CloudServer != "" {
  58. mgr.Cloud = NewMqttClient(cfg.CloudServer, cfg.CloudClientID,
  59. cfg.CloudUser, cfg.CloudPassword, cfg.CloudTimeout, cfg.CloudOnline)
  60. }
  61. if cfg.EdgeServer != "" {
  62. mgr.Edge = NewMqttClient(cfg.EdgeServer, cfg.EdgeClientID,
  63. cfg.EdgeUser, cfg.EdgePassword, cfg.EdgeTimeout, cfg.EdgeOnline)
  64. }
  65. return mgr
  66. }
  67. // SetRestartFn 设置 panic 恢复后的协程重建回调
  68. func (o *MQTTMgr) SetRestartFn(fn RestartFunc) {
  69. o.restartFn = fn
  70. }
  71. // Subscribe 订阅主题,支持方向选择
  72. func (o *MQTTMgr) Subscribe(topic string, qos QOS, handler MessageHandler, tp OptType) {
  73. logrus.Infof("MQTTMgr.Subscribe:topic=%s,qos=%d,direction=%d", topic, qos, tp)
  74. switch tp {
  75. case ToAll:
  76. if o.Cloud != nil {
  77. o.Cloud.Handle(topic, handler)
  78. if err := o.Cloud.Subscribe(topic, qos); err != nil {
  79. logrus.Errorf("MQTTMgr.Subscribe:Cloud订阅失败,topic=%s,err=%v", topic, err)
  80. }
  81. } else {
  82. logrus.Warnf("MQTTMgr.Subscribe:Cloud客户端为nil,无法订阅topic=%s", topic)
  83. }
  84. if o.Edge != nil {
  85. o.Edge.Handle(topic, handler)
  86. if err := o.Edge.Subscribe(topic, qos); err != nil {
  87. logrus.Errorf("MQTTMgr.Subscribe:Edge订阅失败,topic=%s,err=%v", topic, err)
  88. }
  89. } else {
  90. logrus.Warnf("MQTTMgr.Subscribe:Edge客户端为nil,无法订阅topic=%s", topic)
  91. }
  92. case ToCloud:
  93. if o.Cloud != nil {
  94. o.Cloud.Handle(topic, handler)
  95. if err := o.Cloud.Subscribe(topic, qos); err != nil {
  96. logrus.Errorf("MQTTMgr.Subscribe:Cloud订阅失败,topic=%s,err=%v", topic, err)
  97. }
  98. } else {
  99. logrus.Warnf("MQTTMgr.Subscribe:Cloud客户端为nil,无法订阅topic=%s", topic)
  100. }
  101. case ToEdge:
  102. if o.Edge != nil {
  103. o.Edge.Handle(topic, handler)
  104. if err := o.Edge.Subscribe(topic, qos); err != nil {
  105. logrus.Errorf("MQTTMgr.Subscribe:Edge订阅失败,topic=%s,err=%v", topic, err)
  106. }
  107. } else {
  108. logrus.Warnf("MQTTMgr.Subscribe:Edge客户端为nil,无法订阅topic=%s", topic)
  109. }
  110. }
  111. }
  112. // UnSubscribe 退订主题
  113. func (o *MQTTMgr) UnSubscribe(topic string, tp OptType) {
  114. switch tp {
  115. case ToAll:
  116. if o.Cloud != nil {
  117. _ = o.Cloud.Unsubscribe(topic)
  118. }
  119. if o.Edge != nil {
  120. _ = o.Edge.Unsubscribe(topic)
  121. }
  122. case ToCloud:
  123. if o.Cloud != nil {
  124. _ = o.Cloud.Unsubscribe(topic)
  125. }
  126. case ToEdge:
  127. if o.Edge != nil {
  128. _ = o.Edge.Unsubscribe(topic)
  129. }
  130. }
  131. }
  132. // Publish 异步发布(入队列)
  133. func (o *MQTTMgr) Publish(topic, payload string, qos QOS, tp OptType) {
  134. o.Queue.Put(&MQTTMessage{topic: topic, payload: payload, qos: qos, tp: tp})
  135. }
  136. // doPublish 底层同步发布
  137. func (o *MQTTMgr) doPublish(msg *MQTTMessage) error {
  138. var err error
  139. switch msg.tp {
  140. case ToAll:
  141. if o.Cloud != nil {
  142. err = o.Cloud.PublishString(msg.topic, msg.payload, msg.qos)
  143. }
  144. if o.Edge != nil {
  145. _ = o.Edge.PublishString(msg.topic, msg.payload, msg.qos)
  146. }
  147. case ToCloud:
  148. if o.Cloud != nil {
  149. err = o.Cloud.PublishString(msg.topic, msg.payload, msg.qos)
  150. }
  151. case ToEdge:
  152. if o.Edge != nil {
  153. err = o.Edge.PublishString(msg.topic, msg.payload, msg.qos)
  154. }
  155. }
  156. return err
  157. }
  158. // MQTTConnectMgr 连接保持协程(每 10 秒重连)
  159. func (o *MQTTMgr) MQTTConnectMgr(args ...interface{}) interface{} {
  160. for {
  161. time.Sleep(10 * time.Second)
  162. if o.Cloud != nil {
  163. o.Cloud.Connect()
  164. }
  165. if o.Edge != nil {
  166. o.Edge.Connect()
  167. }
  168. }
  169. }
  170. // MQTTMessageHandle 队列消费协程(带 RETRY)
  171. func (o *MQTTMgr) MQTTMessageHandle(args ...interface{}) interface{} {
  172. defer func() {
  173. if err := recover(); err != nil {
  174. logrus.Errorf("MQTTMgr.MQTTMessageHandle发生异常:%v", err)
  175. logrus.Errorf("MQTTMgr.MQTTMessageHandle发生异常,堆栈信息:%s", string(debug.Stack()))
  176. time.Sleep(time.Second)
  177. if o.restartFn != nil {
  178. o.restartFn(o.MQTTMessageHandle, args)
  179. }
  180. }
  181. }()
  182. for {
  183. if m, ok, _ := o.Queue.Get(); ok {
  184. if msg, ok := m.(*MQTTMessage); ok {
  185. for {
  186. if err := o.doPublish(msg); err != nil {
  187. logrus.Errorf("发布主题为%s的消息失败,原因:%s", msg.topic, err.Error())
  188. time.Sleep(time.Second)
  189. continue
  190. }
  191. break
  192. }
  193. }
  194. } else {
  195. time.Sleep(200 * time.Millisecond)
  196. }
  197. }
  198. }