package uhf import ( "errors" "fmt" "net" "sync" "sync/atomic" "time" "go.uber.org/zap" "wails-app/internal/global" ) type TCPReader struct { IP string Port int conn net.Conn connected atomic.Bool buffer []byte // 内部拼包缓冲区(关键修复) ioMu sync.Mutex } // 实现 Reader 接口的继电器方法(不再panic) func (t *TCPReader) CloseRelay1(validTime byte) error { frame := buildRelayFrame(0x0077, RELAY_OP_CLOSE, validTime) return t.sendRelayCommand("open", 1, validTime, frame) } func (t *TCPReader) ReleaseRelay1() error { frame := buildRelayFrame(0x0077, RELAY_OP_RELEASE, 0) return t.sendRelayCommand("release", 1, 0, frame) } func (t *TCPReader) CloseRelay2(validTime byte) error { frame := buildRelayFrame(0x0078, RELAY_OP_CLOSE, validTime) return t.sendRelayCommand("close", 2, validTime, frame) } func (t *TCPReader) ReleaseRelay2() error { frame := buildRelayFrame(0x0078, RELAY_OP_RELEASE, 0) return t.sendRelayCommand("release", 2, 0, frame) } func NewTCPReader(ip string, port int) *TCPReader { return &TCPReader{ IP: ip, Port: port, } } func (t *TCPReader) Connect() error { t.ioMu.Lock() defer t.ioMu.Unlock() addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)) conn, err := net.DialTimeout("tcp", addr, 5*time.Second) if err != nil { t.connected.Store(false) t.logWarn("UHF TCP 连接失败", zap.Error(err)) return err } t.conn = conn t.connected.Store(true) t.logInfo("UHF TCP 连接成功") return nil } func (t *TCPReader) Disconnect() error { t.ioMu.Lock() defer t.ioMu.Unlock() if t.conn != nil { _ = t.conn.Close() } t.connected.Store(false) t.buffer = nil t.logInfo("UHF TCP 连接已关闭") return nil } func (t *TCPReader) IsConnected() bool { return t.connected.Load() } func (t *TCPReader) SendData(data []byte) error { t.ioMu.Lock() defer t.ioMu.Unlock() return t.sendDataLocked(data) } func (t *TCPReader) sendDataLocked(data []byte) error { if !t.connected.Load() { return errors.New("TCP未连接") } _, err := t.conn.Write(data) if err != nil { t.connected.Store(false) t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err)) } return err } // ============================== // 🔥 关键修复:TCP 自动拼包,按协议变长帧提取(标签帧25字节,继电器应答9字节等) // ============================== func (t *TCPReader) ReadData() ([]byte, error) { t.ioMu.Lock() defer t.ioMu.Unlock() return t.readDataLocked() } func (t *TCPReader) readDataLocked() ([]byte, error) { if !t.connected.Load() { return nil, errors.New("TCP未连接") } // 读超时,防止永久阻塞 t.conn.SetReadDeadline(time.Now().Add(3 * time.Second)) buf := make([]byte, 1024) n, err := t.conn.Read(buf) if err != nil { if netErr, ok := err.(net.Error); ok && netErr.Timeout() { return nil, ErrReadTimeout } t.connected.Store(false) t.logWarn("UHF TCP 读取失败", zap.Error(err)) return nil, err } // 加入缓冲区 t.buffer = append(t.buffer, buf[:n]...) t.logDebug("UHF TCP 收到数据块", zap.Int("packet_bytes", n), zap.Int("buffered_bytes", len(t.buffer)), zap.String("packet_hex", fmt.Sprintf("%x", buf[:n])), ) frame, rest, ok := extractUHFFrame(t.buffer) t.buffer = rest if !ok { if len(t.buffer) > 0 { return nil, ErrIncompleteFrame } return nil, ErrReadTimeout } return frame, nil } func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) { t.ioMu.Lock() defer t.ioMu.Unlock() if err := t.sendDataLocked(data); err != nil { return nil, err } t.conn.SetReadDeadline(time.Now().Add(3 * time.Second)) return t.readDataLocked() } // isRelayAckFor 判断响应帧是否为对应指令的应答(设备回显命令字)。 func isRelayAckFor(response, frame []byte) bool { return len(response) >= 7 && response[0] == 0xCF && response[2] == frame[2] && response[3] == frame[3] } func (t *TCPReader) sendRelayCommand(action string, relay int, validTime byte, frame []byte) error { t.logInfo("UHF TCP 道闸指令下发", zap.String("action", action), zap.Int("relay", relay), zap.Uint8("valid_time_seconds", validTime), zap.String("packet_hex", fmt.Sprintf("%x", frame)), ) t.ioMu.Lock() defer t.ioMu.Unlock() if err := t.sendDataLocked(frame); err != nil { t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err)) return err } // 等待设备应答(回显命令字)。期间收到的标签上报帧放回缓冲区交还读循环。 for attempt := 0; attempt < 3; attempt++ { response, err := t.readDataLocked() if err != nil { break } if isRelayAckFor(response, frame) { status := "unknown" if len(response) >= 7 { status = fmt.Sprintf("0x%02x", response[5]) } t.logInfo("UHF TCP 道闸指令收到响应", zap.String("action", action), zap.Int("relay", relay), zap.String("status", status), zap.String("response_hex", fmt.Sprintf("%x", response)), ) if len(response) >= 7 && response[5] != 0x00 { return fmt.Errorf("道闸指令被设备拒绝: 状态码 %s", status) } return nil } t.buffer = append(append(make([]byte, 0, len(response)+len(t.buffer)), response...), t.buffer...) } err := fmt.Errorf("道闸指令未收到设备应答(帧已发出)") t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err)) return err } func (t *TCPReader) logInfo(message string, fields ...zap.Field) { if global.GVA_LOG == nil { return } fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...) global.GVA_LOG.Info(message, fields...) } func (t *TCPReader) logDebug(message string, fields ...zap.Field) { if global.GVA_LOG == nil { return } fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...) global.GVA_LOG.Debug(message, fields...) } func (t *TCPReader) logWarn(message string, fields ...zap.Field) { if global.GVA_LOG == nil { return } fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...) global.GVA_LOG.Warn(message, fields...) }