ipcmgr.go 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247
  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. var _IpcMgrOnce sync.Once
  14. var _IpcMgrSingle *IPCMgr
  15. func GetIPCMgr() *IPCMgr {
  16. _IpcMgrOnce.Do(func() {
  17. _IpcMgrSingle = &IPCMgr{
  18. queue: util.NewQueue(10000),
  19. //mapSosAlarm: make(map[string]int64),
  20. mapIpcState: make(map[string]*StateInfo),
  21. }
  22. })
  23. return _IpcMgrSingle
  24. }
  25. type StateInfo struct {
  26. Time time.Time
  27. State uint8
  28. }
  29. type IPCMgr struct {
  30. queue *util.MlQueue
  31. //mapSosAlarm map[string]int64 //主题到数据库表记录id
  32. mapIpcState map[string]*StateInfo ////0在线,1离线
  33. }
  34. func (o *IPCMgr) SubscribeTopics() {
  35. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_IPC, protocol.TP_ONVIF_ALARM), mqtt.AtMostOnce, o.HandlerData)
  36. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_IPC, protocol.TP_ONVIF_STATE), mqtt.AtMostOnce, o.HandlerData)
  37. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_IPC, protocol.TP_ONVIF_PRESETS_ACK), mqtt.AtMostOnce, o.HandlerData)
  38. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_IPC, protocol.TP_ONVIF_PRESET_ACK), mqtt.AtMostOnce, o.HandlerData)
  39. }
  40. func (o *IPCMgr) Handler(args ...interface{}) interface{} {
  41. defer func() {
  42. if err := recover(); err != nil {
  43. time.Sleep(time.Second)
  44. gopool.Add(o.Handler, args)
  45. logrus.Errorf("IPCMgr.Handler发生异常:%v", err)
  46. logrus.Errorf("IPCMgr.Handler发生异常:%s", string(debug.Stack()))
  47. }
  48. }()
  49. timer := time.NewTicker(time.Duration(CheckOfflineInterval) * time.Minute)
  50. for {
  51. select {
  52. case <-timer.C: //每隔5分钟执行一次
  53. o.UpdateState()
  54. default:
  55. msg, ok, quantity := o.queue.Get()
  56. if !ok {
  57. time.Sleep(100 * time.Millisecond)
  58. continue
  59. } else if quantity > 1000 {
  60. logrus.Warnf("IPCMgr.Handler:数据队列累积过多,请注意优化,当前队列条数:%d", quantity)
  61. }
  62. m, ok := msg.(*mqtt.Message)
  63. if !ok {
  64. continue
  65. }
  66. _, _, _, topic, err := ParseTopic(m.Topic())
  67. if err != nil {
  68. continue
  69. }
  70. switch topic {
  71. case protocol.TP_ONVIF_STATE:
  72. o.HandlerState(m)
  73. case protocol.TP_ONVIF_ALARM:
  74. o.HandlerAlarm(m)
  75. case protocol.TP_ONVIF_PRESETS_ACK:
  76. o.HandlerGetPreset(m)
  77. case protocol.TP_ONVIF_PRESET_ACK:
  78. o.HandlerSetPreset(m)
  79. default:
  80. logrus.Warnf("IPCMgr.Handler:收到暂不支持的主题:%s", topic)
  81. }
  82. }
  83. }
  84. }
  85. func (o *IPCMgr) HandlerData(m mqtt.Message) {
  86. for {
  87. ok, cnt := o.queue.Put(&m)
  88. if ok {
  89. break
  90. } else {
  91. logrus.Errorf("IPCMgr.HandlerData:查询队列失败,队列消息数量:%d", cnt)
  92. runtime.Gosched()
  93. }
  94. }
  95. }
  96. func (o *IPCMgr) HandlerAlarm(m *mqtt.Message) {
  97. var obj protocol.Pack_OnvifAlarm
  98. if err := obj.DeCode(m.PayloadString()); err != nil {
  99. logrus.Errorf("数据解析失败,主题:%s,内容:%s,失败原因:%s", m.Topic(), m.PayloadString(), err.Error())
  100. return
  101. }
  102. t, err := util.MlParseTime(obj.Time)
  103. if err != nil {
  104. logrus.Errorf("时间[%s]解析错误:%s", obj.Time, err.Error())
  105. return
  106. }
  107. if obj.Data.Alarm.AlarmTopic == "tns1:RuleEngine/LineDetector/Crossed" {
  108. oo := models.IPCAlarm{
  109. DID: obj.Id,
  110. TStart: t,
  111. TEnd: t,
  112. AType: "Crossed",
  113. Content: "遮挡告警",
  114. }
  115. if err := models.G_db.Create(&oo).Error; err != nil {
  116. logrus.Errorf("遮挡告警数据:%v", obj)
  117. logrus.Errorf("遮挡告警数据入库失败:%s", err.Error())
  118. }
  119. }
  120. }
  121. func (o *IPCMgr) HandlerState(m *mqtt.Message) {
  122. var obj protocol.Pack_IPCState
  123. if err := obj.DeCode(m.PayloadString()); err != nil {
  124. logrus.Errorf("数据解析失败,主题:%s,内容:%s,失败原因:%s", m.Topic(), m.PayloadString(), err.Error())
  125. return
  126. }
  127. t, err := util.MlParseTime(obj.Time)
  128. if err != nil {
  129. logrus.Errorf("时间[%s]解析错误:%s", obj.Time, err.Error())
  130. return
  131. }
  132. //0在线,1离线
  133. si, ok := o.mapIpcState[obj.Id]
  134. if !ok {
  135. t, s, err := getState(obj.Id)
  136. if err != nil {
  137. cacheState(obj.Id, obj.Time, obj.Data.State)
  138. o.mapIpcState[obj.Id] = &StateInfo{Time: t, State: obj.Data.State}
  139. return
  140. } else {
  141. si = &StateInfo{Time: t, State: s}
  142. o.mapIpcState[obj.Id] = si
  143. }
  144. }
  145. if si.State != obj.Data.State {
  146. if obj.Data.State == 1 { //在线到离线
  147. GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_OFFLINE, Time: t})
  148. } else { //离线到在线
  149. GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_ONLINE, Time: t})
  150. }
  151. }
  152. cacheState(obj.Id, obj.Time, obj.Data.State)
  153. o.mapIpcState[obj.Id].State = obj.Data.State
  154. o.mapIpcState[obj.Id].Time = t
  155. }
  156. func (o *IPCMgr) UpdateState() {
  157. t := util.MlNow()
  158. for k, v := range o.mapIpcState {
  159. if v.State == 0 && t.Sub(v.Time).Minutes() > OfflineInterval { //只检查当前还在线的
  160. GetEventMgr().PushEvent(&EventObject{ID: k, EventType: models.ET_OFFLINE, Time: t})
  161. cacheState(k, t.Format("2006-01-02 15:04:05"), 1)
  162. o.mapIpcState[k].State = 1
  163. o.mapIpcState[k].Time = t
  164. }
  165. }
  166. }
  167. func (o *IPCMgr) HandlerGetPreset(m *mqtt.Message) {
  168. var obj protocol.Pack_PresetInfo
  169. if err := obj.DeCode(m.PayloadString()); err != nil {
  170. logrus.Errorf("数据解析失败,主题:%s,内容:%s,失败原因:%s", m.Topic(), m.PayloadString(), err.Error())
  171. return
  172. }
  173. if obj.Data.State == protocol.FAILED {
  174. logrus.Errorf("读取设备[%s]预置点失败:%s", obj.Id, obj.Data.Error)
  175. return
  176. }
  177. //所有预置点不存库,用于调试
  178. if obj.Data.Flag == 6 {
  179. logrus.Debugf("设备[%s]预置点:%v", obj.Id, obj.Data.Presets)
  180. return
  181. }
  182. if len(obj.Data.Presets) == 0 {
  183. return
  184. }
  185. t := util.MlNow()
  186. for _, v := range obj.Data.Presets {
  187. ps := models.IpcPresets{
  188. ID: obj.Id,
  189. Token: v.Token,
  190. Name: v.Name,
  191. X: v.X,
  192. Y: v.Y,
  193. Z: v.Z,
  194. CreatedAt: t,
  195. }
  196. if err := ps.SaveFromGateway(); err != nil {
  197. logrus.Errorf("保存设置预置位信息失败:%s", err.Error())
  198. logrus.Debugf("预置位信息:%v", ps)
  199. }
  200. }
  201. }
  202. func (o *IPCMgr) HandlerSetPreset(m *mqtt.Message) {
  203. var obj protocol.Pack_IPCSetPresetACK
  204. if err := obj.DeCode(m.PayloadString()); err != nil {
  205. logrus.Errorf("数据解析失败,主题:%s,内容:%s,失败原因:%s", m.Topic(), m.PayloadString(), err.Error())
  206. return
  207. }
  208. if obj.Data.State == protocol.FAILED {
  209. logrus.Errorf("对设备[%s]设置预置点失败:%s", obj.Id, obj.Data.Error)
  210. return
  211. }
  212. ps := models.IpcPresets{
  213. ID: obj.Id,
  214. Token: obj.Data.Token,
  215. Name: obj.Data.Name,
  216. X: obj.Data.X,
  217. Y: obj.Data.Y,
  218. Z: obj.Data.Z,
  219. File: obj.Data.File,
  220. }
  221. if err := ps.SaveFromGateway2(); err != nil {
  222. logrus.Errorf("保存设置预置位信息失败:%s", err.Error())
  223. logrus.Debugf("预置位信息:%v", ps)
  224. }
  225. }