reader.go 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416
  1. package uhf
  2. import (
  3. "context"
  4. "encoding/binary"
  5. "encoding/hex"
  6. "errors"
  7. "fmt"
  8. "server/dao"
  9. "server/global"
  10. "server/model/vehicle/request"
  11. "server/service"
  12. "sync"
  13. "time"
  14. )
  15. // Reader 通用读写接口(串口和TCP都实现它)
  16. type Reader interface {
  17. Connect() error
  18. Disconnect() error
  19. SendData([]byte) error
  20. ReadData() ([]byte, error)
  21. IsConnected() bool
  22. // ===================== 加入继电器接口 =====================
  23. CloseRelay1(validTime byte) error
  24. ReleaseRelay1() error
  25. CloseRelay2(validTime byte) error
  26. ReleaseRelay2() error
  27. }
  28. // 上报数据模型
  29. type ReportData struct {
  30. DeviceCode string `json:"device_code"` // 设备编码
  31. Hex string `json:"hex"` // 原始报文
  32. Epcs []string `json:"epcs"` // 解析出的EPC列表
  33. RSSI int `json:"rssi"` // 信号强度
  34. Antenna int `json:"antenna"` // 天线号
  35. Timestamp int64 `json:"timestamp"` // 上报时间
  36. }
  37. // 全局设备管理器(管理所有已连接的设备)
  38. var DeviceManager = &Manager{
  39. devices: make(map[string]*DeviceHandler),
  40. }
  41. type Manager struct {
  42. mu sync.RWMutex
  43. devices map[string]*DeviceHandler
  44. }
  45. // DeviceHandler 每个设备的处理实例(包含读写器和协程)
  46. type DeviceHandler struct {
  47. reader Reader
  48. device *dao.UHFReader
  49. ctx context.Context
  50. cancel context.CancelFunc
  51. dataChan chan *ReportData
  52. isRunning bool
  53. }
  54. // Register 注册设备
  55. func (m *Manager) Register(code string, h *DeviceHandler) {
  56. m.mu.Lock()
  57. defer m.mu.Unlock()
  58. m.devices[code] = h
  59. }
  60. // Unregister 注销设备
  61. func (m *Manager) Unregister(code string) {
  62. m.mu.Lock()
  63. defer m.mu.Unlock()
  64. delete(m.devices, code)
  65. }
  66. // Get 获取设备
  67. func (m *Manager) Get(code string) (*DeviceHandler, bool) {
  68. m.mu.RLock()
  69. defer m.mu.RUnlock()
  70. h, ok := m.devices[code]
  71. return h, ok
  72. }
  73. // ===================== 复用原有CRC和解析函数(无需修改)=====================
  74. func uiCrc16Cal(pucY []byte, ucX uint8) uint16 {
  75. const PRESET_VALUE = 0xFFFF
  76. const POLYNOMIAL = 0x8408
  77. var uiCrcValue uint16 = PRESET_VALUE
  78. for ucI := uint8(0); ucI < ucX; ucI++ {
  79. uiCrcValue = uiCrcValue ^ uint16(pucY[ucI])
  80. for ucJ := uint8(0); ucJ < 8; ucJ++ {
  81. if uiCrcValue&0x0001 != 0 {
  82. uiCrcValue = (uiCrcValue >> 1) ^ POLYNOMIAL
  83. } else {
  84. uiCrcValue = uiCrcValue >> 1
  85. }
  86. }
  87. }
  88. return (uiCrcValue << 8) | (uiCrcValue >> 8)
  89. }
  90. // StartDeviceHandler 启动单个设备的读写协程
  91. func StartDeviceHandler(device *dao.UHFReader) error {
  92. // 1. 根据连接类型创建读写器
  93. var r Reader
  94. switch device.ConnectType {
  95. case "serial":
  96. r = NewSerialReader(device.COMPort, device.BaudRate)
  97. case "tcp":
  98. r = NewTCPReader(device.IPAddress, device.Port)
  99. default:
  100. return fmt.Errorf("不支持的连接类型: %s", device.ConnectType)
  101. }
  102. // 2. 连接设备
  103. if err := r.Connect(); err != nil {
  104. return err
  105. }
  106. // 3. 创建上下文和handler
  107. ctx, cancel := context.WithCancel(context.Background())
  108. h := &DeviceHandler{
  109. reader: r,
  110. device: device,
  111. ctx: ctx,
  112. cancel: cancel,
  113. dataChan: make(chan *ReportData, 100),
  114. isRunning: true,
  115. }
  116. // 4. 注册到管理器
  117. DeviceManager.Register(device.DeviceCode, h)
  118. // 5. 启动读写协程
  119. go func() {
  120. defer func() {
  121. h.isRunning = false
  122. r.Disconnect()
  123. close(h.dataChan)
  124. DeviceManager.Unregister(device.DeviceCode)
  125. }()
  126. for {
  127. select {
  128. case <-ctx.Done():
  129. return
  130. default:
  131. buf, err := r.ReadData()
  132. if err != nil {
  133. // 读取出错,更新设备状态
  134. global.GVA_DB.Model(device).Update("status", "offline")
  135. // 自动重连
  136. time.Sleep(2 * time.Second)
  137. if err := r.Connect(); err != nil {
  138. continue
  139. }
  140. global.GVA_DB.Model(device).Update("status", "online")
  141. continue
  142. }
  143. if len(buf) == 0 {
  144. continue
  145. }
  146. // 解析上报数据
  147. report, err := parseReportData(device.DeviceCode, buf)
  148. if err != nil {
  149. continue
  150. }
  151. // 推送到数据通道
  152. select {
  153. case h.dataChan <- report:
  154. default:
  155. }
  156. }
  157. }
  158. }()
  159. // 6. 启动业务处理协程
  160. go func() {
  161. for {
  162. select {
  163. case <-ctx.Done():
  164. return
  165. case report, ok := <-h.dataChan:
  166. if !ok {
  167. return
  168. }
  169. handleReportData(report)
  170. }
  171. }
  172. }()
  173. // 更新设备状态
  174. global.GVA_DB.Model(device).Update("status", "online")
  175. return nil
  176. }
  177. // parseReportData 解析上报数据
  178. func parseReportData(deviceCode string, buf []byte) (*ReportData, error) {
  179. if len(buf) < 25 || buf[0] != 0xCF {
  180. return nil, errors.New("无效帧")
  181. }
  182. dataWithoutCRC := buf[:len(buf)-2]
  183. recvCRC := binary.LittleEndian.Uint16(buf[len(buf)-2:])
  184. calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC)))
  185. if recvCRC != calcCRC {
  186. return nil, errors.New("CRC校验失败")
  187. }
  188. rssi := int(buf[6])
  189. antenna := int(buf[22])
  190. epc := hex.EncodeToString(buf[11:21])
  191. return &ReportData{
  192. DeviceCode: deviceCode,
  193. Hex: hex.EncodeToString(buf),
  194. Epcs: []string{epc},
  195. RSSI: rssi,
  196. Antenna: antenna,
  197. Timestamp: time.Now().Unix(),
  198. }, nil
  199. }
  200. // 防抖与状态
  201. var (
  202. epcLastTime = make(map[string]int64)
  203. epcStatus = make(map[string]bool)
  204. epcMutex sync.Mutex
  205. debounceSecond = int64(3)
  206. )
  207. // handleReportData 处理UHF读取到的标签 + 完整权限判断
  208. func handleReportData(report *ReportData) {
  209. if len(report.Epcs) == 0 {
  210. return
  211. }
  212. epc := report.Epcs[0]
  213. fmt.Println("<UNK>", epc)
  214. now := time.Now().Unix()
  215. epcMutex.Lock()
  216. defer epcMutex.Unlock()
  217. // ===================== 修复:3秒防抖 =====================
  218. if lastTime, ok := epcLastTime[epc]; ok && now-lastTime < debounceSecond {
  219. return
  220. }
  221. // 获取设备处理器
  222. handler, exists := DeviceManager.Get(report.DeviceCode)
  223. if !exists {
  224. return
  225. }
  226. // ===================== 查询设备 + 通道 =====================
  227. var device dao.UHFReader
  228. err := global.GVA_DB.Preload("Channel").First(&device, "device_code = ?", report.DeviceCode).Error
  229. if err != nil {
  230. return
  231. }
  232. channel := device.Channel
  233. direction := channel.Direction
  234. allowTemporary := channel.AllowTemporary
  235. vehicleService := service.ServiceGroupApp.VehicleServiceGroup.VehicleService
  236. shortlistService := service.ServiceGroupApp.VehicleServiceGroup.ShortlistService
  237. // ===================== 1. 根据EPC查询车辆 =====================
  238. vehicle, err := vehicleService.GetVehicleByPlateNumber("", epc)
  239. isTempVehicle := false
  240. if err != nil || vehicle.ID == 0 {
  241. isTempVehicle = true
  242. } else {
  243. // ===================== 修复:防止空指针panic =====================
  244. if vehicle.VehicleType != nil {
  245. isTempVehicle = (vehicle.VehicleTypeID == 1)
  246. } else {
  247. isTempVehicle = true
  248. }
  249. }
  250. // ===================== 2. 查询黑白名单 =====================
  251. isBlack, isWhite := shortlistService.CheckVehicleShortlist(epc)
  252. // ===================== 3. 黑名单优先判断 =====================
  253. if isBlack {
  254. fmt.Printf("🚫 禁止通行 | 车辆在黑名单中 | 设备:%s EPC:%s\n", report.DeviceCode, epc)
  255. epcLastTime[epc] = now
  256. return
  257. }
  258. // ===================== 4. 临时车规则 =====================
  259. if isTempVehicle {
  260. if !allowTemporary && !isWhite {
  261. fmt.Printf("🚫 禁止通行 | 临时车不允许此通道且不在白名单 | 设备:%s EPC:%s\n", report.DeviceCode, epc)
  262. epcLastTime[epc] = now
  263. return
  264. }
  265. if !isWhite {
  266. fmt.Printf("🚫 禁止通行 | 临时车不在白名单 | 设备:%s EPC:%s\n", report.DeviceCode, epc)
  267. epcLastTime[epc] = now
  268. return
  269. }
  270. }
  271. // ===================== 5. 入口 =====================
  272. if direction == "in" {
  273. // 关键:只要状态是 false 就允许进入
  274. if !epcStatus[epc] {
  275. fmt.Printf("🟢 入口开闸 | 设备:%s EPC:%s 临时车:%v 白名单:%v\n", report.DeviceCode, epc, isTempVehicle, isWhite)
  276. handler.reader.CloseRelay1(2)
  277. // 记录入场
  278. var vehicleEntry request.VehicleEntry
  279. vehicleEntry.RFIDTag = epc
  280. vehicleEntry.ParkingLotID = channel.ParkingLotID
  281. _, err := vehicleService.VehicleEntry(vehicleEntry)
  282. if err != nil {
  283. fmt.Printf("❌ 入场记录失败:%v | EPC:%s\n", err, epc)
  284. // 关键:不设置 epcStatus[epc] = true
  285. // 让下一次刷卡可以继续进入
  286. } else {
  287. // ✅ 只有成功才修改状态
  288. epcStatus[epc] = true
  289. }
  290. } else {
  291. fmt.Printf("⚠️ 已在场内,禁止重复入场 | EPC:%s\n", epc)
  292. }
  293. }
  294. // ===================== 6. 出口 =====================
  295. if direction == "out" {
  296. fmt.Printf("🔴 出口开闸 | 设备:%s EPC:%s 临时车:%v 白名单:%v\n", report.DeviceCode, epc, isTempVehicle, isWhite)
  297. handler.reader.CloseRelay1(2)
  298. // 记录出场
  299. var vehicleExit request.VehicleExit
  300. vehicleExit.RFIDTag = epc
  301. _, err := vehicleService.VehicleExit(vehicleExit)
  302. if err != nil {
  303. fmt.Printf("❌ 出场记录失败:%v | EPC:%s\n", err, epc)
  304. // 关键:不设置 epcStatus[epc] = false
  305. } else {
  306. // ✅ 只有成功才重置状态
  307. epcStatus[epc] = false
  308. }
  309. }
  310. // 更新最后读卡时间
  311. epcLastTime[epc] = now
  312. }
  313. // ===================== 手动重置离场 =====================
  314. func SetEpcExit(epc string) {
  315. epcMutex.Lock()
  316. defer epcMutex.Unlock()
  317. epcStatus[epc] = false
  318. fmt.Printf("🚙 手动重置离场 | EPC: %s\n", epc)
  319. }
  320. // ==============================
  321. // 继电器控制(完整支持 Relay1 & Relay2)
  322. // ==============================
  323. const (
  324. RELAY_OP_RELEASE = 0x01
  325. RELAY_OP_CLOSE = 0x02
  326. )
  327. // CloseRelay1 开闸
  328. func (s *SerialReader) CloseRelay1(validTime byte) error {
  329. frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
  330. _, err := s.SendAndRecv(frame)
  331. return err
  332. }
  333. // ReleaseRelay1
  334. func (s *SerialReader) ReleaseRelay1() error {
  335. frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
  336. _, err := s.SendAndRecv(frame)
  337. return err
  338. }
  339. // CloseRelay2 关闸
  340. func (s *SerialReader) CloseRelay2(validTime byte) error {
  341. frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
  342. _, err := s.SendAndRecv(frame)
  343. return err
  344. }
  345. // ReleaseRelay2
  346. func (s *SerialReader) ReleaseRelay2() error {
  347. frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
  348. _, err := s.SendAndRecv(frame)
  349. return err
  350. }
  351. // buildRelayFrame 构建指令
  352. func buildRelayFrame(cmd uint16, relayNum byte, option byte, validTime byte) []byte {
  353. frame := []byte{
  354. 0xCF, 0xFF,
  355. byte(cmd >> 8), byte(cmd & 0xFF),
  356. 0x03, // len
  357. relayNum, // 1=继电器1 2=继电器2
  358. option,
  359. validTime,
  360. }
  361. crc := uiCrc16Cal(frame, uint8(len(frame)))
  362. frame = append(frame, byte(crc&0xFF), byte(crc>>8))
  363. return frame
  364. }