sosmgr.go 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190
  1. package main
  2. import (
  3. "runtime"
  4. "runtime/debug"
  5. "sync"
  6. "time"
  7. "github.com/sirupsen/logrus"
  8. "lc/common/models"
  9. "lc/common/mqtt"
  10. "lc/common/protocol"
  11. "lc/common/util"
  12. )
  13. // 一键报警管理
  14. var _SosMgrOnce sync.Once
  15. var _SosMgrSingle *SosMgr
  16. func GetSosMgr() *SosMgr {
  17. _SosMgrOnce.Do(func() {
  18. _SosMgrSingle = &SosMgr{
  19. queue: util.NewQueue(10000),
  20. mapSosAlarm: make(map[string]int64),
  21. mapSosState: make(map[string]*StateInfo),
  22. }
  23. })
  24. return _SosMgrSingle
  25. }
  26. type SosMgr struct {
  27. queue *util.MlQueue
  28. mapSosAlarm map[string]int64 //主题到数据库表记录id
  29. mapSosState map[string]*StateInfo ////0在线,1离线
  30. }
  31. func (o *SosMgr) SubscribeTopics() {
  32. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_SOS, protocol.TP_ONVIF_ALARM), mqtt.AtMostOnce, o.HandlerData)
  33. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_SOS, protocol.TP_ONVIF_STATE), mqtt.AtMostOnce, o.HandlerData)
  34. }
  35. func (o *SosMgr) Handler(args ...interface{}) interface{} {
  36. defer func() {
  37. if err := recover(); err != nil {
  38. time.Sleep(time.Second)
  39. gopool.Add(o.Handler, args)
  40. logrus.Errorf("SosMgr.Handler发生异常:%s", string(debug.Stack()))
  41. }
  42. }()
  43. timer := time.NewTicker(time.Duration(CheckOfflineInterval) * time.Minute)
  44. for {
  45. select {
  46. case <-timer.C: //每隔5分钟执行一次
  47. o.UpdateState()
  48. default:
  49. msg, ok, quantity := o.queue.Get()
  50. if !ok {
  51. time.Sleep(100 * time.Millisecond)
  52. continue
  53. } else if quantity > 1000 {
  54. logrus.Warnf("SosMgr.Handler:数据队列累积过多,请注意优化,当前队列条数:%d", quantity)
  55. }
  56. m, ok := msg.(*mqtt.Message)
  57. if !ok {
  58. continue
  59. }
  60. _, _, _, topic, err := ParseTopic(m.Topic())
  61. if err != nil {
  62. continue
  63. }
  64. switch topic {
  65. case protocol.TP_ONVIF_STATE:
  66. o.HandlerState(m)
  67. case protocol.TP_ONVIF_ALARM:
  68. o.HandlerAlarm(m)
  69. default:
  70. logrus.Warnf("SosMgr.Handler:收到暂不支持的主题:%s", topic)
  71. }
  72. }
  73. }
  74. }
  75. func (o *SosMgr) HandlerData(m mqtt.Message) {
  76. for {
  77. ok, cnt := o.queue.Put(&m)
  78. if ok {
  79. break
  80. } else {
  81. logrus.Errorf("SosMgr.HandlerData:查询队列失败,队列消息数量:%d", cnt)
  82. runtime.Gosched()
  83. }
  84. }
  85. }
  86. func (o *SosMgr) HandlerAlarm(m *mqtt.Message) {
  87. var obj protocol.Pack_OnvifAlarm
  88. if err := obj.DeCode(m.PayloadString()); err != nil {
  89. logrus.Errorf("数据解析失败,主题:%s,内容:%s,失败原因:%s", m.Topic(), m.PayloadString(), err.Error())
  90. return
  91. }
  92. t, err := util.MlParseTime(obj.Time)
  93. if err != nil {
  94. logrus.Errorf("时间[%s]解析错误:%s", obj.Time, err.Error())
  95. return
  96. }
  97. if obj.Data.Alarm.AlarmTopic == "tns1:Device/Trigger/tnshik:AlarmIn" { //海康威视一键报警
  98. state, ok2 := obj.Data.Alarm.Data["State"]
  99. if !ok2 {
  100. return
  101. }
  102. //更新数据库告警结束时间
  103. if state == "false" {
  104. if id, ok := o.mapSosAlarm[obj.Id+obj.Data.Alarm.AlarmTopic]; ok {
  105. oo := models.SosAlarm{ID: id, TEnd: t}
  106. if err := oo.Update(); err != nil {
  107. logrus.Errorf("一键告警数据:%v", obj)
  108. logrus.Errorf("一键告警数据更新失败:%s", err.Error())
  109. }
  110. delete(o.mapSosAlarm, obj.Id+obj.Data.Alarm.AlarmTopic)
  111. }
  112. } else { //新建告警记录
  113. oo := models.SosAlarm{
  114. DID: obj.Id,
  115. TStart: t,
  116. AType: "SOS",
  117. Content: "一键求助",
  118. }
  119. if err := models.G_db.Create(&oo).Error; err != nil {
  120. logrus.Errorf("一键告警数据:%v", obj)
  121. logrus.Errorf("一键告警数据入库失败:%s", err.Error())
  122. } else {
  123. o.mapSosAlarm[obj.Id+obj.Data.Alarm.AlarmTopic] = oo.ID
  124. }
  125. }
  126. }
  127. }
  128. func (o *SosMgr) HandlerState(m *mqtt.Message) {
  129. var obj protocol.Pack_IPCState
  130. if err := obj.DeCode(m.PayloadString()); err != nil {
  131. logrus.Errorf("数据解析失败,主题:%s,内容:%s,失败原因:%s", m.Topic(), m.PayloadString(), err.Error())
  132. return
  133. }
  134. t, err := util.MlParseTime(obj.Time)
  135. if err != nil {
  136. logrus.Errorf("时间[%s]解析错误:%s", obj.Time, err.Error())
  137. return
  138. }
  139. //0在线,1离线
  140. si, ok := o.mapSosState[obj.Id]
  141. if !ok {
  142. t, s, err := getState(obj.Id)
  143. if err != nil {
  144. cacheState(obj.Id, obj.Time, obj.Data.State)
  145. o.mapSosState[obj.Id] = &StateInfo{Time: t, State: obj.Data.State}
  146. return
  147. } else {
  148. si = &StateInfo{Time: t, State: s}
  149. o.mapSosState[obj.Id] = si
  150. }
  151. }
  152. if si.State != obj.Data.State {
  153. if obj.Data.State == 1 { //在线到离线
  154. GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_OFFLINE, Time: t})
  155. } else { //离线到在线
  156. GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_ONLINE, Time: t})
  157. }
  158. }
  159. cacheState(obj.Id, obj.Time, obj.Data.State)
  160. o.mapSosState[obj.Id].State = obj.Data.State
  161. o.mapSosState[obj.Id].Time = t
  162. }
  163. func (o *SosMgr) UpdateState() {
  164. t := util.MlNow()
  165. for k, v := range o.mapSosState {
  166. if v.State == 0 && t.Sub(v.Time).Minutes() > OfflineInterval { //只检查当前还在线的
  167. GetEventMgr().PushEvent(&EventObject{ID: k, EventType: models.ET_OFFLINE, Time: t})
  168. cacheState(k, t.Format("2006-01-02 15:04:05"), 1)
  169. o.mapSosState[k].State = 1
  170. o.mapSosState[k].Time = t
  171. }
  172. }
  173. }