package uhf import ( "errors" "fmt" "net" "sync" "time" "go.uber.org/zap" "wails-app/internal/global" ) type TCPReader struct { IP string Port int conn net.Conn isConnected bool buffer []byte // 内部拼包缓冲区(关键修复) ioMu sync.Mutex } // 实现 Reader 接口的继电器方法(不再panic) func (t *TCPReader) CloseRelay1(validTime byte) error { frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime) return t.sendRelayCommand("open", 1, validTime, frame) } func (t *TCPReader) ReleaseRelay1() error { frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0) return t.sendRelayCommand("release", 1, 0, frame) } func (t *TCPReader) CloseRelay2(validTime byte) error { frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime) return t.sendRelayCommand("close", 2, validTime, frame) } func (t *TCPReader) ReleaseRelay2() error { frame := buildRelayFrame(0x0078, 2, 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.isConnected = false t.logWarn("UHF TCP 连接失败", zap.Error(err)) return err } t.conn = conn t.isConnected = 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.isConnected = false t.buffer = nil t.logInfo("UHF TCP 连接已关闭") return nil } func (t *TCPReader) IsConnected() bool { t.ioMu.Lock() defer t.ioMu.Unlock() return t.isConnected } 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.isConnected { return errors.New("TCP未连接") } _, err := t.conn.Write(data) if err != nil { t.isConnected = false t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err)) } return err } // ============================== // 🔥 关键修复:TCP 自动拼包,返回完整25字节帧 // ============================== func (t *TCPReader) ReadData() ([]byte, error) { t.ioMu.Lock() defer t.ioMu.Unlock() return t.readDataLocked() } func (t *TCPReader) readDataLocked() ([]byte, error) { if !t.isConnected { 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.isConnected = false t.logWarn("UHF TCP 读取失败", zap.Error(err)) return nil, err } // 加入缓冲区 t.buffer = append(t.buffer, buf[:n]...) t.logInfo("UHF TCP 收到数据块", zap.Int("packet_bytes", n), zap.Int("buffered_bytes", len(t.buffer)), zap.String("packet_hex", fmt.Sprintf("%x", buf[:n])), ) // 找完整帧:0xCF 开头 + 25字节 for len(t.buffer) >= 25 { if t.buffer[0] == 0xCF { frame := t.buffer[:25] t.buffer = t.buffer[25:] return frame, nil } else { t.buffer = t.buffer[1:] } } return nil, ErrIncompleteFrame } 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() } 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)), ) response, err := t.SendAndRecv(frame) if err != nil { t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err)) return err } t.logInfo("UHF TCP 道闸指令收到响应", zap.String("action", action), zap.Int("response_bytes", len(response)), zap.String("response_hex", fmt.Sprintf("%x", response)), ) return nil } 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) 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...) }