| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390 |
- 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", "")
- }
|