cableGuardianHandler.go 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. package main
  2. import (
  3. "fmt"
  4. "github.com/jinzhu/gorm"
  5. "runtime"
  6. "runtime/debug"
  7. "strconv"
  8. "sync"
  9. "time"
  10. "github.com/sirupsen/logrus"
  11. "lc/common/models"
  12. "lc/common/mqtt"
  13. "lc/common/protocol"
  14. "lc/common/util"
  15. )
  16. const (
  17. CableStatusNormal = iota //电缆状态 正常
  18. CableStatusBeStolen //电缆状态 被盗
  19. CableStatusOpened //电缆状态 被打开
  20. CableStatusBeStolenAndOpened //电缆状态 被盗、被打开
  21. )
  22. const (
  23. CableStatusNormalStr = "正常" //电缆状态 正常
  24. CableStatusBeStolenStr = "被盗" //电缆状态 被盗
  25. CableStatusOpenedStr = "被打开" //电缆状态 被打开
  26. CableStatusBeStolenAndOpenedStr = "被盗、被打开" //电缆状态 被盗、被打开
  27. )
  28. const cableGuardianDataPrefix = "cable_guardian_data_%s_%s_%d"
  29. // 电缆防盗 mqtt消息处理
  30. var _cableGuardianHandlerOnce sync.Once
  31. var _cableGuardianHandlerSingle *cableGuardianHandler
  32. func GetCableGuardianHandler() *cableGuardianHandler {
  33. _cableGuardianHandlerOnce.Do(func() {
  34. _cableGuardianHandlerSingle = &cableGuardianHandler{
  35. queue: util.NewQueue(10000),
  36. }
  37. })
  38. return _cableGuardianHandlerSingle
  39. }
  40. type cableGuardianHandler struct {
  41. queue *util.MlQueue
  42. }
  43. func (o *cableGuardianHandler) SubscribeTopics() {
  44. //电缆防盗
  45. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_CableGuardian, protocol.TP_MODBUS_DATA), mqtt.AtMostOnce, o.HandlerData)
  46. }
  47. func (o *cableGuardianHandler) HandlerData(m mqtt.Message) {
  48. for {
  49. ok, cnt := o.queue.Put(&m)
  50. if ok {
  51. break
  52. } else {
  53. logrus.Errorf("cableGuardianHandler.HandlerData:查询队列失败,队列消息数量:%d", cnt)
  54. runtime.Gosched()
  55. }
  56. }
  57. }
  58. func (o *cableGuardianHandler) Handler(args ...interface{}) interface{} {
  59. defer func() {
  60. if err := recover(); err != nil {
  61. time.Sleep(time.Second)
  62. gopool.Add(o.Handler, args)
  63. logrus.Errorf("cableGuardianHandler.Handler:%v发生异常:%s", args, string(debug.Stack()))
  64. }
  65. }()
  66. for {
  67. msg, ok, quantity := o.queue.Get()
  68. if !ok {
  69. time.Sleep(10 * time.Millisecond)
  70. continue
  71. } else if quantity > 1000 {
  72. logrus.Warnf("数据队列累积过多,请注意优化,当前队列条数:%d", quantity)
  73. }
  74. m, ok := msg.(*mqtt.Message)
  75. if !ok {
  76. continue
  77. }
  78. _, _, DID, topic, err := ParseTopic(m.Topic())
  79. if err != nil {
  80. continue
  81. }
  82. switch topic {
  83. case protocol.TP_MODBUS_DATA:
  84. var ret protocol.Pack_UploadData
  85. if err := ret.DeCode(m.PayloadString()); err == nil {
  86. if ret.Data.State == protocol.FAILED {
  87. logrus.Warningf("电缆防盗数据不正确 %+v", ret)
  88. continue
  89. }
  90. t, _ := util.MlParseTime(ret.Time)
  91. for id, value := range ret.Data.Data {
  92. key := fmt.Sprintf(cableGuardianDataPrefix, ret.Gid, DID, id)
  93. tId := strconv.Itoa(int(id))
  94. old := getCableData(key)
  95. if value != old {
  96. cableGuardianStatus := &models.CableGuardianStatus{
  97. GID: ret.Gid,
  98. DID: DID,
  99. TerminalID: tId,
  100. Status: int(value),
  101. }
  102. err := cableGuardianStatus.Get()
  103. if err != nil {
  104. if !gorm.IsRecordNotFoundError(err) {
  105. logrus.Warnf("CableGuardianStatus get fail = %v", err)
  106. continue
  107. }
  108. gateway := models.Gateway{ID: ret.Gid}
  109. err = gateway.Get()
  110. if err != nil {
  111. logrus.Warnf("CableGuardianStatus get gateway fail = %v", err)
  112. continue
  113. }
  114. cableGuardianStatus.GatewayName = gateway.Name
  115. }
  116. if cableGuardianStatus.ID > 0 {
  117. cableGuardianStatus.UpdateAt = t
  118. cableGuardianStatus.CreatedAt = t
  119. cableGuardianStatus.Status = int(value)
  120. err = cableGuardianStatus.Update()
  121. } else {
  122. cableGuardianStatus.CreatedAt = t
  123. cableGuardianStatus.UpdateAt = t
  124. err = cableGuardianStatus.Save()
  125. }
  126. cacheCableData(key, value)
  127. if old != -1 {
  128. err = sendSms([]string{cableGuardianStatus.GatewayName, tId, getSmsStr(value)})
  129. }
  130. }
  131. }
  132. }
  133. default:
  134. logrus.Warnf("cableGuardianHandler.Handler:收到暂不支持的主题:%s", topic)
  135. }
  136. }
  137. }
  138. func getSmsStr(status float64) string {
  139. switch status {
  140. case CableStatusNormal:
  141. return CableStatusNormalStr
  142. case CableStatusBeStolen:
  143. return CableStatusBeStolenStr
  144. case CableStatusOpened:
  145. return CableStatusOpenedStr
  146. case CableStatusBeStolenAndOpened:
  147. return CableStatusBeStolenAndOpenedStr
  148. }
  149. return ""
  150. }
  151. // 缓存最新数据到redis
  152. func cacheCableData(key string, data float64) {
  153. if err := redisCltRawData.Set(key, data, 0).Err(); err != nil {
  154. logrus.Errorf("cacheCableData err = ", err.Error())
  155. }
  156. }
  157. // 获取缓存的redis数据
  158. func getCableData(key string) float64 {
  159. var value float64
  160. if err := redisCltRawData.Get(key).Scan(&value); err != nil {
  161. logrus.Warningf("getCableData err = %s", err.Error())
  162. return -1
  163. }
  164. return value
  165. }