websvr.go 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257
  1. package main
  2. import (
  3. "context"
  4. "errors"
  5. "net/http"
  6. "runtime/debug"
  7. "strconv"
  8. "sync"
  9. "time"
  10. "github.com/labstack/echo/v4"
  11. "github.com/labstack/echo/v4/middleware"
  12. "github.com/sirupsen/logrus"
  13. "lc/common/mqtt"
  14. "lc/common/protocol"
  15. "lc/common/util"
  16. )
  17. var _onceWebSvr sync.Once
  18. var _singleWebSvr *WebSvr
  19. func GetWebSvr() *WebSvr {
  20. _onceWebSvr.Do(func() {
  21. _singleWebSvr = NewWebSvr(itsConfig.IPAddr)
  22. })
  23. return _singleWebSvr
  24. }
  25. type H map[string]interface{}
  26. type WebSvr struct {
  27. echo *echo.Echo
  28. IPAddr string
  29. queue *util.MlQueue
  30. ctx context.Context
  31. cancel context.CancelFunc
  32. muEnvdata sync.Mutex
  33. Envdata *EnvData //环境数据
  34. muVSpeed sync.Mutex
  35. VSpeed map[string]*protocol.VehicleSpeed //CamID+Direction->速度
  36. mapTopicHandle map[string]func(m mqtt.Message)
  37. }
  38. func NewWebSvr(addr string) *WebSvr {
  39. ctx, cancel := context.WithCancel(context.Background())
  40. obj := WebSvr{
  41. echo: echo.New(),
  42. IPAddr: addr,
  43. queue: util.NewQueue(200),
  44. ctx: ctx,
  45. cancel: cancel,
  46. VSpeed: make(map[string]*protocol.VehicleSpeed),
  47. mapTopicHandle: make(map[string]func(m mqtt.Message)),
  48. }
  49. //obj.echo.Use(middleware.Logger())
  50. obj.echo.Debug = true
  51. obj.echo.HideBanner = true
  52. obj.echo.Use(middleware.Recover())
  53. obj.echo.Use(middleware.CORS())
  54. obj.echo.HTTPErrorHandler = obj.HTTPErrorHandler
  55. obj.echo.Static("/", "public")
  56. obj.echo.GET("/getdata", obj.GetData)
  57. return &obj
  58. }
  59. func (o *WebSvr) MQTTSubscribe() {
  60. o.mapTopicHandle[GetTopic(protocol.DT_ENVIRONMENT, itsConfig.EnvID, protocol.TP_MODBUS_DATA)] = o.HandleTpModbusData
  61. o.mapTopicHandle[GetTopic(protocol.DT_ITS, itsConfig.RemoteSub, protocol.TP_ITS_VEHICLESPEED)] = o.HandleTpItsVehiclespeed
  62. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_ENVIRONMENT, itsConfig.EnvID, protocol.TP_MODBUS_DATA), mqtt.ExactlyOnce, o.HandleCache, ToEdge)
  63. if len(itsConfig.RemoteSub) > 0 {
  64. GetMQTTMgr().Subscribe(GetTopic(protocol.DT_ITS, itsConfig.RemoteSub, protocol.TP_ITS_VEHICLESPEED), mqtt.ExactlyOnce, o.HandleCache, ToCloud)
  65. }
  66. }
  67. func (o *WebSvr) HandleCache(m mqtt.Message) {
  68. o.queue.Put(m)
  69. }
  70. func (o *WebSvr) HandleTpItsVehiclespeed(m mqtt.Message) {
  71. var obj protocol.Pack_VehicleSpeed
  72. if err := obj.DeCode(m.PayloadString()); err != nil {
  73. logrus.Errorf("解码错误:%s", err.Error())
  74. return
  75. }
  76. key := obj.Id
  77. value := &protocol.VehicleSpeed{
  78. Plate: obj.Data.Plate,
  79. Time: obj.Data.Time,
  80. Type: obj.Data.Type,
  81. Speed: obj.Data.Speed,
  82. }
  83. o.muVSpeed.Lock()
  84. o.VSpeed[key] = value
  85. o.muVSpeed.Unlock()
  86. }
  87. func (o *WebSvr) HandleTpModbusData(m mqtt.Message) {
  88. var obj protocol.Pack_UploadData
  89. if err := obj.DeCode(m.PayloadString()); err != nil {
  90. logrus.Errorf("解码错误:%s", err.Error())
  91. return
  92. }
  93. if obj.Data.Tid != 1 {
  94. return
  95. }
  96. //key = 1,噪声;2,PM2.5;3,PM10;4,温度;5,湿度;6,大气压;7,风速;8,风向
  97. var env EnvData
  98. env.Time = util.MlNow()
  99. pm25, err := GetValue(obj.Data.Data, 2)
  100. if err == nil {
  101. env.PM25 = strconv.FormatFloat(pm25, 'f', 0, 64)
  102. }
  103. pm10, err := GetValue(obj.Data.Data, 3)
  104. if err == nil {
  105. env.PM10 = strconv.FormatFloat(pm10, 'f', 0, 64)
  106. }
  107. temp, err := GetValue(obj.Data.Data, 4)
  108. if err == nil {
  109. env.Wendu = strconv.FormatFloat(temp, 'f', 0, 64)
  110. }
  111. sd, err := GetValue(obj.Data.Data, 5)
  112. if err == nil {
  113. env.Shidu = strconv.FormatFloat(sd, 'f', 0, 64)
  114. }
  115. aqi, degree := util.GetAqiAndDegree(pm25, pm10)
  116. env.Aqi = strconv.FormatFloat(aqi, 'f', 0, 64)
  117. env.Zhishu = degree
  118. //风速,风力
  119. direction, err := GetValue(obj.Data.Data, 8)
  120. if err == nil {
  121. env.Fengxiang = util.GetWindDirection(direction).Namezh
  122. }
  123. speed, err := GetValue(obj.Data.Data, 7)
  124. if err == nil {
  125. x, _ := util.GetWindGrade(speed)
  126. env.Fengji = x
  127. }
  128. o.muEnvdata.Lock()
  129. o.Envdata = &env
  130. o.muEnvdata.Unlock()
  131. }
  132. func (o *WebSvr) Handle(args ...interface{}) interface{} {
  133. defer func() {
  134. if err := recover(); err != nil {
  135. logrus.Errorf("WebSvr.StartSvr发生异常:%v", err)
  136. logrus.Errorf("WebSvr.StartSvr发生异常,堆栈信息:%s", string(debug.Stack()))
  137. time.Sleep(time.Second)
  138. gopool.Add(o.StartSvr, args)
  139. }
  140. }()
  141. exit := false
  142. for {
  143. select {
  144. case <-o.ctx.Done():
  145. logrus.Errorf("WebSvr.HandleData退出,原因:%v", o.ctx.Err())
  146. exit = true
  147. default:
  148. //从队列钟获取指令执行
  149. if m, ok, _ := o.queue.Get(); ok {
  150. if mm, ok := m.(mqtt.Message); ok {
  151. if fn, ok := o.mapTopicHandle[mm.Topic()]; ok {
  152. fn(mm)
  153. }
  154. }
  155. } else {
  156. if exit { //退出前全部恢复时控模式
  157. return 0
  158. }
  159. time.Sleep(300 * time.Millisecond)
  160. }
  161. }
  162. }
  163. }
  164. func (o *WebSvr) Update(data *tagVehicleInfo) {
  165. key := data.CamID
  166. //不用车牌时间,防止抓拍单元时间和网关时间相差太大
  167. value := &protocol.VehicleSpeed{Plate: data.VehiclePlate, Time: protocol.BJTime(util.MlNow()), Type: data.VehicleType, Speed: data.VehicleSpeed}
  168. o.muVSpeed.Lock()
  169. o.VSpeed[key] = value
  170. o.muVSpeed.Unlock()
  171. }
  172. func (o *WebSvr) StartSvr(args ...interface{}) interface{} {
  173. defer func() {
  174. if err := recover(); err != nil {
  175. logrus.Errorf("WebSvr.StartSvr发生异常:%v", err)
  176. logrus.Errorf("WebSvr.StartSvr发生异常,堆栈信息:%s", string(debug.Stack()))
  177. time.Sleep(time.Second)
  178. gopool.Add(o.StartSvr, args)
  179. }
  180. }()
  181. return o.echo.Start(o.IPAddr)
  182. }
  183. func (o *WebSvr) HTTPErrorHandler(err error, c echo.Context) {
  184. c.JSON(http.StatusInternalServerError, nil)
  185. }
  186. func (o *WebSvr) GetData(c echo.Context) error {
  187. obj := GetDataPool().GetData()
  188. obj.Reset()
  189. defer GetDataPool().Release(obj)
  190. key := c.QueryParam("code")
  191. o.muVSpeed.Lock()
  192. if v, ok := o.VSpeed[key]; ok {
  193. diffSec := util.MlNow().Sub(time.Time(v.Time)).Seconds()
  194. if diffSec > float64(itsConfig.Duration) {
  195. obj.Chesu = 0
  196. delete(o.VSpeed, key)
  197. } else {
  198. obj.Chesu = v.Speed
  199. }
  200. }
  201. o.muVSpeed.Unlock()
  202. o.muEnvdata.Lock()
  203. if o.Envdata != nil {
  204. //if util.MlNow().Sub(o.Envdata.Time).Seconds() < 24*3600 {
  205. obj.Fengji = o.Envdata.Fengji
  206. obj.Tianqi = o.Envdata.Tianqi
  207. obj.Aqi = o.Envdata.Aqi
  208. obj.PM25 = o.Envdata.PM25
  209. obj.PM10 = o.Envdata.PM10
  210. obj.Wendu = o.Envdata.Wendu
  211. obj.Zhishu = o.Envdata.Zhishu
  212. obj.Shidu = o.Envdata.Shidu
  213. obj.Fengxiang = o.Envdata.Fengxiang
  214. obj.Fengji = o.Envdata.Fengji
  215. //} else {
  216. // o.Envdata = nil
  217. //}
  218. }
  219. o.muEnvdata.Unlock()
  220. obj.Chesufazhi = int(itsConfig.SuggestSpeed)
  221. return c.JSON(http.StatusOK, H{"data": 0})
  222. }
  223. func GetValue(m map[uint16]float64, key uint16) (float64, error) {
  224. if v, ok := m[key]; ok {
  225. return v, nil
  226. }
  227. return 0.0, errors.New("未找到")
  228. }