mqtthandle.go 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165
  1. package main
  2. import (
  3. "os"
  4. "strings"
  5. "lc/common/mqtt"
  6. "lc/common/protocol"
  7. "lc/common/util"
  8. )
  9. func HandleTpQApp(m mqtt.Message) {
  10. var obj protocol.Pack_IDObject
  11. var ret protocol.Pack_MutilFileObject
  12. if err := obj.DeCode(m.PayloadString()); err == nil {
  13. //读文件内容
  14. ReadMutilFileContent(protocol.TP_GW_APP, obj.Data.Id, &ret)
  15. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  16. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP_ACK), str, 0, ToCloud)
  17. }
  18. }
  19. }
  20. func HandleTpWApp(m mqtt.Message) {
  21. var obj protocol.Pack_SeqFileObject
  22. var ret protocol.Pack_Ack
  23. err := obj.DeCode(m.PayloadString())
  24. if err == nil {
  25. err = HandleFile(protocol.TP_GW_SET_APP, &obj)
  26. }
  27. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  28. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_APP_ACK), str, 0, ToCloud)
  29. }
  30. }
  31. func HandleTpQSerial(m mqtt.Message) {
  32. var obj protocol.Pack_IDObject
  33. var ret protocol.Pack_MutilFileObject
  34. if err := obj.DeCode(m.PayloadString()); err == nil {
  35. //读文件内容
  36. ReadMutilFileContent(protocol.TP_GW_SERIAL, obj.Data.Id, &ret)
  37. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  38. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SERIAL_ACK), str, 0, ToCloud)
  39. }
  40. }
  41. }
  42. func HandleTpWSerial(m mqtt.Message) {
  43. var obj protocol.Pack_SeqFileObject
  44. var ret protocol.Pack_Ack
  45. err := obj.DeCode(m.PayloadString())
  46. if err == nil {
  47. err = HandleFile(protocol.TP_GW_SET_SERIAL, &obj)
  48. }
  49. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  50. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_SERIAL_ACK), str, 0, ToCloud)
  51. }
  52. }
  53. func HandleTpQRtu(m mqtt.Message) {
  54. var obj protocol.Pack_IDObject
  55. var ret protocol.Pack_MutilFileObject
  56. if err := obj.DeCode(m.PayloadString()); err == nil {
  57. //读文件内容
  58. ReadMutilFileContent(protocol.TP_GW_RTU, obj.Data.Id, &ret)
  59. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  60. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_RTU_ACK), str, 0, ToCloud)
  61. }
  62. }
  63. }
  64. func HandleTpWRtu(m mqtt.Message) {
  65. var obj protocol.Pack_SeqFileObject
  66. var ret protocol.Pack_Ack
  67. err := obj.DeCode(m.PayloadString())
  68. if err == nil {
  69. err = HandleFile(protocol.TP_GW_SET_RTU, &obj)
  70. }
  71. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  72. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_RTU_ACK), str, 0, ToCloud)
  73. }
  74. }
  75. func HandleTpQModel(m mqtt.Message) {
  76. var obj protocol.Pack_IDObject
  77. var ret protocol.Pack_MutilFileObject
  78. if err := obj.DeCode(m.PayloadString()); err == nil {
  79. //读文件内容
  80. ReadMutilFileContent(protocol.TP_GW_MODEL, obj.Data.Id, &ret)
  81. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  82. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_MODEL_ACK), str, 0, ToCloud)
  83. }
  84. }
  85. }
  86. func HandleTpWModel(m mqtt.Message) {
  87. var obj protocol.Pack_SeqFileObject
  88. var ret protocol.Pack_Ack
  89. err := obj.DeCode(m.PayloadString())
  90. if err == nil {
  91. err = HandleFile(protocol.TP_GW_SET_MODEL, &obj)
  92. }
  93. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  94. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL_ACK), str, 0, ToCloud)
  95. }
  96. }
  97. func HandleTpQLog(m mqtt.Message) {
  98. var obj protocol.Pack_IDObject
  99. var ret protocol.Pack_MutilFileObject
  100. if err := obj.DeCode(m.PayloadString()); err == nil {
  101. //读文件内容
  102. ReadMutilFileContent(protocol.TP_GW_LOG, obj.Data.Id, &ret)
  103. if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
  104. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_ACK), str, 0, ToCloud)
  105. }
  106. }
  107. }
  108. func HandleTpRLog(m mqtt.Message) {
  109. var obj protocol.Pack_IDObject
  110. var ret protocol.Pack_Ack
  111. var err error
  112. if err = obj.DeCode(m.PayloadString()); err == nil {
  113. rd, _ := os.ReadDir(util.GetPath(3))
  114. for _, fi := range rd {
  115. if ok := strings.HasSuffix(fi.Name(), ".log"); ok {
  116. err = os.Remove(util.GetPath(3) + fi.Name())
  117. }
  118. }
  119. if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil {
  120. GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG_ACK), str, 0, ToCloud)
  121. }
  122. }
  123. }
  124. func HandleTpQSys(m mqtt.Message) {
  125. var obj protocol.Pack_IDObject
  126. if err := obj.DeCode(m.PayloadString()); err == nil {
  127. go SysInfoStat(obj.Seq)
  128. }
  129. }
  130. type MqttOnline struct {
  131. }
  132. func (o *MqttOnline) GetOnlineMsg() (string, string) {
  133. //发布上线消息
  134. var obj protocol.Pack_IDObject
  135. str, err := obj.EnCode(appConfig.GID, GetNextUint64(), 0)
  136. if err != nil {
  137. return "", ""
  138. }
  139. return GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_ONLINE), str
  140. }
  141. func (o *MqttOnline) GetWillMsg() (string, string) {
  142. payload, _ := (&protocol.Pack_IDObject{}).EnCode(appConfig.GID, GetNextUint64(), 0) //遗嘱消息
  143. return GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_WILL), payload
  144. }
  145. // InitCloudMqttSubscribeTopics 初始化网关级别的主题订阅及路由
  146. func InitCloudMqttSubscribeTopics() {
  147. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP), mqtt.AtMostOnce, HandleTpQApp, ToCloud)
  148. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_APP), mqtt.AtMostOnce, HandleTpWApp, ToCloud)
  149. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SERIAL), mqtt.AtMostOnce, HandleTpQSerial, ToCloud)
  150. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_SERIAL), mqtt.AtMostOnce, HandleTpWSerial, ToCloud)
  151. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_RTU), mqtt.AtMostOnce, HandleTpQRtu, ToCloud)
  152. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_RTU), mqtt.AtMostOnce, HandleTpWRtu, ToCloud)
  153. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_MODEL), mqtt.AtMostOnce, HandleTpQModel, ToCloud)
  154. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL), mqtt.AtMostOnce, HandleTpWModel, ToCloud)
  155. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG), mqtt.AtMostOnce, HandleTpQLog, ToCloud)
  156. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG), mqtt.AtMostOnce, HandleTpRLog, ToCloud)
  157. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SYS), mqtt.AtMostOnce, HandleTpQSys, ToCloud)
  158. }