mqtt.go 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390
  1. package item
  2. import (
  3. "encoding/json"
  4. "errors"
  5. "fmt"
  6. "regexp"
  7. "runtime"
  8. "runtime/debug"
  9. "server/global"
  10. "server/model"
  11. "server/utils/mqtt"
  12. "server/utils/protocol"
  13. "strings"
  14. "sync"
  15. "time"
  16. )
  17. func InitMqtt() {
  18. MqttService = GetHandler()
  19. MqttService.SubscribeTopics()
  20. go MqttService.Handler()
  21. }
  22. var MqttService *MqttHandler
  23. var timeoutReg = regexp.MustCompile("Client .* has exceeded timeout")
  24. var connectReg = regexp.MustCompile(`New client connected from .* as .*\(`)
  25. var disconnectReg = regexp.MustCompile("Client mqttx_893e4b7d disconnected")
  26. type MqttHandler struct {
  27. queue *mqtt.MlQueue
  28. }
  29. var _handlerOnce sync.Once
  30. var _handlerSingle *MqttHandler
  31. func GetHandler() *MqttHandler {
  32. _handlerOnce.Do(func() {
  33. _handlerSingle = &MqttHandler{
  34. queue: mqtt.NewQueue(10000),
  35. }
  36. })
  37. return _handlerSingle
  38. }
  39. func (o *MqttHandler) SubscribeTopics() {
  40. mqtt.GetMQTTMgr().Subscribe("/sys/#", mqtt.AtLeastOnce, o.HandlerData)
  41. }
  42. func (o *MqttHandler) HandlerData(m mqtt.Message) {
  43. for {
  44. ok, cnt := o.queue.Put(&m)
  45. if ok {
  46. break
  47. } else {
  48. global.GVA_LOG.Error(fmt.Sprintf("HandlerData:查询队列失败,队列消息数量:%d", cnt))
  49. runtime.Gosched()
  50. }
  51. }
  52. }
  53. func (o *MqttHandler) Handler() interface{} {
  54. defer func() {
  55. if err := recover(); err != nil {
  56. go GetHandler().Handler()
  57. global.GVA_LOG.Error(fmt.Sprintf("MqttHandler.Handler:发生异常:%s", string(debug.Stack())))
  58. }
  59. }()
  60. for {
  61. msg, ok, quantity := o.queue.Get()
  62. if !ok {
  63. time.Sleep(10 * time.Millisecond)
  64. continue
  65. }
  66. if quantity > 1000 {
  67. global.GVA_LOG.Error(fmt.Sprintf("队列堆积: %d", quantity))
  68. }
  69. m, ok := msg.(*mqtt.Message)
  70. if !ok {
  71. continue
  72. }
  73. fmt.Println(m.Topic())
  74. // 1. 解析 Topic
  75. _, _, eventFmt, err := parseTopic(m.Topic())
  76. if err != nil {
  77. global.GVA_LOG.Error("解析Topic失败:" + err.Error())
  78. continue
  79. }
  80. // 2. 解析外层报文
  81. var root model.MsgRoot
  82. if err := json.Unmarshal(m.Payload(), &root); err != nil {
  83. global.GVA_LOG.Error("解析MQTT报文失败:" + err.Error())
  84. continue
  85. }
  86. // 3. 根据事件类型分流解析 value
  87. switch eventFmt {
  88. case protocol.HeartbeatFmt:
  89. var val model.HeartbeatValue
  90. if err := json.Unmarshal(root.Params.Value, &val); err != nil {
  91. global.GVA_LOG.Error("解析心跳报文失败")
  92. continue
  93. }
  94. go HandleHeartBeat(val.DeviceName, val.UUID)
  95. case protocol.TriggerFmt:
  96. var val model.TriggerValue
  97. _ = json.Unmarshal(root.Params.Value, &val)
  98. // 传感器触发业务逻辑
  99. case protocol.AlsFmt:
  100. var val model.AlsValue
  101. _ = json.Unmarshal(root.Params.Value, &val)
  102. // 照度入库
  103. case protocol.ConsumptionFmt:
  104. var val model.ConsumptionValue
  105. _ = json.Unmarshal(root.Params.Value, &val)
  106. // 能耗统计
  107. CreateDeviceConsumption(val)
  108. case protocol.CurrentFmt:
  109. var val model.CurrentValue
  110. _ = json.Unmarshal(root.Params.Value, &val)
  111. // 更新亮度、色温
  112. case protocol.TemperatureHumidityFmt:
  113. var val model.TempHumValue
  114. _ = json.Unmarshal(root.Params.Value, &val)
  115. // 温湿度入库
  116. case protocol.BeaconFmt:
  117. var val model.BeaconValue
  118. _ = json.Unmarshal(root.Params.Value, &val)
  119. go HandleTempScanDeviceData(val)
  120. go HandleHeartBeat(val.DeviceName, val.UUID)
  121. case protocol.SettingFmt:
  122. var val model.SceneValue
  123. _ = json.Unmarshal(root.Params.Value, &val)
  124. go HandleSceneData(val)
  125. default:
  126. global.GVA_LOG.Info("未处理事件:" + eventFmt)
  127. }
  128. }
  129. }
  130. // Publish 发布消息
  131. func (o *MqttHandler) Publish(topic string, data interface{}) error {
  132. return mqtt.GetMQTTMgr().Publish(topic, data, mqtt.AtLeastOnce)
  133. }
  134. // GetTopic 自定义主题
  135. func (o *MqttHandler) GetTopic(deviceSn, protocol string) string {
  136. return fmt.Sprintf("mini/%s/%s", deviceSn, protocol)
  137. }
  138. // parseTopic 解析 /sys/# 主题
  139. func parseTopic(topic string) (string, string, string, error) {
  140. strList := strings.Split(topic, "/")
  141. if len(strList) < 7 {
  142. return "", "", "", errors.New("topic 格式不正确")
  143. }
  144. if strList[1] != "sys" {
  145. return "", "", "", errors.New("不是 sys 主题")
  146. }
  147. productKey := strList[2]
  148. deviceName := strList[3]
  149. eventFmt := "/" + strings.Join(strList[4:], "/")
  150. return productKey, deviceName, eventFmt, nil
  151. }
  152. // ==================== 下发指令结构体 ====================
  153. type ControlCmd struct {
  154. Code int `json:"code"`
  155. DeviceName string `json:"deviceName"`
  156. Area string `json:"area"`
  157. Address string `json:"address"`
  158. Action string `json:"action"`
  159. Params string `json:"params"`
  160. Identity string `json:"identity"`
  161. }
  162. // ==================== 底层发送方法 ====================
  163. func sendCmd(code int, productKey, gatewayDN, netPwd, area, addr, action, params string) error {
  164. topic := fmt.Sprintf("/%s/%s/user/get", productKey, gatewayDN)
  165. cmd := ControlCmd{
  166. Code: code,
  167. DeviceName: gatewayDN,
  168. Area: area,
  169. Address: addr,
  170. Action: action,
  171. Params: params,
  172. Identity: netPwd,
  173. }
  174. payload, err := json.Marshal(cmd)
  175. if err != nil {
  176. return fmt.Errorf("报文序列化失败: %w", err)
  177. }
  178. return MqttService.Publish(topic, payload)
  179. }
  180. func GetSetting(productKey, gatewayDN, netPwd, area, addr string) error {
  181. return sendCmd(400, productKey, gatewayDN, netPwd, area, addr, "getSetting", "")
  182. }
  183. func DeviceSendCmd(productKey, gatewayDN, netPwd, area, addr, action, params string) error {
  184. return sendCmd(400, productKey, gatewayDN, netPwd, area, addr, action, params)
  185. }
  186. // -------------------------------------------------------------------------------------
  187. // ==================== GS200 全量 80+ 指令封装(单灯/群组/全区/空调/参数/上报/情景) ====================
  188. // -------------------------------------------------------------------------------------
  189. // ==================== 1. 单灯控制 ====================
  190. func BlinkLight(productKey, gatewayDN, netPwd, area, number string) error {
  191. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "blink", "")
  192. }
  193. func LightOn(productKey, gatewayDN, netPwd, area, number string) error {
  194. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "lightOn", "")
  195. }
  196. func LightOff(productKey, gatewayDN, netPwd, area, number string) error {
  197. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "lightOff", "")
  198. }
  199. func LightSleep(productKey, gatewayDN, netPwd, area, number string) error {
  200. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "lightSleep", "")
  201. }
  202. func StopBlink(productKey, gatewayDN, netPwd, area, number string) error {
  203. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "stopBlink", "")
  204. }
  205. func SsrOn(productKey, gatewayDN, netPwd, area, number string) error {
  206. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "ssrOn", "")
  207. }
  208. func SsrOff(productKey, gatewayDN, netPwd, area, number string) error {
  209. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "ssrOff", "")
  210. }
  211. func NetOn(productKey, gatewayDN, netPwd, area, number string) error {
  212. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "netOn", "")
  213. }
  214. func NetOff(productKey, gatewayDN, netPwd, area, number string) error {
  215. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "netOff", "")
  216. }
  217. func RelayOn(productKey, gatewayDN, netPwd, area, number string) error {
  218. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "relayOn", "")
  219. }
  220. func RelayOff(productKey, gatewayDN, netPwd, area, number string) error {
  221. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "relayOff", "")
  222. }
  223. // ==================== 2. 群组控制 ====================
  224. func BlinkGroup(productKey, gatewayDN, netPwd, area, cluster string) error {
  225. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "blink", "")
  226. }
  227. func GroupLightOn(productKey, gatewayDN, netPwd, area, cluster string) error {
  228. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "lightOn", "")
  229. }
  230. func GroupLightOff(productKey, gatewayDN, netPwd, area, cluster string) error {
  231. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "lightOff", "")
  232. }
  233. func GroupLightSleep(productKey, gatewayDN, netPwd, area, cluster string) error {
  234. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "lightSleep", "")
  235. }
  236. // ==================== 3. 亮度/色温/延时/模式参数 ====================
  237. func SetHighBright(productKey, gatewayDN, netPwd, area, number, bright string) error {
  238. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setHighBright", bright)
  239. }
  240. func SetStandbyBright(productKey, gatewayDN, netPwd, area, number, bright string) error {
  241. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setStandbyBright", bright)
  242. }
  243. func SetCctBright(productKey, gatewayDN, netPwd, area, number, cct string) error {
  244. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setCctBright", cct)
  245. }
  246. func SetDelayTime(productKey, gatewayDN, netPwd, area, number, sec string) error {
  247. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setDelayTime", sec)
  248. }
  249. func SetDelayTime2(productKey, gatewayDN, netPwd, area, number, sec string) error {
  250. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setDelayTime2", sec)
  251. }
  252. func SetLightMode(productKey, gatewayDN, netPwd, area, number, mode string) error {
  253. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setLightMode", mode)
  254. }
  255. func SetDelayMode(productKey, gatewayDN, netPwd, area, number, mode string) error {
  256. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setDelayMode", mode)
  257. }
  258. func SetAlsMode(productKey, gatewayDN, netPwd, area, number, mode string) error {
  259. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setAlsMode", mode)
  260. }
  261. // ==================== 4. 情景模式 ====================
  262. func CallScene(productKey, gatewayDN, netPwd, area, cluster, sceneNo string) error {
  263. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "callScene", sceneNo)
  264. }
  265. func SaveToScene(productKey, gatewayDN, netPwd, area, cluster, sceneNo string) error {
  266. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "savetoScene", sceneNo)
  267. }
  268. func ReadScene(productKey, gatewayDN, netPwd, area, number, sceneNo string) error {
  269. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readScene", sceneNo)
  270. }
  271. // ==================== 5. 全区/网关指令 ====================
  272. func ScanAll(productKey, gatewayDN, netPwd, area string) error {
  273. return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "scan", "")
  274. }
  275. func StopScanAll(productKey, gatewayDN, netPwd, area string) error {
  276. return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "stopScan", "")
  277. }
  278. func GatewayReboot(productKey, gatewayDN, netPwd, area, waitSec string) error {
  279. return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "reboot", waitSec)
  280. }
  281. func SetRebootTime(productKey, gatewayDN, netPwd, area, timeRange string) error {
  282. return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "rebootSetting", timeRange)
  283. }
  284. func GatewayUpgrade(productKey, gatewayDN, netPwd, area string) error {
  285. return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "upgrade", "")
  286. }
  287. // ==================== 6. 上报控制 ====================
  288. func ReportConsumption(productKey, gatewayDN, netPwd, area, interval string) error {
  289. return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportConsumption", interval)
  290. }
  291. func StopReportConsumption(productKey, gatewayDN, netPwd, area string) error {
  292. return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportConsumptionStop", "")
  293. }
  294. func ReportSetting(productKey, gatewayDN, netPwd, area, interval string) error {
  295. return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportSetting", interval)
  296. }
  297. func StopReportSetting(productKey, gatewayDN, netPwd, area string) error {
  298. return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportSettingStop", "")
  299. }
  300. func ReportConsumptionAck(productKey, gatewayDN, netPwd, area string) error {
  301. return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportConsumptionAck", "")
  302. }
  303. func ReportSettingAck(productKey, gatewayDN, netPwd, area string) error {
  304. return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportSettingAck", "")
  305. }
  306. // ==================== 8. 地址修改 ====================
  307. func SetNumberAddress(productKey, gatewayDN, netPwd, area, number, param string) error {
  308. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setNumberAddress", param)
  309. }
  310. func SetClusterAddress(productKey, gatewayDN, netPwd, area, cluster, newCluster string) error {
  311. return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "setClusterAddress", newCluster)
  312. }
  313. func SetAreaAddress(productKey, gatewayDN, netPwd, area, number, newArea string) error {
  314. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setAreaAddress", newArea)
  315. }
  316. // ==================== 9. 传感器/报警/复位 ====================
  317. func SsrControl(productKey, gatewayDN, netPwd, area, number string) error {
  318. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "ssrControl", "")
  319. }
  320. func AlarmOn(productKey, gatewayDN, netPwd, area, number string) error {
  321. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "alarmOn", "")
  322. }
  323. func AlarmOff(productKey, gatewayDN, netPwd, area, number string) error {
  324. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "alarmOff", "")
  325. }
  326. func ResetDevice(productKey, gatewayDN, netPwd, area, number string) error {
  327. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "reset", "")
  328. }
  329. // ==================== 10. 电量/电压/电流/功率查询 ====================
  330. func ReadMeterData(productKey, gatewayDN, netPwd, area, number string) error {
  331. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readMeterData", "")
  332. }
  333. func ReadVoltage(productKey, gatewayDN, netPwd, area, number string) error {
  334. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readVoltage", "")
  335. }
  336. func ReadCurrent(productKey, gatewayDN, netPwd, area, number string) error {
  337. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readCurrent", "")
  338. }
  339. func ReadPower(productKey, gatewayDN, netPwd, area, number string) error {
  340. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readPower", "")
  341. }
  342. // ==================== 11. 版本/信号查询 ====================
  343. func ReadVersion(productKey, gatewayDN, netPwd, area, number string) error {
  344. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readVersion", "")
  345. }
  346. func ReadRssi(productKey, gatewayDN, netPwd, area, number string) error {
  347. return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readRssi", "")
  348. }