mqtthandle.go 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200
  1. package main
  2. import (
  3. "fmt"
  4. "os"
  5. "strings"
  6. "time"
  7. "lc/common/mqtt"
  8. "lc/common/protocol"
  9. "lc/common/util"
  10. )
  11. func HandleTpQApp(m mqtt.Message) {
  12. var obj protocol.Pack_IDObject
  13. var ret protocol.Pack_MutilFileObject
  14. if err := obj.DeCode(m.PayloadString()); err == nil {
  15. //读文件内容
  16. ReadMutilFileContent(protocol.TP_GW_APP, obj.Data.Id, &ret)
  17. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  18. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP_ACK), str, 0, ToCloud)
  19. }
  20. }
  21. }
  22. func HandleTpWApp(m mqtt.Message) {
  23. var obj protocol.Pack_SeqFileObject
  24. var ret protocol.Pack_Ack
  25. err := obj.DeCode(m.PayloadString())
  26. if err == nil {
  27. err = HandleFile(protocol.TP_GW_SET_APP, &obj)
  28. }
  29. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  30. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_APP_ACK), str, 0, ToCloud)
  31. }
  32. }
  33. func HandleTpQSerial(m mqtt.Message) {
  34. var obj protocol.Pack_IDObject
  35. var ret protocol.Pack_MutilFileObject
  36. if err := obj.DeCode(m.PayloadString()); err == nil {
  37. //读文件内容
  38. ReadMutilFileContent(protocol.TP_GW_SERIAL, obj.Data.Id, &ret)
  39. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  40. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SERIAL_ACK), str, 0, ToCloud)
  41. }
  42. }
  43. }
  44. func HandleTpWSerial(m mqtt.Message) {
  45. var obj protocol.Pack_SeqFileObject
  46. var ret protocol.Pack_Ack
  47. err := obj.DeCode(m.PayloadString())
  48. if err == nil {
  49. err = HandleFile(protocol.TP_GW_SET_SERIAL, &obj)
  50. }
  51. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  52. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_SERIAL_ACK), str, 0, ToCloud)
  53. }
  54. }
  55. func HandleTpQRtu(m mqtt.Message) {
  56. var obj protocol.Pack_IDObject
  57. var ret protocol.Pack_MutilFileObject
  58. if err := obj.DeCode(m.PayloadString()); err == nil {
  59. //读文件内容
  60. ReadMutilFileContent(protocol.TP_GW_RTU, obj.Data.Id, &ret)
  61. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  62. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_RTU_ACK), str, 0, ToCloud)
  63. }
  64. }
  65. }
  66. func HandleTpWRtu(m mqtt.Message) {
  67. var obj protocol.Pack_SeqFileObject
  68. var ret protocol.Pack_Ack
  69. err := obj.DeCode(m.PayloadString())
  70. if err == nil {
  71. err = HandleFile(protocol.TP_GW_SET_RTU, &obj)
  72. }
  73. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  74. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_RTU_ACK), str, 0, ToCloud)
  75. }
  76. }
  77. func HandleTpQModel(m mqtt.Message) {
  78. var obj protocol.Pack_IDObject
  79. var ret protocol.Pack_MutilFileObject
  80. if err := obj.DeCode(m.PayloadString()); err == nil {
  81. //读文件内容
  82. ReadMutilFileContent(protocol.TP_GW_MODEL, obj.Data.Id, &ret)
  83. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  84. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_MODEL_ACK), str, 0, ToCloud)
  85. }
  86. }
  87. }
  88. func HandleTpWModel(m mqtt.Message) {
  89. var obj protocol.Pack_SeqFileObject
  90. var ret protocol.Pack_Ack
  91. err := obj.DeCode(m.PayloadString())
  92. if err == nil {
  93. err = HandleFile(protocol.TP_GW_SET_MODEL, &obj)
  94. }
  95. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  96. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL_ACK), str, 0, ToCloud)
  97. }
  98. }
  99. func HandleTpQLog(m mqtt.Message) {
  100. util.GetTagLog().Infof("sys", "HandleTpQLog:收到日志查询请求,topic=%s,payload=%s", m.Topic(), m.PayloadString())
  101. var obj protocol.Pack_IDObject
  102. var ret protocol.Pack_MutilFileObject
  103. if err := obj.DeCode(m.PayloadString()); err != nil {
  104. util.GetTagLog().Errorf("sys", "HandleTpQLog:DeCode失败,err=%v", err)
  105. return
  106. }
  107. //读文件内容
  108. ReadMutilFileContent(protocol.TP_GW_LOG, obj.Data.Id, &ret)
  109. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  110. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_ACK), str, 0, ToCloud)
  111. util.GetTagLog().Infof("sys", "HandleTpQLog:日志ACK已发布")
  112. } else {
  113. util.GetTagLog().Errorf("sys", "HandleTpQLog:EnCode失败,err=%v", err)
  114. }
  115. }
  116. func HandleTpRLog(m mqtt.Message) {
  117. var obj protocol.Pack_IDObject
  118. var ret protocol.Pack_Ack
  119. var err error
  120. if err = obj.DeCode(m.PayloadString()); err == nil {
  121. rd, _ := os.ReadDir(util.GetPath(3))
  122. for _, fi := range rd {
  123. if ok := strings.HasSuffix(fi.Name(), ".log"); ok {
  124. err = os.Remove(util.GetPath(3) + fi.Name())
  125. }
  126. }
  127. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  128. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG_ACK), str, 0, ToCloud)
  129. }
  130. }
  131. }
  132. func HandleTpQSys(m mqtt.Message) {
  133. var obj protocol.Pack_IDObject
  134. if err := obj.DeCode(m.PayloadString()); err == nil {
  135. go SysInfoStat(obj.Seq)
  136. }
  137. }
  138. type MqttOnline struct {
  139. }
  140. func (o *MqttOnline) GetOnlineMsg() (string, string) {
  141. // 手拼 JSON,避开 json-iterator ConfigFastest 的潜在 marshal 问题
  142. seq := GetNextUint64()
  143. now := protocol.BJNow().Format("2006-01-02 15:04:05")
  144. payload := fmt.Sprintf(`{"id":"%s","seq":%d,"gid":"%s","time":"%s","data":{"id":0}}`,
  145. appConfig.GID, seq, appConfig.GID, now)
  146. topic := GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_ONLINE)
  147. util.GetTagLog().Infof("sys", "GetOnlineMsg: topic=%s", topic)
  148. return topic, payload
  149. }
  150. func (o *MqttOnline) GetWillMsg() (string, string) {
  151. seq := GetNextUint64()
  152. now := protocol.BJNow().Format("2006-01-02 15:04:05")
  153. payload := fmt.Sprintf(`{"id":"%s","seq":%d,"gid":"%s","time":"%s","data":{"id":0}}`,
  154. appConfig.GID, seq, appConfig.GID, now)
  155. return GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_WILL), payload
  156. }
  157. // Heartbeat 定期重发 online 消息,防止 QoS=0 丢包导致云端永久离线
  158. func Heartbeat(args ...interface{}) interface{} {
  159. for {
  160. time.Sleep(60 * time.Second)
  161. mgr := GetMQTTMgr()
  162. if mgr.Cloud != nil && mgr.Cloud.IsConnected() {
  163. topic, str := (&MqttOnline{}).GetOnlineMsg()
  164. if topic != "" {
  165. if err := mgr.Cloud.PublishString(topic, str, 0); err != nil {
  166. util.GetTagLog().Errorf("sys", "Heartbeat:发布online失败,topic=%s,err=%v", topic, err)
  167. }
  168. }
  169. }
  170. }
  171. }
  172. // HandleTpWDeploy 转换签名供 Subscribe 使用
  173. func HandleTpWDeploy(m mqtt.Message) {
  174. HandleTpDeploy(m)
  175. }
  176. // InitCloudMqttSubscribeTopics 初始化网关级别的主题订阅及路由
  177. func InitCloudMqttSubscribeTopics() {
  178. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP), mqtt.AtMostOnce, HandleTpQApp, ToCloud)
  179. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_APP), mqtt.AtMostOnce, HandleTpWApp, ToCloud)
  180. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SERIAL), mqtt.AtMostOnce, HandleTpQSerial, ToCloud)
  181. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_SERIAL), mqtt.AtMostOnce, HandleTpWSerial, ToCloud)
  182. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_RTU), mqtt.AtMostOnce, HandleTpQRtu, ToCloud)
  183. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_RTU), mqtt.AtMostOnce, HandleTpWRtu, ToCloud)
  184. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_MODEL), mqtt.AtMostOnce, HandleTpQModel, ToCloud)
  185. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL), mqtt.AtMostOnce, HandleTpWModel, ToCloud)
  186. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG), mqtt.AtMostOnce, HandleTpQLog, ToCloud)
  187. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG), mqtt.AtMostOnce, HandleTpRLog, ToCloud)
  188. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_CFG), mqtt.AtMostOnce, HandleTpLogCfg, ToCloud)
  189. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SYS), mqtt.AtMostOnce, HandleTpQSys, ToCloud)
  190. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_DEPLOY), mqtt.AtMostOnce, HandleTpWDeploy, ToCloud)
  191. }