mqttmgr.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187
  1. package main
  2. import (
  3. "runtime/debug"
  4. "sync"
  5. "time"
  6. "github.com/sirupsen/logrus"
  7. "lc/common/mqtt"
  8. "lc/common/util"
  9. )
  10. type OptType uint8
  11. const (
  12. ToAll OptType = 0 //发布和订阅边缘端与云端的消息
  13. ToCloud OptType = 1 //发布和订阅云端的消息
  14. ToEdge OptType = 2 //发布和订阅边缘端的消息
  15. )
  16. var _mqttMgronce sync.Once
  17. var _mqttMgrsingle *MQTTMgr
  18. // GetMQTTMgr 单态
  19. func GetMQTTMgr() *MQTTMgr {
  20. _mqttMgronce.Do(func() {
  21. _mqttMgrsingle = _newMQTTMgr()
  22. })
  23. return _mqttMgrsingle
  24. }
  25. type MQTTMgr struct {
  26. Cloud *MqttClient
  27. Edge *MqttClient
  28. Queue *util.MlQueue
  29. }
  30. // 建两个client
  31. func _newMQTTMgr() *MQTTMgr {
  32. mgr := &MQTTMgr{
  33. Queue: util.NewQueue(2000),
  34. }
  35. if appConfig.Edge.Mqtt.Server != "" {
  36. mgr.Edge = NewMqttClient(appConfig.Edge.Mqtt.Server,
  37. appConfig.GID+"@"+appname+appversion,
  38. appConfig.Edge.Mqtt.User,
  39. appConfig.Edge.Mqtt.Password,
  40. appConfig.Edge.Mqtt.Timeout,
  41. &EmptyMqttOnline{})
  42. }
  43. if appConfig.Cloud.Mqtt.Server != "" {
  44. mgr.Cloud = NewMqttClient(appConfig.Cloud.Mqtt.Server,
  45. appConfig.GID+"@"+appname+appversion,
  46. appConfig.Cloud.Mqtt.User,
  47. appConfig.Cloud.Mqtt.Password,
  48. appConfig.Cloud.Mqtt.Timeout,
  49. &EmptyMqttOnline{})
  50. }
  51. return mgr
  52. }
  53. // Subscribe 定阅
  54. func (o *MQTTMgr) Subscribe(topic string, qos mqtt.QOS, handler mqtt.MessageHandler, tp OptType) {
  55. switch tp {
  56. case ToAll:
  57. if o.Cloud != nil {
  58. o.Cloud.Handle(topic, handler)
  59. o.Cloud.Subscribe(topic, qos)
  60. }
  61. if o.Edge != nil {
  62. o.Edge.Handle(topic, handler)
  63. o.Edge.Subscribe(topic, qos)
  64. }
  65. case ToCloud:
  66. if o.Cloud != nil {
  67. o.Cloud.Handle(topic, handler)
  68. o.Cloud.Subscribe(topic, qos)
  69. }
  70. case ToEdge:
  71. if o.Edge != nil {
  72. o.Edge.Handle(topic, handler)
  73. o.Edge.Subscribe(topic, qos)
  74. }
  75. }
  76. }
  77. // UnSubscribe 退定
  78. func (o *MQTTMgr) UnSubscribe(topic string, tp OptType) {
  79. switch tp {
  80. case ToAll:
  81. if o.Cloud != nil {
  82. o.Cloud.Unsubscribe(topic)
  83. }
  84. if o.Edge != nil {
  85. o.Edge.Unsubscribe(topic)
  86. }
  87. case ToCloud:
  88. if o.Cloud != nil {
  89. o.Cloud.Unsubscribe(topic)
  90. }
  91. case ToEdge:
  92. if o.Edge != nil {
  93. o.Edge.Unsubscribe(topic)
  94. }
  95. }
  96. }
  97. // Publish 发布进队列
  98. func (o *MQTTMgr) Publish(topic string, payload []byte, qos mqtt.QOS, tp OptType) {
  99. msg := MQTTMessage{
  100. topic: topic,
  101. payload: payload,
  102. qos: qos,
  103. tp: tp,
  104. }
  105. o.Queue.Put(&msg)
  106. }
  107. // 发布低
  108. func (o *MQTTMgr) _publish(msg *MQTTMessage) error {
  109. var err error
  110. switch msg.tp {
  111. case ToAll:
  112. if o.Cloud != nil {
  113. err = o.Cloud.Publish(msg.topic, msg.payload, msg.qos)
  114. }
  115. if o.Edge != nil {
  116. o.Edge.Publish(msg.topic, msg.payload, msg.qos)
  117. }
  118. case ToCloud:
  119. if o.Cloud != nil {
  120. err = o.Cloud.Publish(msg.topic, msg.payload, msg.qos)
  121. }
  122. case ToEdge:
  123. if o.Edge != nil {
  124. o.Edge.Publish(msg.topic, msg.payload, msg.qos)
  125. }
  126. }
  127. return err
  128. }
  129. // MQTTConnectMgr 连接保持
  130. func (o *MQTTMgr) MQTTConnectMgr(args ...interface{}) interface{} {
  131. for {
  132. time.Sleep(10 * time.Second)
  133. if o.Cloud != nil {
  134. o.Cloud.Connect()
  135. }
  136. if o.Edge != nil {
  137. o.Edge.Connect()
  138. }
  139. }
  140. }
  141. func (o *MQTTMgr) MQTTMessageHandle(args ...interface{}) interface{} {
  142. defer func() {
  143. if err := recover(); err != nil {
  144. logrus.Errorf("MQTTMgr.MQTTMessageHandle发生异常:%v", err)
  145. logrus.Errorf("MQTTMgr.MQTTMessageHandle发生异常,堆栈信息:%s", string(debug.Stack()))
  146. gopool.Add(o.MQTTMessageHandle, args)
  147. }
  148. }()
  149. var err error
  150. for { //队列中所有发布
  151. if m, ok, _ := o.Queue.Get(); ok {
  152. if msg, ok := m.(*MQTTMessage); ok {
  153. RETRY:
  154. err = o._publish(msg)
  155. if err != nil {
  156. logrus.Errorf("发布主题为%s的消息失败,原因:%s", msg.topic, err.Error())
  157. time.Sleep(time.Second)
  158. goto RETRY
  159. }
  160. }
  161. } else {
  162. time.Sleep(200 * time.Millisecond)
  163. }
  164. }
  165. }
  166. type MQTTMessage struct {
  167. topic string
  168. payload []byte
  169. qos mqtt.QOS
  170. tp OptType
  171. }