reader.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441
  1. package uhf
  2. import (
  3. "context"
  4. "encoding/binary"
  5. "encoding/hex"
  6. "errors"
  7. "fmt"
  8. "sync"
  9. "time"
  10. "wails-app/internal/dao"
  11. "wails-app/internal/global"
  12. common "wails-app/internal/model/common"
  13. incidentService "wails-app/internal/modules/incident/service"
  14. "wails-app/internal/service"
  15. "wails-app/internal/service/parking"
  16. )
  17. // Reader 通用读写接口(串口和TCP都实现它)
  18. type Reader interface {
  19. Connect() error
  20. Disconnect() error
  21. SendData([]byte) error
  22. ReadData() ([]byte, error)
  23. IsConnected() bool
  24. // ===================== 加入继电器接口 =====================
  25. CloseRelay1(validTime byte) error
  26. ReleaseRelay1() error
  27. CloseRelay2(validTime byte) error
  28. ReleaseRelay2() error
  29. }
  30. // ChannelEvent 通道事件(RFID/车牌识别后推送给前端)
  31. type ChannelEvent struct {
  32. ID uint64 `json:"id"`
  33. DeviceCode string `json:"device_code"`
  34. RFIDTag string `json:"rfid_tag"`
  35. PlateNumber string `json:"plate_number"`
  36. Direction string `json:"direction"`
  37. ChannelName string `json:"channel_name"`
  38. Timestamp int64 `json:"timestamp"`
  39. Status string `json:"status"`
  40. Message string `json:"message"`
  41. Fee float64 `json:"fee"`
  42. StayTime int64 `json:"stay_time"`
  43. }
  44. var (
  45. channelEventMu sync.Mutex
  46. ChannelEventQueue = make([]ChannelEvent, 0, 50)
  47. channelEventNextID uint64
  48. )
  49. // PushChannelEvent 追加通道事件(保留最近 50 条)
  50. func PushChannelEvent(e ChannelEvent) {
  51. channelEventMu.Lock()
  52. defer channelEventMu.Unlock()
  53. channelEventNextID++
  54. e.ID = channelEventNextID
  55. if e.Timestamp == 0 {
  56. e.Timestamp = time.Now().Unix()
  57. }
  58. if e.Status == "" {
  59. e.Status = "completed"
  60. }
  61. if len(ChannelEventQueue) >= 50 {
  62. ChannelEventQueue = ChannelEventQueue[1:]
  63. }
  64. ChannelEventQueue = append(ChannelEventQueue, e)
  65. }
  66. // GetChannelEvents 获取事件编号大于 afterID 的通道结果。
  67. func GetChannelEvents(afterID uint64) []ChannelEvent {
  68. channelEventMu.Lock()
  69. defer channelEventMu.Unlock()
  70. result := make([]ChannelEvent, 0)
  71. for _, e := range ChannelEventQueue {
  72. if e.ID > afterID {
  73. result = append(result, e)
  74. }
  75. }
  76. return result
  77. }
  78. // 上报数据模型
  79. type ReportData struct {
  80. DeviceCode string `json:"device_code"` // 设备编码
  81. Hex string `json:"hex"` // 原始报文
  82. Epcs []string `json:"epcs"` // 解析出的EPC列表
  83. RSSI int `json:"rssi"` // 信号强度
  84. Antenna int `json:"antenna"` // 天线号
  85. Timestamp int64 `json:"timestamp"` // 上报时间
  86. }
  87. // 全局设备管理器(管理所有已连接的设备)
  88. var DeviceManager = &Manager{
  89. devices: make(map[string]*DeviceHandler),
  90. }
  91. // init 注入统一道闸控制器(避免停车业务依赖具体设备协议)。
  92. func init() {
  93. parking.SetGateController(DeviceManager)
  94. }
  95. type Manager struct {
  96. mu sync.RWMutex
  97. devices map[string]*DeviceHandler
  98. }
  99. // DeviceHandler 每个设备的处理实例(包含读写器和协程)
  100. type DeviceHandler struct {
  101. reader Reader
  102. device *dao.UHFReader
  103. ctx context.Context
  104. cancel context.CancelFunc
  105. dataChan chan *ReportData
  106. isRunning bool
  107. }
  108. // Register 注册设备
  109. func (m *Manager) Register(code string, h *DeviceHandler) {
  110. m.mu.Lock()
  111. defer m.mu.Unlock()
  112. m.devices[code] = h
  113. }
  114. // Unregister 注销设备
  115. func (m *Manager) Unregister(code string) {
  116. m.mu.Lock()
  117. defer m.mu.Unlock()
  118. delete(m.devices, code)
  119. }
  120. // Get 获取设备
  121. func (m *Manager) Get(code string) (*DeviceHandler, bool) {
  122. m.mu.RLock()
  123. defer m.mu.RUnlock()
  124. h, ok := m.devices[code]
  125. return h, ok
  126. }
  127. // OpenGate 开闸(公开方法,供外部触发源调用,如摄像头识别、手动进出场)
  128. func (h *DeviceHandler) OpenGate(validTime byte) error {
  129. return h.reader.CloseRelay1(validTime)
  130. }
  131. // CloseGate 关闸。
  132. func (h *DeviceHandler) CloseGate(validTime byte) error {
  133. return h.reader.CloseRelay2(validTime)
  134. }
  135. // OpenGate 通过设备编码开闸,实现 parking.GateController。
  136. func (m *Manager) OpenGate(deviceCode string, validTime byte) error {
  137. h, ok := m.Get(deviceCode)
  138. if !ok {
  139. return fmt.Errorf("设备未连接: %s", deviceCode)
  140. }
  141. return h.OpenGate(validTime)
  142. }
  143. // CloseGate 通过设备编码关闸,实现 parking.GateController。
  144. func (m *Manager) CloseGate(deviceCode string, validTime byte) error {
  145. h, ok := m.Get(deviceCode)
  146. if !ok {
  147. return fmt.Errorf("设备未连接: %s", deviceCode)
  148. }
  149. return h.CloseGate(validTime)
  150. }
  151. // IsGateConnected 判断设备是否已注册到当前进程。
  152. func (m *Manager) IsGateConnected(deviceCode string) bool {
  153. _, ok := m.Get(deviceCode)
  154. return ok
  155. }
  156. // OpenGateByDeviceCode 通过设备编码直接开闸(便捷方法)
  157. func (m *Manager) OpenGateByDeviceCode(deviceCode string, validTime byte) error {
  158. return m.OpenGate(deviceCode, validTime)
  159. }
  160. // CloseGateByDeviceCode 通过设备编码直接关闸。
  161. func (m *Manager) CloseGateByDeviceCode(deviceCode string, validTime byte) error {
  162. return m.CloseGate(deviceCode, validTime)
  163. }
  164. // ===================== 复用原有CRC和解析函数(无需修改)=====================
  165. func uiCrc16Cal(pucY []byte, ucX uint8) uint16 {
  166. const PRESET_VALUE = 0xFFFF
  167. const POLYNOMIAL = 0x8408
  168. var uiCrcValue uint16 = PRESET_VALUE
  169. for ucI := uint8(0); ucI < ucX; ucI++ {
  170. uiCrcValue = uiCrcValue ^ uint16(pucY[ucI])
  171. for ucJ := uint8(0); ucJ < 8; ucJ++ {
  172. if uiCrcValue&0x0001 != 0 {
  173. uiCrcValue = (uiCrcValue >> 1) ^ POLYNOMIAL
  174. } else {
  175. uiCrcValue = uiCrcValue >> 1
  176. }
  177. }
  178. }
  179. return (uiCrcValue << 8) | (uiCrcValue >> 8)
  180. }
  181. // StartDeviceHandler 启动单个设备的读写协程
  182. func StartDeviceHandler(device *dao.UHFReader) error {
  183. // 1. 根据连接类型创建读写器
  184. var r Reader
  185. switch device.ConnectType {
  186. case "serial":
  187. r = NewSerialReader(device.COMPort, device.BaudRate)
  188. case "tcp":
  189. r = NewTCPReader(device.IPAddress, device.Port)
  190. default:
  191. return fmt.Errorf("不支持的连接类型: %s", device.ConnectType)
  192. }
  193. // 2. 连接设备
  194. if err := r.Connect(); err != nil {
  195. return err
  196. }
  197. // 3. 创建上下文和handler
  198. ctx, cancel := context.WithCancel(context.Background())
  199. h := &DeviceHandler{
  200. reader: r,
  201. device: device,
  202. ctx: ctx,
  203. cancel: cancel,
  204. dataChan: make(chan *ReportData, 200), // 适当加大
  205. isRunning: true,
  206. }
  207. // 4. 注册到管理器
  208. DeviceManager.Register(device.DeviceCode, h)
  209. // 5. 启动读写协程
  210. go func() {
  211. defer func() {
  212. h.isRunning = false
  213. _ = r.Disconnect()
  214. close(h.dataChan)
  215. DeviceManager.Unregister(device.DeviceCode)
  216. global.GVA_LOG.Info("设备已停止 code:" + device.DeviceCode)
  217. }()
  218. for {
  219. select {
  220. case <-ctx.Done():
  221. return
  222. default:
  223. buf, err := r.ReadData()
  224. if err != nil {
  225. global.GVA_DB.Model(device).Update("status", "offline")
  226. // 离线异常埋点(best-effort,同设备未关闭不重复生成)
  227. incidentService.NewIncidentService().RecordIncident(incidentService.RecordIncidentRequest{
  228. Category: incidentService.CategoryDeviceOffline,
  229. Source: incidentService.SourceDevice,
  230. ParkingLotID: device.ParkingLotID,
  231. DeviceCode: device.DeviceCode,
  232. Description: "UHF 读卡器离线",
  233. Detail: err.Error(),
  234. })
  235. _ = r.Disconnect() // 先断开
  236. time.Sleep(2 * time.Second)
  237. if err := r.Connect(); err != nil {
  238. continue
  239. }
  240. global.GVA_DB.Model(device).Update("status", "online")
  241. // 设备恢复在线:自动关闭未处理的离线异常
  242. incidentService.NewIncidentService().ResolveDeviceOffline(device.DeviceCode)
  243. continue
  244. }
  245. if len(buf) == 0 {
  246. continue
  247. }
  248. // 解析上报数据
  249. report, err := parseReportData(device.DeviceCode, buf)
  250. if err != nil {
  251. continue
  252. }
  253. // 推送(不丢死,也不阻塞)
  254. select {
  255. case h.dataChan <- report:
  256. default:
  257. global.GVA_LOG.Warn("通道已满,丢弃数据 device" + device.DeviceCode)
  258. }
  259. }
  260. }
  261. }()
  262. // 6. 启动业务处理协程
  263. go func() {
  264. for {
  265. select {
  266. case <-ctx.Done():
  267. return
  268. case report, ok := <-h.dataChan:
  269. if !ok {
  270. return
  271. }
  272. handleReportData(report)
  273. }
  274. }
  275. }()
  276. global.GVA_DB.Model(device).Update("status", "online")
  277. // 启动时兜底:关闭历史遗留的同设备离线异常
  278. incidentService.NewIncidentService().ResolveDeviceOffline(device.DeviceCode)
  279. return nil
  280. }
  281. // parseReportData 解析上报数据
  282. func parseReportData(deviceCode string, buf []byte) (*ReportData, error) {
  283. if len(buf) < 25 || buf[0] != 0xCF {
  284. return nil, errors.New("无效帧")
  285. }
  286. dataWithoutCRC := buf[:len(buf)-2]
  287. recvCRC := binary.LittleEndian.Uint16(buf[len(buf)-2:])
  288. calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC)))
  289. if recvCRC != calcCRC {
  290. return nil, errors.New("CRC校验失败")
  291. }
  292. rssi := int(buf[6])
  293. antenna := int(buf[22])
  294. epc := hex.EncodeToString(buf[11:23])
  295. fmt.Printf("%s", hex.EncodeToString(buf))
  296. return &ReportData{
  297. DeviceCode: deviceCode,
  298. Hex: hex.EncodeToString(buf),
  299. Epcs: []string{epc},
  300. RSSI: rssi,
  301. Antenna: antenna,
  302. Timestamp: time.Now().Unix(),
  303. }, nil
  304. }
  305. // handleReportData 处理UHF读取到的标签(精简版:只做解析+防抖,业务逻辑委托给 PassageService)
  306. func handleReportData(report *ReportData) {
  307. if len(report.Epcs) == 0 {
  308. return
  309. }
  310. epc := report.Epcs[0]
  311. fmt.Println("<UNK>", epc)
  312. // 委托给统一的进出场服务(所有触发方式共用同一入口)
  313. result, err := service.ServiceGroupApp.ParkingServiceGroup.PassageService.HandlePassage(common.PassageRequest{
  314. RFIDTag: epc,
  315. DeviceCode: report.DeviceCode,
  316. TriggerSource: "rfid",
  317. })
  318. if err != nil {
  319. fmt.Printf("⚠️ 进出场处理失败 | 设备:%s EPC:%s 错误:%v\n", report.DeviceCode, epc, err)
  320. if result != nil {
  321. PushChannelEvent(ChannelEvent{
  322. DeviceCode: report.DeviceCode, RFIDTag: epc, PlateNumber: result.PlateNumber,
  323. Direction: result.Direction, Timestamp: time.Now().Unix(), Status: result.GateStatus,
  324. Message: result.Message, Fee: result.Fee, StayTime: result.StayTime,
  325. })
  326. }
  327. return
  328. }
  329. fmt.Printf("✅ 进出场成功 | 方向:%s 车牌:%s 费用:%.2f\n", result.Direction, result.PlateNumber, result.Fee)
  330. // 推送通道事件给前端
  331. var device dao.UHFReader
  332. if err := global.GVA_DB.Preload("Channel").First(&device, "device_code = ?", report.DeviceCode).Error; err == nil {
  333. PushChannelEvent(ChannelEvent{
  334. DeviceCode: report.DeviceCode,
  335. RFIDTag: epc,
  336. PlateNumber: result.PlateNumber,
  337. Direction: result.Direction,
  338. ChannelName: device.Channel.ChannelName,
  339. Timestamp: time.Now().Unix(),
  340. Status: "completed",
  341. Message: result.Message,
  342. Fee: result.Fee,
  343. StayTime: result.StayTime,
  344. })
  345. }
  346. }
  347. // ==============================
  348. // 继电器控制(完整支持 Relay1 & Relay2)
  349. // ==============================
  350. const (
  351. RELAY_OP_RELEASE = 0x01
  352. RELAY_OP_CLOSE = 0x02
  353. )
  354. // CloseRelay1 开闸
  355. func (s *SerialReader) CloseRelay1(validTime byte) error {
  356. frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
  357. _, err := s.SendAndRecv(frame)
  358. return err
  359. }
  360. // ReleaseRelay1
  361. func (s *SerialReader) ReleaseRelay1() error {
  362. frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
  363. _, err := s.SendAndRecv(frame)
  364. return err
  365. }
  366. // CloseRelay2 关闸
  367. func (s *SerialReader) CloseRelay2(validTime byte) error {
  368. frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
  369. _, err := s.SendAndRecv(frame)
  370. return err
  371. }
  372. // ReleaseRelay2
  373. func (s *SerialReader) ReleaseRelay2() error {
  374. frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
  375. _, err := s.SendAndRecv(frame)
  376. return err
  377. }
  378. // buildRelayFrame 构建指令
  379. func buildRelayFrame(cmd uint16, relayNum byte, option byte, validTime byte) []byte {
  380. frame := []byte{
  381. 0xCF, 0xFF,
  382. byte(cmd >> 8), byte(cmd & 0xFF),
  383. 0x03, // len
  384. relayNum, // 1=继电器1 2=继电器2
  385. option,
  386. validTime,
  387. }
  388. crc := uiCrc16Cal(frame, uint8(len(frame)))
  389. frame = append(frame, byte(crc&0xFF), byte(crc>>8))
  390. return frame
  391. }