ymlampcontrollermgr.go 7.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245
  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.AtMostOnce, o.HandlerData)
  38. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_ALARM), mqtt.AtMostOnce, o.HandlerData)
  39. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_SET_SWITCH_ACK), mqtt.AtMostOnce, o.HandlerData)
  40. GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_LAMPCONTROLLER, protocol.TP_YM_SET_ONOFFTIME_ACK), mqtt.AtMostOnce, 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. return quantity
  107. }
  108. pymlc, ok := o.mapYMLampController[DID]
  109. if !ok {
  110. pymlc = &YMLampController{}
  111. pymlc.Set(Tenant, DID)
  112. o.mapYMLampController[DID] = pymlc
  113. }
  114. switch topic {
  115. case protocol.TP_YM_DATA:
  116. o.handleDATA(pymlc, m)
  117. case protocol.TP_YM_ALARM:
  118. o.handleALARM(pymlc, m)
  119. case protocol.TP_YM_SET_SWITCH_ACK, protocol.TP_YM_SET_ONOFFTIME_ACK:
  120. o.handleACK(m)
  121. }
  122. return quantity
  123. }
  124. func (o *YMLampControllerMgr) handleDATA(lp *YMLampController, m *mqtt.Message) {
  125. var obj protocol.Pack_CHZB_UploadData
  126. if err := obj.DeCode(m.PayloadString()); err != nil {
  127. return
  128. }
  129. t, err := util.MlParseTime(obj.Time)
  130. if err != nil {
  131. logrus.Errorf("时间[%s]解析错误:%s", obj.Time, err.Error())
  132. return
  133. }
  134. for _, v := range obj.Data.Data {
  135. lp.HandleData(obj.Gid, obj.Data.TID, t, v)
  136. }
  137. }
  138. func (o *YMLampControllerMgr) handleALARM(lp *YMLampController, m *mqtt.Message) {
  139. var obj protocol.Pack_CHZB_LampAlarm
  140. if err := obj.DeCode(m.PayloadString()); err != nil {
  141. return
  142. }
  143. lp.HandleAlarm(obj.Data)
  144. }
  145. func (o *YMLampControllerMgr) handleACK(m *mqtt.Message) {
  146. var obj protocol.Pack_Ack
  147. if err := obj.DeCode(m.PayloadString()); err != nil {
  148. return
  149. }
  150. oo := models.DeviceCmdRecord{ID: obj.Seq, State: 1, Resp: obj.Data.Error}
  151. if err := oo.Update(); err != nil {
  152. logrus.Errorf("收到设备[%s]的响应[seq:%d],主题:%s,但更新数据库失败[%s]", obj.Id, obj.Seq, m.Topic(), err.Error())
  153. }
  154. }
  155. // SyncSunset 统一更新裕明485灯控的日出日落时间
  156. func (o *YMLampControllerMgr) SyncSunset() error {
  157. arr, err := models.GetYm485Lampstrategy(nil)
  158. if err != nil {
  159. logrus.Errorf("从数据库读取设置为日出日落时间的485灯控发生错误:%s", err.Error())
  160. return err
  161. }
  162. if len(arr) == 0 {
  163. return nil
  164. }
  165. //分别计算日出日落
  166. mapTime := make(map[string]*protocol.CHZB_OnOffTime) //策略时间
  167. for _, v := range arr {
  168. if _, ok := mapTime[v.Strategy]; ok {
  169. //已计算日出日落时间的,不再重复计算
  170. continue
  171. }
  172. var oot protocol.CHZB_OnOffTime
  173. var ls []models.LampStrategy
  174. if err := json.UnmarshalFromString(v.TimeInfo, &ls); err == nil && len(ls) > 0 {
  175. oot.Brightness = uint8(ls[0].Brightness)
  176. }
  177. //计算时间
  178. if rise, set, err := util.SunriseSunsetForChina(v.Latitude, v.Longitude); err == nil {
  179. onHour, _ := strconv.Atoi(strings.Split(set, ":")[0])
  180. onMinute, _ := strconv.Atoi(strings.Split(set, ":")[1])
  181. offHour, _ := strconv.Atoi(strings.Split(rise, ":")[0])
  182. offMinute, _ := strconv.Atoi(strings.Split(rise, ":")[1])
  183. oot.OnHour = uint8(onHour)
  184. oot.OnMinite = uint8(onMinute)
  185. oot.OffHour = uint8(offHour)
  186. oot.OffMinite = uint8(offMinute)
  187. }
  188. mapTime[v.Strategy] = &oot
  189. }
  190. //发布mqtt消息
  191. for _, v := range arr {
  192. if oot, ok := mapTime[v.Strategy]; ok {
  193. var obj protocol.Pack_SetOnOffTime
  194. seq := GetNextSeq()
  195. if str, err := obj.EnCode(v.ID, v.GID, seq, nil, []protocol.CHZB_OnOffTime{*oot}); err == nil {
  196. topic := GetTopic(v.Tenant, protocol.DT_LAMPCONTROLLER, v.ID, protocol.TP_YM_SET_ONOFFTIME)
  197. err = GetMQTTMgr().Publish(topic, str, mqtt.AtLeastOnce)
  198. if err != nil {
  199. logrus.Errorf("SyncSunset:对灯控[%s]发布日出日落消息错误:%s", v.ID, err.Error())
  200. }
  201. var msg string
  202. if msg0, errmsg := json.MarshalIndent(obj, "", " "); errmsg == nil {
  203. msg = string(msg0)
  204. } else {
  205. msg = str
  206. }
  207. odb := models.DeviceCmdRecord{
  208. ID: seq,
  209. GID: v.GID,
  210. DID: v.ID,
  211. Topic: topic,
  212. Message: msg,
  213. State: 0,
  214. }
  215. if err := models.G_db.Create(&odb).Error; err != nil {
  216. logrus.Errorf("对灯控[%s]发布日出日落时间时指令入库错误:%s", v.ID, err.Error())
  217. } else {
  218. logrus.Errorf("对灯控[%s]发布日出日落时间时指令入库成功", v.ID)
  219. }
  220. }
  221. }
  222. }
  223. return nil
  224. }
  225. func (o *YMLampControllerMgr) UpdateLampControllerState() {
  226. for _, v := range o.mapYMLampController {
  227. v.UpdateState()
  228. }
  229. }