package item import ( "encoding/json" "errors" "fmt" "regexp" "runtime" "runtime/debug" "server/global" "server/model" "server/utils/mqtt" "server/utils/protocol" "strings" "sync" "time" ) func InitMqtt() { MqttService = GetHandler() MqttService.SubscribeTopics() go MqttService.Handler() } var MqttService *MqttHandler var timeoutReg = regexp.MustCompile("Client .* has exceeded timeout") var connectReg = regexp.MustCompile(`New client connected from .* as .*\(`) var disconnectReg = regexp.MustCompile("Client mqttx_893e4b7d disconnected") type MqttHandler struct { queue *mqtt.MlQueue } var _handlerOnce sync.Once var _handlerSingle *MqttHandler func GetHandler() *MqttHandler { _handlerOnce.Do(func() { _handlerSingle = &MqttHandler{ queue: mqtt.NewQueue(10000), } }) return _handlerSingle } func (o *MqttHandler) SubscribeTopics() { mqtt.GetMQTTMgr().Subscribe("/sys/#", mqtt.AtLeastOnce, o.HandlerData) } func (o *MqttHandler) HandlerData(m mqtt.Message) { for { ok, cnt := o.queue.Put(&m) if ok { break } else { global.GVA_LOG.Error(fmt.Sprintf("HandlerData:查询队列失败,队列消息数量:%d", cnt)) runtime.Gosched() } } } func (o *MqttHandler) Handler() interface{} { defer func() { if err := recover(); err != nil { go GetHandler().Handler() global.GVA_LOG.Error(fmt.Sprintf("MqttHandler.Handler:发生异常:%s", string(debug.Stack()))) } }() for { msg, ok, quantity := o.queue.Get() if !ok { time.Sleep(10 * time.Millisecond) continue } if quantity > 1000 { global.GVA_LOG.Error(fmt.Sprintf("队列堆积: %d", quantity)) } m, ok := msg.(*mqtt.Message) if !ok { continue } fmt.Println(m.Topic()) // 1. 解析 Topic _, _, eventFmt, err := parseTopic(m.Topic()) if err != nil { global.GVA_LOG.Error("解析Topic失败:" + err.Error()) continue } // 2. 解析外层报文 var root model.MsgRoot if err := json.Unmarshal(m.Payload(), &root); err != nil { global.GVA_LOG.Error("解析MQTT报文失败:" + err.Error()) continue } // 3. 根据事件类型分流解析 value switch eventFmt { case protocol.HeartbeatFmt: var val model.HeartbeatValue if err := json.Unmarshal(root.Params.Value, &val); err != nil { global.GVA_LOG.Error("解析心跳报文失败") continue } go HandleHeartBeat(val.DeviceName, val.UUID) case protocol.TriggerFmt: var val model.TriggerValue _ = json.Unmarshal(root.Params.Value, &val) // 传感器触发业务逻辑 case protocol.AlsFmt: var val model.AlsValue _ = json.Unmarshal(root.Params.Value, &val) // 照度入库 case protocol.ConsumptionFmt: var val model.ConsumptionValue _ = json.Unmarshal(root.Params.Value, &val) // 能耗统计 CreateDeviceConsumption(val) case protocol.CurrentFmt: var val model.CurrentValue _ = json.Unmarshal(root.Params.Value, &val) // 更新亮度、色温 case protocol.TemperatureHumidityFmt: var val model.TempHumValue _ = json.Unmarshal(root.Params.Value, &val) // 温湿度入库 case protocol.BeaconFmt: var val model.BeaconValue _ = json.Unmarshal(root.Params.Value, &val) go HandleTempScanDeviceData(val) go HandleHeartBeat(val.DeviceName, val.UUID) case protocol.SettingFmt: var val model.SceneValue _ = json.Unmarshal(root.Params.Value, &val) go HandleSceneData(val) default: global.GVA_LOG.Info("未处理事件:" + eventFmt) } } } // Publish 发布消息 func (o *MqttHandler) Publish(topic string, data interface{}) error { return mqtt.GetMQTTMgr().Publish(topic, data, mqtt.AtLeastOnce) } // GetTopic 自定义主题 func (o *MqttHandler) GetTopic(deviceSn, protocol string) string { return fmt.Sprintf("mini/%s/%s", deviceSn, protocol) } // parseTopic 解析 /sys/# 主题 func parseTopic(topic string) (string, string, string, error) { strList := strings.Split(topic, "/") if len(strList) < 7 { return "", "", "", errors.New("topic 格式不正确") } if strList[1] != "sys" { return "", "", "", errors.New("不是 sys 主题") } productKey := strList[2] deviceName := strList[3] eventFmt := "/" + strings.Join(strList[4:], "/") return productKey, deviceName, eventFmt, nil } // ==================== 下发指令结构体 ==================== type ControlCmd struct { Code int `json:"code"` DeviceName string `json:"deviceName"` Area string `json:"area"` Address string `json:"address"` Action string `json:"action"` Params string `json:"params"` Identity string `json:"identity"` } // ==================== 底层发送方法 ==================== func sendCmd(code int, productKey, gatewayDN, netPwd, area, addr, action, params string) error { topic := fmt.Sprintf("/%s/%s/user/get", productKey, gatewayDN) cmd := ControlCmd{ Code: code, DeviceName: gatewayDN, Area: area, Address: addr, Action: action, Params: params, Identity: netPwd, } payload, err := json.Marshal(cmd) if err != nil { return fmt.Errorf("报文序列化失败: %w", err) } return MqttService.Publish(topic, payload) } func GetSetting(productKey, gatewayDN, netPwd, area, addr string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, addr, "getSetting", "") } func DeviceSendCmd(productKey, gatewayDN, netPwd, area, addr, action, params string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, addr, action, params) } // ------------------------------------------------------------------------------------- // ==================== GS200 全量 80+ 指令封装(单灯/群组/全区/空调/参数/上报/情景) ==================== // ------------------------------------------------------------------------------------- // ==================== 1. 单灯控制 ==================== func BlinkLight(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "blink", "") } func LightOn(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "lightOn", "") } func LightOff(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "lightOff", "") } func LightSleep(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "lightSleep", "") } func StopBlink(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "stopBlink", "") } func SsrOn(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "ssrOn", "") } func SsrOff(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "ssrOff", "") } func NetOn(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "netOn", "") } func NetOff(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "netOff", "") } func RelayOn(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "relayOn", "") } func RelayOff(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "relayOff", "") } // ==================== 2. 群组控制 ==================== func BlinkGroup(productKey, gatewayDN, netPwd, area, cluster string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "blink", "") } func GroupLightOn(productKey, gatewayDN, netPwd, area, cluster string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "lightOn", "") } func GroupLightOff(productKey, gatewayDN, netPwd, area, cluster string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "lightOff", "") } func GroupLightSleep(productKey, gatewayDN, netPwd, area, cluster string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "lightSleep", "") } // ==================== 3. 亮度/色温/延时/模式参数 ==================== func SetHighBright(productKey, gatewayDN, netPwd, area, number, bright string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setHighBright", bright) } func SetStandbyBright(productKey, gatewayDN, netPwd, area, number, bright string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setStandbyBright", bright) } func SetCctBright(productKey, gatewayDN, netPwd, area, number, cct string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setCctBright", cct) } func SetDelayTime(productKey, gatewayDN, netPwd, area, number, sec string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setDelayTime", sec) } func SetDelayTime2(productKey, gatewayDN, netPwd, area, number, sec string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setDelayTime2", sec) } func SetLightMode(productKey, gatewayDN, netPwd, area, number, mode string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setLightMode", mode) } func SetDelayMode(productKey, gatewayDN, netPwd, area, number, mode string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setDelayMode", mode) } func SetAlsMode(productKey, gatewayDN, netPwd, area, number, mode string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setAlsMode", mode) } // ==================== 4. 情景模式 ==================== func CallScene(productKey, gatewayDN, netPwd, area, cluster, sceneNo string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "callScene", sceneNo) } func SaveToScene(productKey, gatewayDN, netPwd, area, cluster, sceneNo string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "savetoScene", sceneNo) } func ReadScene(productKey, gatewayDN, netPwd, area, number, sceneNo string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readScene", sceneNo) } // ==================== 5. 全区/网关指令 ==================== func ScanAll(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "scan", "") } func StopScanAll(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "stopScan", "") } func GatewayReboot(productKey, gatewayDN, netPwd, area, waitSec string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "reboot", waitSec) } func SetRebootTime(productKey, gatewayDN, netPwd, area, timeRange string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "rebootSetting", timeRange) } func GatewayUpgrade(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "00 00", "upgrade", "") } // ==================== 6. 上报控制 ==================== func ReportConsumption(productKey, gatewayDN, netPwd, area, interval string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportConsumption", interval) } func StopReportConsumption(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportConsumptionStop", "") } func ReportSetting(productKey, gatewayDN, netPwd, area, interval string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportSetting", interval) } func StopReportSetting(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportSettingStop", "") } func ReportConsumptionAck(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportConsumptionAck", "") } func ReportSettingAck(productKey, gatewayDN, netPwd, area string) error { return sendCmd(400, productKey, gatewayDN, netPwd, area, "FF FF", "reportSettingAck", "") } // ==================== 8. 地址修改 ==================== func SetNumberAddress(productKey, gatewayDN, netPwd, area, number, param string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setNumberAddress", param) } func SetClusterAddress(productKey, gatewayDN, netPwd, area, cluster, newCluster string) error { return sendCmd(200, productKey, gatewayDN, netPwd, area, cluster, "setClusterAddress", newCluster) } func SetAreaAddress(productKey, gatewayDN, netPwd, area, number, newArea string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "setAreaAddress", newArea) } // ==================== 9. 传感器/报警/复位 ==================== func SsrControl(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "ssrControl", "") } func AlarmOn(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "alarmOn", "") } func AlarmOff(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "alarmOff", "") } func ResetDevice(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "reset", "") } // ==================== 10. 电量/电压/电流/功率查询 ==================== func ReadMeterData(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readMeterData", "") } func ReadVoltage(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readVoltage", "") } func ReadCurrent(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readCurrent", "") } func ReadPower(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readPower", "") } // ==================== 11. 版本/信号查询 ==================== func ReadVersion(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readVersion", "") } func ReadRssi(productKey, gatewayDN, netPwd, area, number string) error { return sendCmd(100, productKey, gatewayDN, netPwd, area, number, "readRssi", "") }