package parking import ( "context" "encoding/binary" "encoding/hex" "errors" "fmt" "wails-app/internal/dao" "wails-app/internal/global" "wails-app/internal/model/vehicle/request" "wails-app/internal/service" "sync" "time" ) // Reader 通用读写接口(串口和TCP都实现它) type Reader interface { Connect() error Disconnect() error SendData([]byte) error ReadData() ([]byte, error) IsConnected() bool // ===================== 加入继电器接口 ===================== CloseRelay1(validTime byte) error ReleaseRelay1() error CloseRelay2(validTime byte) error ReleaseRelay2() error } // 上报数据模型 type ReportData struct { DeviceCode string `json:"device_code"` // 设备编码 Hex string `json:"hex"` // 原始报文 Epcs []string `json:"epcs"` // 解析出的EPC列表 RSSI int `json:"rssi"` // 信号强度 Antenna int `json:"antenna"` // 天线号 Timestamp int64 `json:"timestamp"` // 上报时间 } // 全局设备管理器(管理所有已连接的设备) var DeviceManager = &Manager{ devices: make(map[string]*DeviceHandler), } type Manager struct { mu sync.RWMutex devices map[string]*DeviceHandler } // DeviceHandler 每个设备的处理实例(包含读写器和协程) type DeviceHandler struct { reader Reader device *dao.UHFReader ctx context.Context cancel context.CancelFunc dataChan chan *ReportData isRunning bool } // Register 注册设备 func (m *Manager) Register(code string, h *DeviceHandler) { m.mu.Lock() defer m.mu.Unlock() m.devices[code] = h } // Unregister 注销设备 func (m *Manager) Unregister(code string) { m.mu.Lock() defer m.mu.Unlock() delete(m.devices, code) } // Get 获取设备 func (m *Manager) Get(code string) (*DeviceHandler, bool) { m.mu.RLock() defer m.mu.RUnlock() h, ok := m.devices[code] return h, ok } // ===================== 复用原有CRC和解析函数(无需修改)===================== func uiCrc16Cal(pucY []byte, ucX uint8) uint16 { const PRESET_VALUE = 0xFFFF const POLYNOMIAL = 0x8408 var uiCrcValue uint16 = PRESET_VALUE for ucI := uint8(0); ucI < ucX; ucI++ { uiCrcValue = uiCrcValue ^ uint16(pucY[ucI]) for ucJ := uint8(0); ucJ < 8; ucJ++ { if uiCrcValue&0x0001 != 0 { uiCrcValue = (uiCrcValue >> 1) ^ POLYNOMIAL } else { uiCrcValue = uiCrcValue >> 1 } } } return (uiCrcValue << 8) | (uiCrcValue >> 8) } // StartDeviceHandler 启动单个设备的读写协程 func StartDeviceHandler(device *dao.UHFReader) error { // 1. 根据连接类型创建读写器 var r Reader switch device.ConnectType { case "serial": r = NewSerialReader(device.COMPort, device.BaudRate) case "tcp": r = NewTCPReader(device.IPAddress, device.Port) default: return fmt.Errorf("不支持的连接类型: %s", device.ConnectType) } // 2. 连接设备 if err := r.Connect(); err != nil { return err } // 3. 创建上下文和handler ctx, cancel := context.WithCancel(context.Background()) h := &DeviceHandler{ reader: r, device: device, ctx: ctx, cancel: cancel, dataChan: make(chan *ReportData, 200), // 适当加大 isRunning: true, } // 4. 注册到管理器 DeviceManager.Register(device.DeviceCode, h) // 5. 启动读写协程 go func() { defer func() { h.isRunning = false _ = r.Disconnect() close(h.dataChan) DeviceManager.Unregister(device.DeviceCode) global.GVA_LOG.Info("设备已停止 code:" + device.DeviceCode) }() for { select { case <-ctx.Done(): return default: buf, err := r.ReadData() if err != nil { global.GVA_DB.Model(device).Update("status", "offline") _ = r.Disconnect() // 先断开 time.Sleep(2 * time.Second) if err := r.Connect(); err != nil { continue } global.GVA_DB.Model(device).Update("status", "online") continue } if len(buf) == 0 { continue } // 解析上报数据 report, err := parseReportData(device.DeviceCode, buf) if err != nil { continue } // 推送(不丢死,也不阻塞) select { case h.dataChan <- report: default: global.GVA_LOG.Warn("通道已满,丢弃数据 device" + device.DeviceCode) } } } }() // 6. 启动业务处理协程 go func() { for { select { case <-ctx.Done(): return case report, ok := <-h.dataChan: if !ok { return } handleReportData(report) } } }() global.GVA_DB.Model(device).Update("status", "online") return nil } // parseReportData 解析上报数据 func parseReportData(deviceCode string, buf []byte) (*ReportData, error) { if len(buf) < 25 || buf[0] != 0xCF { return nil, errors.New("无效帧") } dataWithoutCRC := buf[:len(buf)-2] recvCRC := binary.LittleEndian.Uint16(buf[len(buf)-2:]) calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC))) if recvCRC != calcCRC { return nil, errors.New("CRC校验失败") } rssi := int(buf[6]) antenna := int(buf[22]) epc := hex.EncodeToString(buf[11:23]) fmt.Printf(hex.EncodeToString(buf)) return &ReportData{ DeviceCode: deviceCode, Hex: hex.EncodeToString(buf), Epcs: []string{epc}, RSSI: rssi, Antenna: antenna, Timestamp: time.Now().Unix(), }, nil } // 防抖与状态 var ( epcLastTime = make(map[string]int64) epcStatus = make(map[string]bool) epcMutex sync.Mutex debounceSecond = int64(3) ) // handleReportData 处理UHF读取到的标签 + 完整权限判断 func handleReportData(report *ReportData) { if len(report.Epcs) == 0 { return } epc := report.Epcs[0] fmt.Println("", epc) now := time.Now().Unix() // ===================== 【修复】锁只包防抖,绝不包业务!===================== epcMutex.Lock() // 3秒防抖 if lastTime, ok := epcLastTime[epc]; ok && now-lastTime < debounceSecond { epcMutex.Unlock() return } epcLastTime[epc] = now epcMutex.Unlock() // ===================== 锁已释放,后面逻辑不会阻塞 ===================== // 获取设备处理器 _, exists := DeviceManager.Get(report.DeviceCode) if !exists { return } // ===================== 查询设备 + 通道 ===================== var device dao.UHFReader err := global.GVA_DB.Preload("Channel").First(&device, "device_code = ?", report.DeviceCode).Error if err != nil { return } channel := device.Channel direction := channel.Direction allowTemporary := channel.AllowTemporary vehicleService := service.ServiceGroupApp.VehicleServiceGroup.VehicleService shortlistService := service.ServiceGroupApp.VehicleServiceGroup.ShortlistService // ===================== 1. 根据EPC查询车辆 ===================== vehicle, err := vehicleService.GetVehicleByPlateNumber("", epc) isTempVehicle := false if err != nil || vehicle.ID == 0 { isTempVehicle = true } else { if vehicle.VehicleType != nil { isTempVehicle = (vehicle.VehicleTypeID == 1) } else { isTempVehicle = true } } // ===================== 2. 查询黑白名单 ===================== isBlack, isWhite := shortlistService.CheckVehicleShortlist(epc) // ===================== 3. 黑名单优先判断 ===================== if isBlack { fmt.Printf("🚫 禁止通行 | 车辆在黑名单中 | 设备:%s EPC:%s\n", report.DeviceCode, epc) return } // ===================== 4. 临时车规则 ===================== if isTempVehicle { if !allowTemporary && !isWhite { fmt.Printf("🚫 禁止通行 | 临时车不允许此通道且不在白名单 | 设备:%s EPC:%s\n", report.DeviceCode, epc) return } if !isWhite { fmt.Printf("🚫 禁止通行 | 临时车不在白名单 | 设备:%s EPC:%s\n", report.DeviceCode, epc) return } } // ===================== 5. 入口 ===================== if direction == "in" { epcMutex.Lock() canIn := !epcStatus[epc] epcMutex.Unlock() if canIn { fmt.Printf("🟢 入口开闸 | 设备:%s EPC:%s\n", report.DeviceCode, epc) //handler.reader.CloseRelay1(2) var vehicleEntry request.VehicleEntry vehicleEntry.RFIDTag = epc vehicleEntry.ParkingLotID = channel.ParkingLotID _, err := vehicleService.VehicleEntry(vehicleEntry) if err == nil { epcMutex.Lock() epcStatus[epc] = true epcMutex.Unlock() } } } // ===================== 6. 出口 ===================== if direction == "out" { fmt.Printf("🔴 出口开闸 | 设备:%s EPC:%s\n", report.DeviceCode, epc) //handler.reader.CloseRelay1(2) var vehicleExit request.VehicleExit vehicleExit.RFIDTag = epc _, err := vehicleService.VehicleExit(vehicleExit) if err == nil { epcMutex.Lock() epcStatus[epc] = false epcMutex.Unlock() } } } // ===================== 手动重置离场 ===================== func SetEpcExit(epc string) { epcMutex.Lock() defer epcMutex.Unlock() epcStatus[epc] = false fmt.Printf("🚙 手动重置离场 | EPC: %s\n", epc) } // ============================== // 继电器控制(完整支持 Relay1 & Relay2) // ============================== const ( RELAY_OP_RELEASE = 0x01 RELAY_OP_CLOSE = 0x02 ) // CloseRelay1 开闸 func (s *SerialReader) CloseRelay1(validTime byte) error { frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime) _, err := s.SendAndRecv(frame) return err } // ReleaseRelay1 func (s *SerialReader) ReleaseRelay1() error { frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0) _, err := s.SendAndRecv(frame) return err } // CloseRelay2 关闸 func (s *SerialReader) CloseRelay2(validTime byte) error { frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime) _, err := s.SendAndRecv(frame) return err } // ReleaseRelay2 func (s *SerialReader) ReleaseRelay2() error { frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0) _, err := s.SendAndRecv(frame) return err } // buildRelayFrame 构建指令 func buildRelayFrame(cmd uint16, relayNum byte, option byte, validTime byte) []byte { frame := []byte{ 0xCF, 0xFF, byte(cmd >> 8), byte(cmd & 0xFF), 0x03, // len relayNum, // 1=继电器1 2=继电器2 option, validTime, } crc := uiCrc16Cal(frame, uint8(len(frame))) frame = append(frame, byte(crc&0xFF), byte(crc>>8)) return frame }