gwhandler.go 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382
  1. package main
  2. import (
  3. "os"
  4. "path/filepath"
  5. "runtime"
  6. "runtime/debug"
  7. "sync"
  8. "time"
  9. "github.com/sirupsen/logrus"
  10. "lc/common/models"
  11. "lc/common/mqtt"
  12. "lc/common/protocol"
  13. "lc/common/util"
  14. )
  15. var _gwHandlerOnce sync.Once
  16. var _gwHandlerSingle *GwHandler
  17. func GetGwHandler() *GwHandler {
  18. _gwHandlerOnce.Do(func() {
  19. _gwHandlerSingle = &GwHandler{
  20. queue: util.NewQueue(10000),
  21. }
  22. })
  23. return _gwHandlerSingle
  24. }
  25. type GwHandler struct {
  26. queue *util.MlQueue
  27. }
  28. func (o *GwHandler) SubscribeTopics() {
  29. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ONLINE), mqtt.AtLeastOnce, o.HandlerData)
  30. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_WILL), mqtt.AtLeastOnce, o.HandlerData)
  31. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_APP_ACK), mqtt.AtMostOnce, o.HandlerData)
  32. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SET_APP_ACK), mqtt.AtMostOnce, o.HandlerData)
  33. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SERIAL_ACK), mqtt.AtMostOnce, o.HandlerData)
  34. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SET_SERIAL_ACK), mqtt.AtMostOnce, o.HandlerData)
  35. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_RTU_ACK), mqtt.AtMostOnce, o.HandlerData)
  36. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SET_RTU_ACK), mqtt.AtMostOnce, o.HandlerData)
  37. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_MODEL_ACK), mqtt.AtMostOnce, o.HandlerData)
  38. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SET_MODEL_ACK), mqtt.AtMostOnce, o.HandlerData)
  39. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_LOG_ACK), mqtt.AtMostOnce, o.HandlerData)
  40. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_REMOVE_LOG_ACK), mqtt.AtMostOnce, o.HandlerData)
  41. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_LOG_CFG_ACK), mqtt.AtMostOnce, o.HandlerData)
  42. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SYS_ACK), mqtt.AtMostOnce, o.HandlerData)
  43. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ITS_ACK), mqtt.AtMostOnce, o.HandlerData)
  44. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ONVIFDEV_ACK), mqtt.AtMostOnce, o.HandlerData)
  45. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_DEPLOY_ACK), mqtt.AtMostOnce, o.HandlerData)
  46. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_DEPLOY_HEALTH), mqtt.AtMostOnce, o.HandlerData)
  47. }
  48. func (o *GwHandler) HandlerData(m mqtt.Message) {
  49. for {
  50. ok, cnt := o.queue.Put(&m)
  51. if ok {
  52. break
  53. } else {
  54. logrus.Errorf("GwHandler.HandlerData:查询队列失败,队列消息数量:%d", cnt)
  55. runtime.Gosched()
  56. }
  57. }
  58. }
  59. func (o *GwHandler) Handler(args ...interface{}) interface{} {
  60. defer func() {
  61. if err := recover(); err != nil {
  62. time.Sleep(time.Second)
  63. gopool.Add(o.Handler, args)
  64. logrus.Errorf("GwHandler.Handler:%v发生异常:%s", args, string(debug.Stack()))
  65. }
  66. }()
  67. for {
  68. msg, ok, quantity := o.queue.Get()
  69. if !ok {
  70. time.Sleep(10 * time.Millisecond)
  71. continue
  72. } else if quantity > 1000 {
  73. logrus.Warnf("数据队列累积过多,请注意优化,当前队列条数:%d", quantity)
  74. }
  75. m, ok := msg.(*mqtt.Message)
  76. if !ok {
  77. continue
  78. }
  79. Tenant, _, GID, topic, err := ParseTopic(m.Topic())
  80. if err != nil {
  81. continue
  82. }
  83. switch topic {
  84. case protocol.TP_GW_ONLINE: //上线
  85. var obj protocol.Pack_IDObject
  86. if err := obj.DeCode(m.PayloadString()); err == nil { //网关在线
  87. cacheState(obj.Id, obj.Time, 0)
  88. GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_ONLINE, Time: util.MlNow()})
  89. }
  90. case protocol.TP_GW_WILL: //下线
  91. var obj protocol.Pack_IDObject
  92. if err := obj.DeCode(m.PayloadString()); err == nil { //网关离线
  93. cacheState(obj.Id, obj.Time, 1)
  94. GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_OFFLINE, Time: util.MlNow()})
  95. }
  96. case protocol.TP_GW_SET_APP_ACK, protocol.TP_GW_SET_MODEL_ACK, protocol.TP_GW_SET_RTU_ACK, protocol.TP_GW_SET_SERIAL_ACK, protocol.TP_GW_REMOVE_LOG_ACK, protocol.TP_GW_LOG_CFG_ACK:
  97. var obj protocol.Pack_Ack
  98. if err := obj.DeCode(m.PayloadString()); err == nil {
  99. o := models.DeviceCmdRecord{
  100. ID: obj.Seq,
  101. State: 1,
  102. Resp: obj.Data.Error,
  103. }
  104. if err := o.Update(); err != nil {
  105. logrus.Errorf("收到网关[%s]的响应[seq:%d],主题:%s,但更新数据库失败[%s]",
  106. obj.Id, obj.Seq, m.Topic(), err.Error())
  107. }
  108. }
  109. case protocol.TP_GW_APP:
  110. var ret protocol.Pack_MutilFileObject
  111. if err := ret.DeCode(m.PayloadString()); err == nil && len(ret.Data.Files) == 1 {
  112. SaveFile(Tenant, GID, "conf", ret.Data.Files)
  113. var obj protocol.AppConfig
  114. if err := json.UnmarshalFromString(ret.Data.Files[0].Content, &obj); err == nil {
  115. o := models.Gateway{
  116. ID: ret.Id,
  117. Name: obj.Name,
  118. Tenant: obj.Tenant,
  119. Sn: obj.SN,
  120. Upgrade: obj.Upgrade,
  121. MqttEdgeServer: obj.Edge.Mqtt.Server,
  122. MqttEdgeUser: obj.Edge.Mqtt.User,
  123. MqttEdgePassword: obj.Edge.Mqtt.Password,
  124. MqttCloudServer: obj.Cloud.Mqtt.Server,
  125. MqttCloudUser: obj.Cloud.Mqtt.User,
  126. MqttCloudPassword: obj.Cloud.Mqtt.Password,
  127. State: 1,
  128. }
  129. if err := o.SaveFromGateway(); err != nil {
  130. logrus.Errorf("插入数据库失败:%s", err.Error())
  131. }
  132. }
  133. }
  134. case protocol.TP_GW_SERIAL_ACK:
  135. var ret protocol.Pack_MutilFileObject
  136. if err := ret.DeCode(m.PayloadString()); err == nil && len(ret.Data.Files) == 1 {
  137. SaveFile(Tenant, GID, "conf", ret.Data.Files)
  138. var obj protocol.SerialConfig
  139. if err := json.UnmarshalFromString(ret.Data.Files[0].Content, &obj); err == nil {
  140. for _, v := range obj.Serial {
  141. o := models.GatewaySerial{
  142. ID: ret.Id,
  143. ComID: int(v.Code),
  144. Interface: v.Interface,
  145. Address: v.Address,
  146. BaudRate: v.BaudRate,
  147. DataBits: int(v.DataBits),
  148. StopBits: int(v.StopBits),
  149. Parity: v.Parity,
  150. Timeout: int(v.Timeout),
  151. ProtocolType: int(v.ProtocolType),
  152. }
  153. if err := models.G_db.Save(&o).Error; err != nil {
  154. logrus.Errorf("插入数据库失败:%s", err.Error())
  155. }
  156. }
  157. }
  158. }
  159. case protocol.TP_GW_RTU_ACK:
  160. var ret protocol.Pack_MutilFileObject
  161. if err := ret.DeCode(m.PayloadString()); err == nil {
  162. SaveFile(Tenant, GID, "dev", ret.Data.Files)
  163. for _, v := range ret.Data.Files {
  164. var obj protocol.MapDevConfig
  165. if err := json.UnmarshalFromString(v.Content, &obj); err == nil {
  166. for _, v1 := range obj.Rtu {
  167. o := models.GatewayDevice{
  168. ID: v1.DevCode,
  169. Name: v1.Name,
  170. GID: ret.Id,
  171. ComID: int(v1.Code),
  172. RtuID: int(v1.DevID),
  173. TID: int(v1.TID),
  174. SendCloud: v1.SendCloud,
  175. WaitTime: int(v1.WaitTime),
  176. ProtocolType: int(v1.ProtocolType),
  177. DevType: int(v1.DevType),
  178. Tenant: Tenant,
  179. State: 1,
  180. }
  181. if err := o.SaveFromGateway(); err != nil {
  182. logrus.Errorf("插入数据库失败:%s", err.Error())
  183. }
  184. //设备类型,1-灯控类设备 2-环境监测类设备 3-裕明485单灯控制器 4-液位计 5-路面状况传感器
  185. if v1.DevType == 4 || v1.DevType == 5 {
  186. if err := models.UpdateDeviceSensorTID(v1.DevCode, int(v1.TID)); err != nil {
  187. logrus.Errorf("更新传感器物模型失败:%s", err.Error())
  188. }
  189. } else if v1.DevType == 2 { //环境监测设备
  190. if err := models.UpdateDeviceEnvironmentTID(v1.DevCode, int(v1.TID)); err != nil {
  191. logrus.Errorf("更新环境传感器物模型失败:%s", err.Error())
  192. }
  193. } else if v1.DevType == 3 { //裕明单灯控制器
  194. if err := models.UpdateDeviceLampControllerTID(v1.DevCode, int(v1.TID)); err != nil {
  195. logrus.Errorf("更新灯控物模型失败:%s", err.Error())
  196. }
  197. } else if v1.DevType == 1 { //集控器
  198. if err := models.UpdateTID(v1.DevCode, int(v1.TID)); err != nil {
  199. logrus.Errorf("更新单灯集控器物模型失败:%s", err.Error())
  200. }
  201. }
  202. }
  203. }
  204. }
  205. }
  206. case protocol.TP_GW_MODEL_ACK:
  207. var ret protocol.Pack_MutilFileObject
  208. if err := ret.DeCode(m.PayloadString()); err == nil {
  209. SaveFile(Tenant, GID, "model", ret.Data.Files)
  210. }
  211. case protocol.TP_GW_LOG_ACK:
  212. var ret protocol.Pack_MutilFileObject
  213. if err := ret.DeCode(m.PayloadString()); err == nil {
  214. SaveFile(Tenant, GID, "log", ret.Data.Files)
  215. }
  216. case protocol.TP_GW_SYS_ACK:
  217. var ret protocol.Pack_SysInfo
  218. if err := ret.DeCode(m.PayloadString()); err == nil {
  219. o := models.GatewaySysInfo{
  220. GID: GID,
  221. AppName: ret.Data.Appinfo.Name,
  222. AppVersion: ret.Data.Appinfo.Version,
  223. CpuCnt: ret.Data.Cpuinfo.Cpus,
  224. CpuCores: ret.Data.Cpuinfo.Cores,
  225. CpuModelName: ret.Data.Cpuinfo.ModelName,
  226. CpuPercent: ret.Data.Cpuinfo.Percent,
  227. MemTotal: ret.Data.Meminfo.Total,
  228. MemAvailable: ret.Data.Meminfo.Available,
  229. MemUsed: ret.Data.Meminfo.Used,
  230. MemPercent: ret.Data.Meminfo.Percent,
  231. }
  232. if str, err := json.MarshalIndent(&ret.Data.Diskinfos, "", " "); err == nil {
  233. o.DiskInfos = string(str)
  234. }
  235. if str, err := json.MarshalIndent(&ret.Data.Ifs, "", " "); err == nil {
  236. o.NetIfs = string(str)
  237. }
  238. if str, err := json.MarshalIndent(&ret.Data.Pis, "", " "); err == nil {
  239. o.Process = string(str)
  240. }
  241. if str, err := json.MarshalIndent(&ret.Data.TcpListen, "", " "); err == nil {
  242. o.TcpListen = string(str)
  243. }
  244. if str, err := json.MarshalIndent(&ret.Data.TcpConn, "", " "); err == nil {
  245. o.TcpConn = string(str)
  246. }
  247. if str, err := json.MarshalIndent(&ret.Data.Udp, "", " "); err == nil {
  248. o.Udp = string(str)
  249. }
  250. if err := models.G_db.Save(&o).Error; err != nil {
  251. logrus.Errorf("插入数据库失败:%s", err.Error())
  252. }
  253. }
  254. case protocol.TP_GW_ITS_ACK:
  255. var ret protocol.Pack_ITSDev
  256. if err := ret.DeCode(m.PayloadString()); err == nil {
  257. for _, v := range ret.Data.Its {
  258. o := models.ItsDevice{
  259. ID: v.ID,
  260. Name: v.Name,
  261. GID: ret.Gid,
  262. Brand: v.Brand,
  263. Model: v.Model,
  264. DevType: v.DevType,
  265. User: v.User,
  266. Password: v.Password,
  267. IP: v.IP,
  268. Port: v.Port,
  269. HttpAddr: ret.Data.IPAddr,
  270. SuggestSpeed: int(ret.Data.SuggestSpeed),
  271. Duration: ret.Data.Duration,
  272. EnvID: ret.Data.EnvID,
  273. TollgateID: v.TollgateID,
  274. Tenant: Tenant,
  275. State: 1,
  276. }
  277. if err := o.SaveFromGateway(); err != nil {
  278. logrus.Errorf("抓拍单元数据入库失败:%s", err.Error())
  279. logrus.Errorf("抓拍单元数据入库失败:%v", o)
  280. }
  281. }
  282. }
  283. case protocol.TP_GW_ONVIFDEV_ACK:
  284. var ret protocol.Pack_OnvifDev
  285. if err := ret.DeCode(m.PayloadString()); err == nil {
  286. for _, v := range ret.Data {
  287. o := models.CameraDevice{
  288. ID: v.Code,
  289. Name: v.Name,
  290. GID: ret.Id,
  291. IP: v.IP,
  292. SN: v.SN,
  293. Brand: v.Brand,
  294. Model: v.Model,
  295. DevType: v.DevType,
  296. User: v.User,
  297. Password: v.Password,
  298. RtmpServer: v.RtmpServer,
  299. WebServer: v.WebServer,
  300. Event: v.Event,
  301. Gb28181: v.Gb28181,
  302. State: 1,
  303. }
  304. if err := o.SaveFromGateway(); err != nil {
  305. logrus.Errorf("摄像头数据入库失败:%s", err.Error())
  306. logrus.Errorf("摄像头数据入库失败:%v", o)
  307. }
  308. }
  309. }
  310. case protocol.TP_GW_DEPLOY_ACK:
  311. var ack protocol.Pack_DeployAck
  312. if err := ack.DeCode(m.PayloadString()); err != nil {
  313. util.GetTagLog().Errorf("sys", "DEPLOY_ACK DeCode失败: %v", err)
  314. break
  315. }
  316. util.GetTagLog().Infof("sys", "DEPLOY_ACK: gid=%s version=%s success=%v", ack.Header.Gid, ack.Data.Version, ack.Data.Success)
  317. status := uint8(1)
  318. if !ack.Data.Success {
  319. status = 2
  320. }
  321. models.G_db.Model(&models.GatewayDeploy{}).
  322. Where("gid = ? AND to_version = ? AND status = 0", ack.Header.Gid, ack.Data.Version).
  323. Updates(map[string]interface{}{
  324. "status": status,
  325. "error_msg": ack.Data.Error,
  326. "update_time": time.Now(),
  327. })
  328. case protocol.TP_GW_DEPLOY_HEALTH:
  329. var health protocol.Pack_HealthCheck
  330. if err := health.DeCode(m.PayloadString()); err != nil {
  331. util.GetTagLog().Errorf("sys", "DEPLOY_HEALTH DeCode失败: %v", err)
  332. break
  333. }
  334. util.GetTagLog().Infof("sys", "DEPLOY_HEALTH: gid=%s version=%s phase=%d success=%v",
  335. health.Header.Gid, health.Data.Version, health.Data.Phase, health.Data.Success)
  336. resultJSON, _ := json.Marshal(health.Data)
  337. updateFields := map[string]interface{}{
  338. "update_time": time.Now(),
  339. }
  340. if health.Data.Phase == 1 {
  341. updateFields["phase1_result"] = string(resultJSON)
  342. if !health.Data.Success {
  343. updateFields["status"] = uint8(3)
  344. }
  345. } else {
  346. updateFields["phase2_result"] = string(resultJSON)
  347. }
  348. models.G_db.Model(&models.GatewayDeploy{}).
  349. Where("gid = ? AND to_version = ?", health.Header.Gid, health.Data.Version).
  350. Updates(updateFields)
  351. default:
  352. logrus.Warnf("GwHandler.Handler:收到暂不支持的主题:%s", topic)
  353. }
  354. }
  355. }
  356. func SaveFile(tenant, GID, strType string, fo []protocol.FileObject) {
  357. Dir := "/opt/ipole_data/gateway_logs" + string(filepath.Separator) + tenant + string(filepath.Separator) + GID + string(filepath.Separator) + strType + string(filepath.Separator)
  358. err := os.MkdirAll(Dir, os.ModePerm)
  359. if err != nil {
  360. return
  361. }
  362. for _, v := range fo {
  363. if err := os.WriteFile(Dir+v.File, []byte(v.Content), os.ModePerm); err != nil {
  364. logrus.Errorf("SaveFile:保存模型文件失败,文件名:%s,原因:%s", Dir+v.File, err.Error())
  365. }
  366. }
  367. }