reader.go 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620
  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. "go.uber.org/zap"
  17. )
  18. // Reader 通用读写接口(串口和TCP都实现它)
  19. type Reader interface {
  20. Connect() error
  21. Disconnect() error
  22. SendData([]byte) error
  23. ReadData() ([]byte, error)
  24. IsConnected() bool
  25. // ===================== 加入继电器接口 =====================
  26. CloseRelay1(validTime byte) error
  27. ReleaseRelay1() error
  28. CloseRelay2(validTime byte) error
  29. ReleaseRelay2() error
  30. }
  31. // ChannelEvent 通道事件(RFID/车牌识别后推送给前端)
  32. type ChannelEvent struct {
  33. ID uint64 `json:"id"`
  34. DeviceCode string `json:"device_code"`
  35. RFIDTag string `json:"rfid_tag"`
  36. PlateNumber string `json:"plate_number"`
  37. Direction string `json:"direction"`
  38. ChannelName string `json:"channel_name"`
  39. Timestamp int64 `json:"timestamp"`
  40. Status string `json:"status"`
  41. Message string `json:"message"`
  42. Fee float64 `json:"fee"`
  43. StayTime int64 `json:"stay_time"`
  44. }
  45. var (
  46. channelEventMu sync.Mutex
  47. ChannelEventQueue = make([]ChannelEvent, 0, 50)
  48. channelEventNextID uint64
  49. )
  50. // PushChannelEvent 追加通道事件(保留最近 50 条)
  51. func PushChannelEvent(e ChannelEvent) {
  52. channelEventMu.Lock()
  53. defer channelEventMu.Unlock()
  54. channelEventNextID++
  55. e.ID = channelEventNextID
  56. if e.Timestamp == 0 {
  57. e.Timestamp = time.Now().Unix()
  58. }
  59. if e.Status == "" {
  60. e.Status = "completed"
  61. }
  62. if len(ChannelEventQueue) >= 50 {
  63. ChannelEventQueue = ChannelEventQueue[1:]
  64. }
  65. ChannelEventQueue = append(ChannelEventQueue, e)
  66. }
  67. // GetChannelEvents 获取事件编号大于 afterID 的通道结果。
  68. func GetChannelEvents(afterID uint64) []ChannelEvent {
  69. channelEventMu.Lock()
  70. defer channelEventMu.Unlock()
  71. result := make([]ChannelEvent, 0)
  72. for _, e := range ChannelEventQueue {
  73. if e.ID > afterID {
  74. result = append(result, e)
  75. }
  76. }
  77. return result
  78. }
  79. // 上报数据模型
  80. type ReportData struct {
  81. DeviceCode string `json:"device_code"` // 设备编码
  82. Hex string `json:"hex"` // 原始报文
  83. Epcs []string `json:"epcs"` // 解析出的EPC列表
  84. RSSI int `json:"rssi"` // 信号强度
  85. Antenna int `json:"antenna"` // 天线号
  86. Timestamp int64 `json:"timestamp"` // 上报时间
  87. }
  88. // 全局设备管理器(管理所有已连接的设备)
  89. var DeviceManager = &Manager{
  90. devices: make(map[string]*DeviceHandler),
  91. }
  92. // init 注入统一道闸控制器(避免停车业务依赖具体设备协议)。
  93. func init() {
  94. parking.SetGateController(DeviceManager)
  95. }
  96. type Manager struct {
  97. mu sync.RWMutex
  98. devices map[string]*DeviceHandler
  99. }
  100. // DeviceHandler 每个设备的处理实例(包含读写器和协程)
  101. type DeviceHandler struct {
  102. reader Reader
  103. device *dao.UHFReader
  104. ctx context.Context
  105. cancel context.CancelFunc
  106. dataChan chan *ReportData
  107. ioMu sync.Mutex
  108. stopOnce sync.Once
  109. done chan struct{}
  110. isRunning bool
  111. }
  112. // Register 注册设备
  113. func (m *Manager) Register(code string, h *DeviceHandler) {
  114. m.mu.Lock()
  115. defer m.mu.Unlock()
  116. m.devices[code] = h
  117. }
  118. // Unregister 注销设备
  119. func (m *Manager) Unregister(code string) {
  120. m.mu.Lock()
  121. defer m.mu.Unlock()
  122. delete(m.devices, code)
  123. }
  124. // Get 获取设备
  125. func (m *Manager) Get(code string) (*DeviceHandler, bool) {
  126. m.mu.RLock()
  127. defer m.mu.RUnlock()
  128. h, ok := m.devices[code]
  129. return h, ok
  130. }
  131. // StopDeviceHandler 停止设备读循环、关闭底层连接并等待退出。
  132. func StopDeviceHandler(code string) error {
  133. h, ok := DeviceManager.Get(code)
  134. if !ok {
  135. return nil
  136. }
  137. h.stopOnce.Do(func() {
  138. if h.cancel != nil {
  139. h.cancel()
  140. }
  141. if h.reader != nil {
  142. _ = h.reader.Disconnect()
  143. }
  144. })
  145. if h.done == nil {
  146. DeviceManager.Unregister(code)
  147. return nil
  148. }
  149. select {
  150. case <-h.done:
  151. case <-time.After(5 * time.Second):
  152. return fmt.Errorf("设备停止超时: %s", code)
  153. }
  154. return nil
  155. }
  156. // OpenGate 开闸(公开方法,供外部触发源调用,如摄像头识别、手动进出场)
  157. func (h *DeviceHandler) OpenGate(validTime byte) error {
  158. h.ioMu.Lock()
  159. defer h.ioMu.Unlock()
  160. return h.reader.CloseRelay1(validTime)
  161. }
  162. // CloseGate 关闸。
  163. func (h *DeviceHandler) CloseGate(validTime byte) error {
  164. h.ioMu.Lock()
  165. defer h.ioMu.Unlock()
  166. return h.reader.CloseRelay2(validTime)
  167. }
  168. // OpenGate 通过设备编码开闸,实现 parking.GateController。
  169. func (m *Manager) OpenGate(deviceCode string, validTime byte) error {
  170. h, ok := m.Get(deviceCode)
  171. if !ok {
  172. return fmt.Errorf("设备未连接: %s", deviceCode)
  173. }
  174. return h.OpenGate(validTime)
  175. }
  176. // CloseGate 通过设备编码关闸,实现 parking.GateController。
  177. func (m *Manager) CloseGate(deviceCode string, validTime byte) error {
  178. h, ok := m.Get(deviceCode)
  179. if !ok {
  180. return fmt.Errorf("设备未连接: %s", deviceCode)
  181. }
  182. return h.CloseGate(validTime)
  183. }
  184. // IsGateConnected 判断设备是否已注册到当前进程。
  185. func (m *Manager) IsGateConnected(deviceCode string) bool {
  186. h, ok := m.Get(deviceCode)
  187. return ok && h.reader != nil && h.reader.IsConnected()
  188. }
  189. // OpenGateByDeviceCode 通过设备编码直接开闸(便捷方法)
  190. func (m *Manager) OpenGateByDeviceCode(deviceCode string, validTime byte) error {
  191. return m.OpenGate(deviceCode, validTime)
  192. }
  193. // CloseGateByDeviceCode 通过设备编码直接关闸。
  194. func (m *Manager) CloseGateByDeviceCode(deviceCode string, validTime byte) error {
  195. return m.CloseGate(deviceCode, validTime)
  196. }
  197. // ===================== 复用原有CRC和解析函数(无需修改)=====================
  198. func uiCrc16Cal(pucY []byte, ucX uint8) uint16 {
  199. const PRESET_VALUE = 0xFFFF
  200. const POLYNOMIAL = 0x8408
  201. var uiCrcValue uint16 = PRESET_VALUE
  202. for ucI := uint8(0); ucI < ucX; ucI++ {
  203. uiCrcValue = uiCrcValue ^ uint16(pucY[ucI])
  204. for ucJ := uint8(0); ucJ < 8; ucJ++ {
  205. if uiCrcValue&0x0001 != 0 {
  206. uiCrcValue = (uiCrcValue >> 1) ^ POLYNOMIAL
  207. } else {
  208. uiCrcValue = uiCrcValue >> 1
  209. }
  210. }
  211. }
  212. return (uiCrcValue << 8) | (uiCrcValue >> 8)
  213. }
  214. // StartDeviceHandler 启动单个设备的读写协程
  215. func StartDeviceHandler(device *dao.UHFReader) error {
  216. if device == nil || device.DeviceCode == "" {
  217. return errors.New("设备编码不能为空")
  218. }
  219. if existing, ok := DeviceManager.Get(device.DeviceCode); ok {
  220. if existing.reader != nil && existing.reader.IsConnected() {
  221. return fmt.Errorf("设备已连接: %s", device.DeviceCode)
  222. }
  223. if err := StopDeviceHandler(device.DeviceCode); err != nil {
  224. return err
  225. }
  226. }
  227. // MQTT 设备由设备主动连接 broker,系统不主动建立连接。
  228. if device.ConnectType == dao.ConnectTypeMQTT {
  229. return nil
  230. }
  231. if global.GVA_LOG != nil {
  232. global.GVA_LOG.Info("启动 UHF/RFID 设备连接",
  233. zap.String("device_code", device.DeviceCode),
  234. zap.String("device_name", device.DeviceName),
  235. zap.String("connect_type", device.ConnectType),
  236. zap.String("tcp_address", readerTCPAddress(device)),
  237. zap.Uint("channel_id", device.ChannelID),
  238. )
  239. }
  240. // 1. 根据连接类型创建读写器
  241. var r Reader
  242. switch device.ConnectType {
  243. case "serial":
  244. r = NewSerialReader(device.COMPort, device.BaudRate)
  245. case "tcp":
  246. r = NewTCPReader(device.IPAddress, device.Port)
  247. default:
  248. return fmt.Errorf("不支持的连接类型: %s", device.ConnectType)
  249. }
  250. // 2. 连接设备
  251. if err := r.Connect(); err != nil {
  252. if global.GVA_LOG != nil {
  253. global.GVA_LOG.Error("UHF/RFID 设备连接失败",
  254. zap.String("device_code", device.DeviceCode),
  255. zap.String("connect_type", device.ConnectType),
  256. zap.String("tcp_address", readerTCPAddress(device)),
  257. zap.Error(err),
  258. )
  259. }
  260. return err
  261. }
  262. if global.GVA_LOG != nil {
  263. global.GVA_LOG.Info("UHF/RFID 设备连接成功",
  264. zap.String("device_code", device.DeviceCode),
  265. zap.String("connect_type", device.ConnectType),
  266. zap.String("tcp_address", readerTCPAddress(device)),
  267. )
  268. }
  269. // 3. 创建上下文和handler
  270. ctx, cancel := context.WithCancel(context.Background())
  271. h := &DeviceHandler{
  272. reader: r,
  273. device: device,
  274. ctx: ctx,
  275. cancel: cancel,
  276. dataChan: make(chan *ReportData, 200),
  277. done: make(chan struct{}),
  278. isRunning: true,
  279. }
  280. // 4. 注册到管理器
  281. DeviceManager.Register(device.DeviceCode, h)
  282. // 5. 启动读写协程
  283. go func() {
  284. defer func() {
  285. h.stopOnce.Do(func() { h.cancel() })
  286. h.isRunning = false
  287. _ = r.Disconnect()
  288. close(h.dataChan)
  289. DeviceManager.Unregister(device.DeviceCode)
  290. close(h.done)
  291. if global.GVA_LOG != nil {
  292. global.GVA_LOG.Info("设备已停止: " + device.DeviceCode)
  293. }
  294. }()
  295. for {
  296. select {
  297. case <-ctx.Done():
  298. return
  299. default:
  300. h.ioMu.Lock()
  301. buf, err := r.ReadData()
  302. h.ioMu.Unlock()
  303. if errors.Is(err, ErrReadTimeout) || errors.Is(err, ErrIncompleteFrame) {
  304. continue
  305. }
  306. if err != nil {
  307. if global.GVA_LOG != nil {
  308. global.GVA_LOG.Warn("UHF/RFID 读取失败,准备重连",
  309. zap.String("device_code", device.DeviceCode),
  310. zap.String("connect_type", device.ConnectType),
  311. zap.String("tcp_address", readerTCPAddress(device)),
  312. zap.Error(err),
  313. )
  314. }
  315. if global.GVA_DB != nil {
  316. global.GVA_DB.Model(device).Update("status", "offline")
  317. }
  318. // 离线异常埋点(best-effort,同设备未关闭不重复生成)
  319. incidentService.NewIncidentService().RecordIncident(incidentService.RecordIncidentRequest{
  320. Category: incidentService.CategoryDeviceOffline,
  321. Source: incidentService.SourceDevice,
  322. ParkingLotID: device.ParkingLotID,
  323. DeviceCode: device.DeviceCode,
  324. Description: "UHF 读卡器离线",
  325. Detail: err.Error(),
  326. })
  327. _ = r.Disconnect() // 先断开
  328. if !waitDeviceRetry(ctx, 2*time.Second) {
  329. return
  330. }
  331. if err := r.Connect(); err != nil {
  332. if global.GVA_LOG != nil {
  333. global.GVA_LOG.Warn("UHF/RFID 设备重连失败",
  334. zap.String("device_code", device.DeviceCode),
  335. zap.String("tcp_address", readerTCPAddress(device)),
  336. zap.Error(err),
  337. )
  338. }
  339. continue
  340. }
  341. if global.GVA_LOG != nil {
  342. global.GVA_LOG.Info("UHF/RFID 设备重连成功",
  343. zap.String("device_code", device.DeviceCode),
  344. zap.String("tcp_address", readerTCPAddress(device)),
  345. )
  346. }
  347. if global.GVA_DB != nil {
  348. global.GVA_DB.Model(device).Updates(map[string]interface{}{"status": "online", "last_online_time": time.Now()})
  349. }
  350. // 设备恢复在线:自动关闭未处理的离线异常
  351. incidentService.NewIncidentService().ResolveDeviceOffline(device.DeviceCode)
  352. continue
  353. }
  354. if len(buf) == 0 {
  355. continue
  356. }
  357. if global.GVA_LOG != nil {
  358. global.GVA_LOG.Info("收到 UHF/RFID 原始帧",
  359. zap.String("device_code", device.DeviceCode),
  360. zap.Int("packet_bytes", len(buf)),
  361. zap.String("packet_hex", hex.EncodeToString(buf)),
  362. )
  363. }
  364. // 解析上报数据
  365. report, err := parseReportData(device.DeviceCode, buf)
  366. if err != nil {
  367. if global.GVA_LOG != nil {
  368. global.GVA_LOG.Warn("UHF/RFID 报文解析失败",
  369. zap.String("device_code", device.DeviceCode),
  370. zap.String("packet_hex", hex.EncodeToString(buf)),
  371. zap.Error(err),
  372. )
  373. }
  374. continue
  375. }
  376. if global.GVA_LOG != nil {
  377. global.GVA_LOG.Info("UHF/RFID 标签解析成功",
  378. zap.String("device_code", report.DeviceCode),
  379. zap.Strings("epcs", report.Epcs),
  380. zap.Int("rssi", report.RSSI),
  381. zap.Int("antenna", report.Antenna),
  382. )
  383. }
  384. // 推送(不丢死,也不阻塞)
  385. select {
  386. case h.dataChan <- report:
  387. default:
  388. if global.GVA_LOG != nil {
  389. global.GVA_LOG.Warn("通道已满,丢弃数据: " + device.DeviceCode)
  390. }
  391. }
  392. }
  393. }
  394. }()
  395. // 6. 启动业务处理协程
  396. go func() {
  397. for {
  398. select {
  399. case <-ctx.Done():
  400. return
  401. case report, ok := <-h.dataChan:
  402. if !ok {
  403. return
  404. }
  405. handleReportData(report)
  406. }
  407. }
  408. }()
  409. if global.GVA_DB != nil {
  410. global.GVA_DB.Model(device).Updates(map[string]interface{}{"status": "online", "last_online_time": time.Now()})
  411. }
  412. // 启动时兜底:关闭历史遗留的同设备离线异常
  413. incidentService.NewIncidentService().ResolveDeviceOffline(device.DeviceCode)
  414. return nil
  415. }
  416. // waitDeviceRetry 等待重连间隔,同时允许停止设备时立即退出。
  417. func waitDeviceRetry(ctx context.Context, delay time.Duration) bool {
  418. timer := time.NewTimer(delay)
  419. defer timer.Stop()
  420. select {
  421. case <-ctx.Done():
  422. return false
  423. case <-timer.C:
  424. return true
  425. }
  426. }
  427. // parseReportData 解析上报数据。
  428. func parseReportData(deviceCode string, buf []byte) (*ReportData, error) {
  429. if len(buf) < 25 || buf[0] != 0xCF {
  430. return nil, errors.New("无效帧")
  431. }
  432. dataWithoutCRC := buf[:len(buf)-2]
  433. recvCRC := binary.LittleEndian.Uint16(buf[len(buf)-2:])
  434. calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC)))
  435. if recvCRC != calcCRC {
  436. return nil, errors.New("CRC校验失败")
  437. }
  438. rssi := int(buf[6])
  439. antenna := int(buf[22])
  440. epc := hex.EncodeToString(buf[11:23])
  441. return &ReportData{
  442. DeviceCode: deviceCode,
  443. Hex: hex.EncodeToString(buf),
  444. Epcs: []string{epc},
  445. RSSI: rssi,
  446. Antenna: antenna,
  447. Timestamp: time.Now().Unix(),
  448. }, nil
  449. }
  450. // handleReportData 处理UHF读取到的标签(精简版:只做解析+防抖,业务逻辑委托给 PassageService)
  451. func handleReportData(report *ReportData) {
  452. if len(report.Epcs) == 0 {
  453. return
  454. }
  455. epc := report.Epcs[0]
  456. if global.GVA_LOG != nil {
  457. global.GVA_LOG.Info("开始处理 RFID 通行",
  458. zap.String("device_code", report.DeviceCode),
  459. zap.String("epc", epc),
  460. zap.Int("rssi", report.RSSI),
  461. zap.Int("antenna", report.Antenna),
  462. )
  463. }
  464. // 委托给统一的进出场服务(所有触发方式共用同一入口)
  465. result, err := service.ServiceGroupApp.ParkingServiceGroup.PassageService.HandlePassage(common.PassageRequest{
  466. RFIDTag: epc,
  467. DeviceCode: report.DeviceCode,
  468. TriggerSource: "rfid",
  469. })
  470. if err != nil {
  471. if global.GVA_LOG != nil {
  472. global.GVA_LOG.Warn("RFID 通行处理失败",
  473. zap.String("device_code", report.DeviceCode),
  474. zap.String("epc", epc),
  475. zap.Error(err),
  476. )
  477. }
  478. if result != nil {
  479. PushChannelEvent(ChannelEvent{
  480. DeviceCode: report.DeviceCode, RFIDTag: epc, PlateNumber: result.PlateNumber,
  481. Direction: result.Direction, Timestamp: time.Now().Unix(), Status: result.GateStatus,
  482. Message: result.Message, Fee: result.Fee, StayTime: result.StayTime,
  483. })
  484. }
  485. return
  486. }
  487. if global.GVA_LOG != nil {
  488. global.GVA_LOG.Info("RFID 通行处理完成",
  489. zap.String("device_code", report.DeviceCode),
  490. zap.String("epc", epc),
  491. zap.String("direction", result.Direction),
  492. zap.String("plate_number", result.PlateNumber),
  493. zap.Uint("session_id", result.SessionID),
  494. zap.Float64("fee", result.Fee),
  495. zap.Bool("gate_opened", result.GateOpened),
  496. zap.String("gate_status", result.GateStatus),
  497. )
  498. }
  499. // 推送通道事件给前端
  500. var device dao.UHFReader
  501. if err := global.GVA_DB.Preload("Channel").First(&device, "device_code = ?", report.DeviceCode).Error; err == nil {
  502. PushChannelEvent(ChannelEvent{
  503. DeviceCode: report.DeviceCode,
  504. RFIDTag: epc,
  505. PlateNumber: result.PlateNumber,
  506. Direction: result.Direction,
  507. ChannelName: device.Channel.ChannelName,
  508. Timestamp: time.Now().Unix(),
  509. Status: "completed",
  510. Message: result.Message,
  511. Fee: result.Fee,
  512. StayTime: result.StayTime,
  513. })
  514. }
  515. }
  516. func readerTCPAddress(device *dao.UHFReader) string {
  517. if device.ConnectType != dao.ConnectTypeTCP {
  518. return ""
  519. }
  520. return fmt.Sprintf("%s:%d", device.IPAddress, device.Port)
  521. }
  522. // ==============================
  523. // 继电器控制(完整支持 Relay1 & Relay2)
  524. // ==============================
  525. const (
  526. RELAY_OP_RELEASE = 0x01
  527. RELAY_OP_CLOSE = 0x02
  528. )
  529. // CloseRelay1 开闸
  530. func (s *SerialReader) CloseRelay1(validTime byte) error {
  531. frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
  532. _, err := s.SendAndRecv(frame)
  533. return err
  534. }
  535. // ReleaseRelay1
  536. func (s *SerialReader) ReleaseRelay1() error {
  537. frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
  538. _, err := s.SendAndRecv(frame)
  539. return err
  540. }
  541. // CloseRelay2 关闸
  542. func (s *SerialReader) CloseRelay2(validTime byte) error {
  543. frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
  544. _, err := s.SendAndRecv(frame)
  545. return err
  546. }
  547. // ReleaseRelay2
  548. func (s *SerialReader) ReleaseRelay2() error {
  549. frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
  550. _, err := s.SendAndRecv(frame)
  551. return err
  552. }
  553. // buildRelayFrame 构建指令
  554. func buildRelayFrame(cmd uint16, relayNum byte, option byte, validTime byte) []byte {
  555. frame := []byte{
  556. 0xCF, 0xFF,
  557. byte(cmd >> 8), byte(cmd & 0xFF),
  558. 0x03, // len
  559. relayNum, // 1=继电器1 2=继电器2
  560. option,
  561. validTime,
  562. }
  563. crc := uiCrc16Cal(frame, uint8(len(frame)))
  564. frame = append(frame, byte(crc&0xFF), byte(crc>>8))
  565. return frame
  566. }