package test import ( "context" "encoding/binary" "encoding/hex" "fmt" "server/service/uhf" "sync" "testing" "time" ) // ===================== 复用原有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) } func parseResp(resp []byte) (bool, error) { if len(resp) < 7 { return false, fmt.Errorf("返回报文长度不足,实际长度: %d", len(resp)) } if resp[0] != 0xCF { return false, fmt.Errorf("帧头错误,期望0xCF,实际0x%02X", resp[0]) } dataWithoutCRC := resp[:len(resp)-2] recvCRC := binary.LittleEndian.Uint16(resp[len(resp)-2:]) calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC))) if recvCRC != calcCRC { return false, fmt.Errorf("CRC校验失败 | 接收:0x%04X | 计算:0x%04X", recvCRC, calcCRC) } fmt.Println("✅ 设备上报原始Hex:", hex.EncodeToString(resp)) fmt.Println("=====================================") cmd := binary.BigEndian.Uint16(resp[2:4]) switch cmd { case 0x0072: fmt.Println("📦 解析【设备版本信息】") if len(resp) >= 8 { fmt.Printf("硬件版本: %d.%d\n", resp[6], resp[7]) if len(resp) >= 9 { fmt.Printf("软件版本: %d.%d\n", resp[8], resp[9]) } } case 0x0064: fmt.Println("📦 解析【设备远程网口信息】") if len(resp) >= 16 { status := resp[6] option := resp[7] enable := resp[8] ip := fmt.Sprintf("%d.%d.%d.%d", resp[9], resp[10], resp[11], resp[12]) port := binary.BigEndian.Uint16([]byte{resp[13], resp[14]}) heartTime := resp[15] fmt.Printf("执行状态: 0x%02X (%s)\n", status, map[byte]string{0x00: "成功", 0xFF: "失败"}[status]) fmt.Printf("操作类型: 0x%02X (0x02=读取)\n", option) fmt.Printf("上报使能: %d (1=开启, 0=关闭)\n", enable) fmt.Printf("远程IP: %s\n", ip) fmt.Printf("远程端口: %d\n", port) fmt.Printf("心跳时间: %d*5S = %dS\n", heartTime, heartTime*5) } default: fmt.Println("📦 解析【设备主动上报数据】") // 这里根据设备上报的自定义指令格式扩展解析逻辑 } fmt.Println("=====================================") return true, nil } // ===================== 串口持续读取核心逻辑 ===================== // SerialReaderWrapper 封装串口读取器,增加协程控制 type SerialReaderWrapper struct { reader *uhf.SerialReader ctx context.Context cancel context.CancelFunc dataChan chan []byte // 上报数据通道 errChan chan error // 错误通道 wg sync.WaitGroup // 等待组,确保协程优雅退出 isRunning bool // 读取协程运行状态 comPort string // 串口号 baudRate int // 波特率 reconnect bool // 是否自动重连 } // NewSerialReaderWrapper 创建串口读取包装器 func NewSerialReaderWrapper(comPort string, baudRate int, reconnect bool) *SerialReaderWrapper { ctx, cancel := context.WithCancel(context.Background()) return &SerialReaderWrapper{ reader: uhf.NewSerialReader(comPort, baudRate), ctx: ctx, cancel: cancel, dataChan: make(chan []byte, 100), // 缓冲通道,避免阻塞(根据上报频率调整大小) errChan: make(chan error, 10), comPort: comPort, baudRate: baudRate, reconnect: reconnect, } } // StartReadLoop 启动串口持续读取协程 func (w *SerialReaderWrapper) StartReadLoop() error { if w.isRunning { return fmt.Errorf("读取协程已在运行") } // 先连接串口 if err := w.reader.Connect(); err != nil { return fmt.Errorf("串口初始连接失败: %v", err) } w.isRunning = true w.wg.Add(1) // 启动读取协程 go func() { defer func() { // 协程退出时清理 w.isRunning = false w.reader.Disconnect() w.wg.Done() close(w.dataChan) close(w.errChan) fmt.Println("🔌 串口读取协程已退出") }() fmt.Println("🔄 串口持续读取协程已启动,等待设备上报数据...") for { select { // 检测退出信号 case <-w.ctx.Done(): fmt.Println("📢 收到退出信号,停止读取串口") return default: // 循环读取串口数据 buf, err := w.reader.ReadData() if err != nil { w.errChan <- fmt.Errorf("读取串口失败: %v", err) // 自动重连逻辑 if w.reconnect { fmt.Println("🔁 尝试重新连接串口...") w.reader.Disconnect() // 重连间隔2秒,避免频繁重试 time.Sleep(2 * time.Second) if reconnectErr := w.reader.Connect(); reconnectErr != nil { w.errChan <- fmt.Errorf("串口重连失败: %v", reconnectErr) } else { w.errChan <- fmt.Errorf("串口重连成功") } } else { // 不重连则退出协程 return } continue } // 读取到有效数据,发送到数据通道 if len(buf) > 0 { select { case w.dataChan <- buf: // 数据成功写入通道 default: // 通道满时丢弃(或根据业务调整,比如阻塞/报错) w.errChan <- fmt.Errorf("数据通道已满,丢弃上报数据: %x", buf) } } } } }() return nil } // StopReadLoop 停止读取协程(优雅退出) func (w *SerialReaderWrapper) StopReadLoop() { if !w.isRunning { return } // 发送退出信号 w.cancel() // 等待协程退出 w.wg.Wait() fmt.Println("✅ 串口读取协程已优雅停止") } // ProcessData 处理上报数据(可在主线程/另一个协程中运行) func (w *SerialReaderWrapper) ProcessData() { for { select { case <-w.ctx.Done(): return // 处理上报数据 case data, ok := <-w.dataChan: if !ok { return } fmt.Println("\n📥 收到设备上报数据:") // 解析数据(复用你已有的parseResp函数) if _, err := parseResp(data); err != nil { fmt.Printf("❌ 解析上报数据失败: %v | 原始数据: %x\n", err, data) } // 处理错误 case err, ok := <-w.errChan: if !ok { return } fmt.Printf("⚠️ 串口读取错误: %v\n", err) } } } // ===================== 测试主函数 ===================== func TestContinuousSerialRead(t *testing.T) { // 1. 创建串口包装器(COM9,115200波特率,开启自动重连) serialWrapper := NewSerialReaderWrapper("COM9", 115200, true) defer serialWrapper.StopReadLoop() // 测试结束时停止协程 // 2. 启动持续读取协程 if err := serialWrapper.StartReadLoop(); err != nil { t.Fatalf("❌ 启动串口读取协程失败: %v", err) } // 3. 启动数据处理协程(主线程也可以直接调用ProcessData,这里用协程避免阻塞) go serialWrapper.ProcessData() // 4. 主线程保持运行(模拟业务逻辑,比如运行30秒后退出) fmt.Println("🚀 主线程运行中,按Ctrl+C退出...") time.Sleep(30 * time.Second) // 运行30秒,可替换为无限循环(for {}) // 5. 主动停止(可选,defer也会自动停止) serialWrapper.StopReadLoop() }