reader.go 9.3 KB

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