package uhf import ( "context" "encoding/binary" "encoding/hex" "errors" "fmt" "sync" "time" "wails-app/internal/dao" "wails-app/internal/global" common "wails-app/internal/model/common" incidentService "wails-app/internal/modules/incident/service" "wails-app/internal/service" "wails-app/internal/service/parking" "go.uber.org/zap" ) // 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 } // ChannelEvent 通道事件(RFID/车牌识别后推送给前端) type ChannelEvent struct { ID uint64 `json:"id"` DeviceCode string `json:"device_code"` RFIDTag string `json:"rfid_tag"` PlateNumber string `json:"plate_number"` Direction string `json:"direction"` ChannelName string `json:"channel_name"` Timestamp int64 `json:"timestamp"` Status string `json:"status"` Message string `json:"message"` Fee float64 `json:"fee"` StayTime int64 `json:"stay_time"` } var ( channelEventMu sync.Mutex ChannelEventQueue = make([]ChannelEvent, 0, 50) channelEventNextID uint64 ) // PushChannelEvent 追加通道事件(保留最近 50 条) func PushChannelEvent(e ChannelEvent) { channelEventMu.Lock() defer channelEventMu.Unlock() channelEventNextID++ e.ID = channelEventNextID if e.Timestamp == 0 { e.Timestamp = time.Now().Unix() } if e.Status == "" { e.Status = "completed" } if len(ChannelEventQueue) >= 50 { ChannelEventQueue = ChannelEventQueue[1:] } ChannelEventQueue = append(ChannelEventQueue, e) } // GetChannelEvents 获取事件编号大于 afterID 的通道结果。 func GetChannelEvents(afterID uint64) []ChannelEvent { channelEventMu.Lock() defer channelEventMu.Unlock() result := make([]ChannelEvent, 0) for _, e := range ChannelEventQueue { if e.ID > afterID { result = append(result, e) } } return result } // 上报数据模型 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), } // init 注入统一道闸控制器(避免停车业务依赖具体设备协议)。 func init() { parking.SetGateController(DeviceManager) } 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 ioMu sync.Mutex stopOnce sync.Once done chan struct{} 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 } // StopDeviceHandler 停止设备读循环、关闭底层连接并等待退出。 func StopDeviceHandler(code string) error { h, ok := DeviceManager.Get(code) if !ok { return nil } h.stopOnce.Do(func() { if h.cancel != nil { h.cancel() } if h.reader != nil { _ = h.reader.Disconnect() } }) if h.done == nil { DeviceManager.Unregister(code) return nil } select { case <-h.done: case <-time.After(5 * time.Second): return fmt.Errorf("设备停止超时: %s", code) } return nil } // OpenGate 开闸(公开方法,供外部触发源调用,如摄像头识别、手动进出场) func (h *DeviceHandler) OpenGate(validTime byte) error { h.ioMu.Lock() defer h.ioMu.Unlock() return h.reader.CloseRelay1(validTime) } // CloseGate 关闸。 func (h *DeviceHandler) CloseGate(validTime byte) error { h.ioMu.Lock() defer h.ioMu.Unlock() return h.reader.CloseRelay2(validTime) } // OpenGate 通过设备编码开闸,实现 parking.GateController。 func (m *Manager) OpenGate(deviceCode string, validTime byte) error { h, ok := m.Get(deviceCode) if !ok { return fmt.Errorf("设备未连接: %s", deviceCode) } return h.OpenGate(validTime) } // CloseGate 通过设备编码关闸,实现 parking.GateController。 func (m *Manager) CloseGate(deviceCode string, validTime byte) error { h, ok := m.Get(deviceCode) if !ok { return fmt.Errorf("设备未连接: %s", deviceCode) } return h.CloseGate(validTime) } // IsGateConnected 判断设备是否已注册到当前进程。 func (m *Manager) IsGateConnected(deviceCode string) bool { h, ok := m.Get(deviceCode) return ok && h.reader != nil && h.reader.IsConnected() } // OpenGateByDeviceCode 通过设备编码直接开闸(便捷方法) func (m *Manager) OpenGateByDeviceCode(deviceCode string, validTime byte) error { return m.OpenGate(deviceCode, validTime) } // CloseGateByDeviceCode 通过设备编码直接关闸。 func (m *Manager) CloseGateByDeviceCode(deviceCode string, validTime byte) error { return m.CloseGate(deviceCode, validTime) } // ===================== 复用原有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 { if device == nil || device.DeviceCode == "" { return errors.New("设备编码不能为空") } if existing, ok := DeviceManager.Get(device.DeviceCode); ok { if existing.reader != nil && existing.reader.IsConnected() { return fmt.Errorf("设备已连接: %s", device.DeviceCode) } if err := StopDeviceHandler(device.DeviceCode); err != nil { return err } } // MQTT 设备由设备主动连接 broker,系统不主动建立连接。 if device.ConnectType == dao.ConnectTypeMQTT { return nil } if global.GVA_LOG != nil { global.GVA_LOG.Info("启动 UHF/RFID 设备连接", zap.String("device_code", device.DeviceCode), zap.String("device_name", device.DeviceName), zap.String("connect_type", device.ConnectType), zap.String("tcp_address", readerTCPAddress(device)), zap.Uint("channel_id", device.ChannelID), ) } // 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 { if global.GVA_LOG != nil { global.GVA_LOG.Error("UHF/RFID 设备连接失败", zap.String("device_code", device.DeviceCode), zap.String("connect_type", device.ConnectType), zap.String("tcp_address", readerTCPAddress(device)), zap.Error(err), ) } return err } if global.GVA_LOG != nil { global.GVA_LOG.Info("UHF/RFID 设备连接成功", zap.String("device_code", device.DeviceCode), zap.String("connect_type", device.ConnectType), zap.String("tcp_address", readerTCPAddress(device)), ) } // 3. 创建上下文和handler ctx, cancel := context.WithCancel(context.Background()) h := &DeviceHandler{ reader: r, device: device, ctx: ctx, cancel: cancel, dataChan: make(chan *ReportData, 200), done: make(chan struct{}), isRunning: true, } // 4. 注册到管理器 DeviceManager.Register(device.DeviceCode, h) // 5. 启动读写协程 go func() { defer func() { h.stopOnce.Do(func() { h.cancel() }) h.isRunning = false _ = r.Disconnect() close(h.dataChan) DeviceManager.Unregister(device.DeviceCode) close(h.done) if global.GVA_LOG != nil { global.GVA_LOG.Info("设备已停止: " + device.DeviceCode) } }() for { select { case <-ctx.Done(): return default: h.ioMu.Lock() buf, err := r.ReadData() h.ioMu.Unlock() if errors.Is(err, ErrReadTimeout) || errors.Is(err, ErrIncompleteFrame) { continue } if err != nil { if global.GVA_LOG != nil { global.GVA_LOG.Warn("UHF/RFID 读取失败,准备重连", zap.String("device_code", device.DeviceCode), zap.String("connect_type", device.ConnectType), zap.String("tcp_address", readerTCPAddress(device)), zap.Error(err), ) } if global.GVA_DB != nil { global.GVA_DB.Model(device).Update("status", "offline") } // 离线异常埋点(best-effort,同设备未关闭不重复生成) incidentService.NewIncidentService().RecordIncident(incidentService.RecordIncidentRequest{ Category: incidentService.CategoryDeviceOffline, Source: incidentService.SourceDevice, ParkingLotID: device.ParkingLotID, DeviceCode: device.DeviceCode, Description: "UHF 读卡器离线", Detail: err.Error(), }) _ = r.Disconnect() // 先断开 if !waitDeviceRetry(ctx, 2*time.Second) { return } if err := r.Connect(); err != nil { if global.GVA_LOG != nil { global.GVA_LOG.Warn("UHF/RFID 设备重连失败", zap.String("device_code", device.DeviceCode), zap.String("tcp_address", readerTCPAddress(device)), zap.Error(err), ) } continue } if global.GVA_LOG != nil { global.GVA_LOG.Info("UHF/RFID 设备重连成功", zap.String("device_code", device.DeviceCode), zap.String("tcp_address", readerTCPAddress(device)), ) } if global.GVA_DB != nil { global.GVA_DB.Model(device).Updates(map[string]interface{}{"status": "online", "last_online_time": time.Now()}) } // 设备恢复在线:自动关闭未处理的离线异常 incidentService.NewIncidentService().ResolveDeviceOffline(device.DeviceCode) continue } if len(buf) == 0 { continue } if global.GVA_LOG != nil { global.GVA_LOG.Info("收到 UHF/RFID 原始帧", zap.String("device_code", device.DeviceCode), zap.Int("packet_bytes", len(buf)), zap.String("packet_hex", hex.EncodeToString(buf)), ) } // 解析上报数据 report, err := parseReportData(device.DeviceCode, buf) if err != nil { if global.GVA_LOG != nil { global.GVA_LOG.Warn("UHF/RFID 报文解析失败", zap.String("device_code", device.DeviceCode), zap.String("packet_hex", hex.EncodeToString(buf)), zap.Error(err), ) } continue } if global.GVA_LOG != nil { global.GVA_LOG.Info("UHF/RFID 标签解析成功", zap.String("device_code", report.DeviceCode), zap.Strings("epcs", report.Epcs), zap.Int("rssi", report.RSSI), zap.Int("antenna", report.Antenna), ) } // 推送(不丢死,也不阻塞) select { case h.dataChan <- report: default: if global.GVA_LOG != nil { global.GVA_LOG.Warn("通道已满,丢弃数据: " + device.DeviceCode) } } } } }() // 6. 启动业务处理协程 go func() { for { select { case <-ctx.Done(): return case report, ok := <-h.dataChan: if !ok { return } handleReportData(report) } } }() if global.GVA_DB != nil { global.GVA_DB.Model(device).Updates(map[string]interface{}{"status": "online", "last_online_time": time.Now()}) } // 启动时兜底:关闭历史遗留的同设备离线异常 incidentService.NewIncidentService().ResolveDeviceOffline(device.DeviceCode) return nil } // waitDeviceRetry 等待重连间隔,同时允许停止设备时立即退出。 func waitDeviceRetry(ctx context.Context, delay time.Duration) bool { timer := time.NewTimer(delay) defer timer.Stop() select { case <-ctx.Done(): return false case <-timer.C: return true } } // 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]) return &ReportData{ DeviceCode: deviceCode, Hex: hex.EncodeToString(buf), Epcs: []string{epc}, RSSI: rssi, Antenna: antenna, Timestamp: time.Now().Unix(), }, nil } // handleReportData 处理UHF读取到的标签(精简版:只做解析+防抖,业务逻辑委托给 PassageService) func handleReportData(report *ReportData) { if len(report.Epcs) == 0 { return } epc := report.Epcs[0] if global.GVA_LOG != nil { global.GVA_LOG.Info("开始处理 RFID 通行", zap.String("device_code", report.DeviceCode), zap.String("epc", epc), zap.Int("rssi", report.RSSI), zap.Int("antenna", report.Antenna), ) } // 委托给统一的进出场服务(所有触发方式共用同一入口) result, err := service.ServiceGroupApp.ParkingServiceGroup.PassageService.HandlePassage(common.PassageRequest{ RFIDTag: epc, DeviceCode: report.DeviceCode, TriggerSource: "rfid", }) if err != nil { if global.GVA_LOG != nil { global.GVA_LOG.Warn("RFID 通行处理失败", zap.String("device_code", report.DeviceCode), zap.String("epc", epc), zap.Error(err), ) } if result != nil { PushChannelEvent(ChannelEvent{ DeviceCode: report.DeviceCode, RFIDTag: epc, PlateNumber: result.PlateNumber, Direction: result.Direction, Timestamp: time.Now().Unix(), Status: result.GateStatus, Message: result.Message, Fee: result.Fee, StayTime: result.StayTime, }) } return } if global.GVA_LOG != nil { global.GVA_LOG.Info("RFID 通行处理完成", zap.String("device_code", report.DeviceCode), zap.String("epc", epc), zap.String("direction", result.Direction), zap.String("plate_number", result.PlateNumber), zap.Uint("session_id", result.SessionID), zap.Float64("fee", result.Fee), zap.Bool("gate_opened", result.GateOpened), zap.String("gate_status", result.GateStatus), ) } // 推送通道事件给前端 var device dao.UHFReader if err := global.GVA_DB.Preload("Channel").First(&device, "device_code = ?", report.DeviceCode).Error; err == nil { PushChannelEvent(ChannelEvent{ DeviceCode: report.DeviceCode, RFIDTag: epc, PlateNumber: result.PlateNumber, Direction: result.Direction, ChannelName: device.Channel.ChannelName, Timestamp: time.Now().Unix(), Status: "completed", Message: result.Message, Fee: result.Fee, StayTime: result.StayTime, }) } } func readerTCPAddress(device *dao.UHFReader) string { if device.ConnectType != dao.ConnectTypeTCP { return "" } return fmt.Sprintf("%s:%d", device.IPAddress, device.Port) } // ============================== // 继电器控制(完整支持 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 }