reader.go 11 KB

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