package main import ( "fmt" "os" "strings" "time" "lc/common/mqtt" "lc/common/protocol" "lc/common/util" ) func HandleTpQApp(m mqtt.Message) { var obj protocol.Pack_IDObject var ret protocol.Pack_MutilFileObject if err := obj.DeCode(m.PayloadString()); err == nil { //读文件内容 ReadMutilFileContent(protocol.TP_GW_APP, obj.Data.Id, &ret) if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP_ACK), str, 0, ToCloud) } } } func HandleTpWApp(m mqtt.Message) { var obj protocol.Pack_SeqFileObject var ret protocol.Pack_Ack err := obj.DeCode(m.PayloadString()) if err == nil { err = HandleFile(protocol.TP_GW_SET_APP, &obj) } if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_APP_ACK), str, 0, ToCloud) } } func HandleTpQSerial(m mqtt.Message) { var obj protocol.Pack_IDObject var ret protocol.Pack_MutilFileObject if err := obj.DeCode(m.PayloadString()); err == nil { //读文件内容 ReadMutilFileContent(protocol.TP_GW_SERIAL, obj.Data.Id, &ret) if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SERIAL_ACK), str, 0, ToCloud) } } } func HandleTpWSerial(m mqtt.Message) { var obj protocol.Pack_SeqFileObject var ret protocol.Pack_Ack err := obj.DeCode(m.PayloadString()) if err == nil { err = HandleFile(protocol.TP_GW_SET_SERIAL, &obj) } if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_SERIAL_ACK), str, 0, ToCloud) } } func HandleTpQRtu(m mqtt.Message) { var obj protocol.Pack_IDObject var ret protocol.Pack_MutilFileObject if err := obj.DeCode(m.PayloadString()); err == nil { //读文件内容 ReadMutilFileContent(protocol.TP_GW_RTU, obj.Data.Id, &ret) if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_RTU_ACK), str, 0, ToCloud) } } } func HandleTpWRtu(m mqtt.Message) { var obj protocol.Pack_SeqFileObject var ret protocol.Pack_Ack err := obj.DeCode(m.PayloadString()) if err == nil { err = HandleFile(protocol.TP_GW_SET_RTU, &obj) } if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_RTU_ACK), str, 0, ToCloud) } } func HandleTpQModel(m mqtt.Message) { var obj protocol.Pack_IDObject var ret protocol.Pack_MutilFileObject if err := obj.DeCode(m.PayloadString()); err == nil { //读文件内容 ReadMutilFileContent(protocol.TP_GW_MODEL, obj.Data.Id, &ret) if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_MODEL_ACK), str, 0, ToCloud) } } } func HandleTpWModel(m mqtt.Message) { var obj protocol.Pack_SeqFileObject var ret protocol.Pack_Ack err := obj.DeCode(m.PayloadString()) if err == nil { err = HandleFile(protocol.TP_GW_SET_MODEL, &obj) } if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL_ACK), str, 0, ToCloud) } } func HandleTpQLog(m mqtt.Message) { util.GetTagLog().Infof("sys", "HandleTpQLog:收到日志查询请求,topic=%s,payload=%s", m.Topic(), m.PayloadString()) var obj protocol.Pack_IDObject var ret protocol.Pack_MutilFileObject if err := obj.DeCode(m.PayloadString()); err != nil { util.GetTagLog().Errorf("sys", "HandleTpQLog:DeCode失败,err=%v", err) return } //读文件内容 ReadMutilFileContent(protocol.TP_GW_LOG, obj.Data.Id, &ret) if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_ACK), str, 0, ToCloud) util.GetTagLog().Infof("sys", "HandleTpQLog:日志ACK已发布") } else { util.GetTagLog().Errorf("sys", "HandleTpQLog:EnCode失败,err=%v", err) } } func HandleTpRLog(m mqtt.Message) { var obj protocol.Pack_IDObject var ret protocol.Pack_Ack var err error if err = obj.DeCode(m.PayloadString()); err == nil { rd, _ := os.ReadDir(util.GetPath(3)) for _, fi := range rd { if ok := strings.HasSuffix(fi.Name(), ".log"); ok { err = os.Remove(util.GetPath(3) + fi.Name()) } } if str, err := ret.EnCode(appConfig.GID, appConfig.GID, obj.Seq, err); err == nil { GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG_ACK), str, 0, ToCloud) } } } func HandleTpQSys(m mqtt.Message) { var obj protocol.Pack_IDObject if err := obj.DeCode(m.PayloadString()); err == nil { go SysInfoStat(obj.Seq) } } type MqttOnline struct { } func (o *MqttOnline) GetOnlineMsg() (string, string) { // 手拼 JSON,避开 json-iterator ConfigFastest 的潜在 marshal 问题 seq := GetNextUint64() now := protocol.BJNow().Format("2006-01-02 15:04:05") payload := fmt.Sprintf(`{"id":"%s","seq":%d,"gid":"%s","time":"%s","data":{"id":0}}`, appConfig.GID, seq, appConfig.GID, now) topic := GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_ONLINE) util.GetTagLog().Infof("sys", "GetOnlineMsg: topic=%s", topic) return topic, payload } func (o *MqttOnline) GetWillMsg() (string, string) { seq := GetNextUint64() now := protocol.BJNow().Format("2006-01-02 15:04:05") payload := fmt.Sprintf(`{"id":"%s","seq":%d,"gid":"%s","time":"%s","data":{"id":0}}`, appConfig.GID, seq, appConfig.GID, now) return GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_WILL), payload } // Heartbeat 定期重发 online 消息,防止 QoS=0 丢包导致云端永久离线 func Heartbeat(args ...interface{}) interface{} { for { time.Sleep(60 * time.Second) mgr := GetMQTTMgr() if mgr.Cloud != nil && mgr.Cloud.IsConnected() { topic, str := (&MqttOnline{}).GetOnlineMsg() if topic != "" { if err := mgr.Cloud.PublishString(topic, str, 0); err != nil { util.GetTagLog().Errorf("sys", "Heartbeat:发布online失败,topic=%s,err=%v", topic, err) } } } } } // HandleTpWDeploy 转换签名供 Subscribe 使用 func HandleTpWDeploy(m mqtt.Message) { HandleTpDeploy(m) } // InitCloudMqttSubscribeTopics 初始化网关级别的主题订阅及路由 func InitCloudMqttSubscribeTopics() { GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP), mqtt.AtMostOnce, HandleTpQApp, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_APP), mqtt.AtMostOnce, HandleTpWApp, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SERIAL), mqtt.AtMostOnce, HandleTpQSerial, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_SERIAL), mqtt.AtMostOnce, HandleTpWSerial, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_RTU), mqtt.AtMostOnce, HandleTpQRtu, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_RTU), mqtt.AtMostOnce, HandleTpWRtu, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_MODEL), mqtt.AtMostOnce, HandleTpQModel, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL), mqtt.AtMostOnce, HandleTpWModel, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG), mqtt.AtMostOnce, HandleTpQLog, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG), mqtt.AtMostOnce, HandleTpRLog, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_CFG), mqtt.AtMostOnce, HandleTpLogCfg, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SYS), mqtt.AtMostOnce, HandleTpQSys, ToCloud) GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_DEPLOY), mqtt.AtMostOnce, HandleTpWDeploy, ToCloud) }