ymlampcontrollermgr.go 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251
  1. package main
  2. import (
  3. "context"
  4. "runtime"
  5. "runtime/debug"
  6. "strconv"
  7. "strings"
  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. var _YMLampControllerMgrOnce sync.Once
  17. var _YMLampControllerMgrSingle *YMLampControllerMgr
  18. func GetYMLampControllerMgr() *YMLampControllerMgr {
  19. _YMLampControllerMgrOnce.Do(func() {
  20. ctx, cancel := context.WithCancel(context.Background())
  21. _YMLampControllerMgrSingle = &YMLampControllerMgr{
  22. queue: util.NewQueue(100),
  23. mapYMLampController: make(map[string]*YMLampController),
  24. ctx: ctx,
  25. cancel: cancel,
  26. }
  27. })
  28. return _YMLampControllerMgrSingle
  29. }
  30. type YMLampControllerMgr struct {
  31. queue *util.MlQueue
  32. mapYMLampController map[string]*YMLampController
  33. ctx context.Context
  34. cancel context.CancelFunc
  35. }
  36. func (o *YMLampControllerMgr) SubscribeTopics() {
  37. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_DATA), mqtt.AtLeastOnce, o.HandlerData)
  38. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_ALARM), mqtt.AtLeastOnce, o.HandlerData)
  39. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_SET_SWITCH_ACK), mqtt.AtLeastOnce, o.HandlerData)
  40. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_SET_ONOFFTIME_ACK), mqtt.AtLeastOnce, o.HandlerData)
  41. }
  42. func (o *YMLampControllerMgr) HandlerData(m mqtt.Message) {
  43. for {
  44. ok, cnt := o.queue.Put(&m)
  45. if ok {
  46. break
  47. } else {
  48. logrus.Errorf("YMLampControllerMgr.HandlerData:查询队列失败,队列消息数量:%d", cnt)
  49. runtime.Gosched()
  50. }
  51. }
  52. }
  53. func (o *YMLampControllerMgr) Stop() {
  54. o.cancel()
  55. }
  56. func (o *YMLampControllerMgr) Handler(args ...interface{}) interface{} {
  57. defer func() {
  58. if err := recover(); err != nil {
  59. time.Sleep(time.Second)
  60. gopool.Add(o.Handler, args)
  61. logrus.Errorf("YMLampControllerMgr.Handler发生异常:%v", err)
  62. logrus.Errorf("YMLampControllerMgr.Handler发生异常,堆栈信息:%s", string(debug.Stack()))
  63. }
  64. }()
  65. exit := false
  66. timer := time.NewTicker(1 * time.Minute)
  67. //每天15点半同步日出日落时间
  68. var SyncSunset = util.New(util.MlNow()).BeginningOfDay().Add(10*time.Hour + 30*time.Minute)
  69. for {
  70. select {
  71. case <-o.ctx.Done():
  72. logrus.Error("YMLampControllerMgr.HandleQueue即将退出,原因:", o.ctx.Err())
  73. exit = true
  74. case <-timer.C: //每隔1分钟执行一次
  75. //更新灯控状态,防止无数据状态不更新
  76. o.UpdateLampControllerState()
  77. //同步日出日落时间
  78. if util.MlNow().After(SyncSunset) {
  79. if err := o.SyncSunset(); err == nil {
  80. SyncSunset = SyncSunset.AddDate(0, 0, 1)
  81. }
  82. }
  83. default:
  84. if o.handleQueue() == 0 {
  85. if exit {
  86. return 0
  87. }
  88. time.Sleep(100 * time.Millisecond)
  89. }
  90. }
  91. }
  92. }
  93. func (o *YMLampControllerMgr) handleQueue() uint32 {
  94. msg, ok, quantity := o.queue.Get()
  95. if !ok {
  96. return quantity
  97. } else if quantity > 1000 {
  98. logrus.Warnf("YMLampControllerMgr.Handler:数据队列累积过多,请注意优化,当前队列条数:%d", quantity)
  99. }
  100. m, ok := msg.(*mqtt.Message)
  101. if !ok {
  102. return quantity
  103. }
  104. Tenant, _, DID, topic, err := ParseTopic(m.Topic())
  105. if err != nil {
  106. logrus.Errorf("YMLampControllerMgr.handleQueue:ParseTopic失败,topic=%s,err=%v", m.Topic(), err)
  107. return quantity
  108. }
  109. pymlc, ok := o.mapYMLampController[DID]
  110. if !ok {
  111. pymlc = &YMLampController{}
  112. pymlc.Set(Tenant, DID)
  113. o.mapYMLampController[DID] = pymlc
  114. }
  115. switch topic {
  116. case protocol.TP_YM_DATA:
  117. o.handleDATA(pymlc, m)
  118. case protocol.TP_YM_ALARM:
  119. o.handleALARM(pymlc, m)
  120. case protocol.TP_YM_SET_SWITCH_ACK, protocol.TP_YM_SET_ONOFFTIME_ACK:
  121. o.handleACK(m)
  122. default:
  123. logrus.Errorf("YMLampControllerMgr.handleQueue:未知topic=%s,全topic=%s", topic, m.Topic())
  124. }
  125. return quantity
  126. }
  127. func (o *YMLampControllerMgr) handleDATA(lp *YMLampController, m *mqtt.Message) {
  128. var obj protocol.Pack_CHZB_UploadData
  129. if err := obj.DeCode(m.PayloadString()); err != nil {
  130. logrus.Errorf("YMLampControllerMgr.handleDATA:DeCode失败,DID=%s,topic=%s,err=%v", lp.did, m.Topic(), err)
  131. return
  132. }
  133. t, err := util.MlParseTime(obj.Time)
  134. if err != nil {
  135. logrus.Errorf("时间[%s]解析错误:%s", obj.Time, err.Error())
  136. return
  137. }
  138. for _, v := range obj.Data.Data {
  139. lp.HandleData(obj.Gid, obj.Data.TID, t, v)
  140. }
  141. }
  142. func (o *YMLampControllerMgr) handleALARM(lp *YMLampController, m *mqtt.Message) {
  143. var obj protocol.Pack_CHZB_LampAlarm
  144. if err := obj.DeCode(m.PayloadString()); err != nil {
  145. logrus.Errorf("YMLampControllerMgr.handleALARM:DeCode失败,DID=%s,topic=%s,err=%v", lp.did, m.Topic(), err)
  146. return
  147. }
  148. lp.HandleAlarm(obj.Data)
  149. }
  150. func (o *YMLampControllerMgr) handleACK(m *mqtt.Message) {
  151. var obj protocol.Pack_Ack
  152. if err := obj.DeCode(m.PayloadString()); err != nil {
  153. logrus.Errorf("YMLampControllerMgr.handleACK:DeCode失败,topic=%s,err=%v", m.Topic(), err)
  154. return
  155. }
  156. oo := models.DeviceCmdRecord{ID: obj.Seq, State: 1, Resp: obj.Data.Error}
  157. if err := oo.Update(); err != nil {
  158. logrus.Errorf("收到设备[%s]的响应[seq:%d],主题:%s,但更新数据库失败[%s]", obj.Id, obj.Seq, m.Topic(), err.Error())
  159. }
  160. }
  161. // SyncSunset 统一更新裕明485灯控的日出日落时间
  162. func (o *YMLampControllerMgr) SyncSunset() error {
  163. arr, err := models.GetYm485Lampstrategy(nil)
  164. if err != nil {
  165. logrus.Errorf("从数据库读取设置为日出日落时间的485灯控发生错误:%s", err.Error())
  166. return err
  167. }
  168. if len(arr) == 0 {
  169. return nil
  170. }
  171. //分别计算日出日落
  172. mapTime := make(map[string]*protocol.CHZB_OnOffTime) //策略时间
  173. for _, v := range arr {
  174. if _, ok := mapTime[v.Strategy]; ok {
  175. //已计算日出日落时间的,不再重复计算
  176. continue
  177. }
  178. var oot protocol.CHZB_OnOffTime
  179. var ls []models.LampStrategy
  180. if err := json.UnmarshalFromString(v.TimeInfo, &ls); err == nil && len(ls) > 0 {
  181. oot.Brightness = uint8(ls[0].Brightness)
  182. }
  183. //计算时间
  184. if rise, set, err := util.SunriseSunsetForChina(v.Latitude, v.Longitude); err == nil {
  185. onHour, _ := strconv.Atoi(strings.Split(set, ":")[0])
  186. onMinute, _ := strconv.Atoi(strings.Split(set, ":")[1])
  187. offHour, _ := strconv.Atoi(strings.Split(rise, ":")[0])
  188. offMinute, _ := strconv.Atoi(strings.Split(rise, ":")[1])
  189. oot.OnHour = uint8(onHour)
  190. oot.OnMinite = uint8(onMinute)
  191. oot.OffHour = uint8(offHour)
  192. oot.OffMinite = uint8(offMinute)
  193. }
  194. mapTime[v.Strategy] = &oot
  195. }
  196. //发布mqtt消息
  197. for _, v := range arr {
  198. if oot, ok := mapTime[v.Strategy]; ok {
  199. var obj protocol.Pack_SetOnOffTime
  200. seq := GetNextSeq()
  201. if str, err := obj.EnCode(v.ID, v.GID, seq, nil, []protocol.CHZB_OnOffTime{*oot}); err == nil {
  202. topic := GetTopic(v.Tenant, protocol.DT_LAMPCONTROLLER, v.ID, protocol.TP_YM_SET_ONOFFTIME)
  203. err = GetMQTTMgr().Publish(topic, str, mqtt.AtLeastOnce)
  204. if err != nil {
  205. logrus.Errorf("SyncSunset:对灯控[%s]发布日出日落消息错误:%s", v.ID, err.Error())
  206. }
  207. var msg string
  208. if msg0, errmsg := json.MarshalIndent(obj, "", " "); errmsg == nil {
  209. msg = string(msg0)
  210. } else {
  211. msg = str
  212. }
  213. odb := models.DeviceCmdRecord{
  214. ID: seq,
  215. GID: v.GID,
  216. DID: v.ID,
  217. Topic: topic,
  218. Message: msg,
  219. State: 0,
  220. }
  221. if err := models.G_db.Create(&odb).Error; err != nil {
  222. logrus.Errorf("对灯控[%s]发布日出日落时间时指令入库错误:%s", v.ID, err.Error())
  223. } else {
  224. logrus.Errorf("对灯控[%s]发布日出日落时间时指令入库成功", v.ID)
  225. }
  226. }
  227. }
  228. }
  229. return nil
  230. }
  231. func (o *YMLampControllerMgr) UpdateLampControllerState() {
  232. for _, v := range o.mapYMLampController {
  233. v.UpdateState()
  234. }
  235. }