uhf_reader_test.go 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248
  1. package test
  2. import (
  3. "context"
  4. "encoding/binary"
  5. "encoding/hex"
  6. "fmt"
  7. "internal/service/uhf"
  8. "sync"
  9. "testing"
  10. "time"
  11. )
  12. // ===================== 复用原有CRC和解析函数(无需修改)=====================
  13. func uiCrc16Cal(pucY []byte, ucX uint8) uint16 {
  14. const PRESET_VALUE = 0xFFFF
  15. const POLYNOMIAL = 0x8408
  16. var uiCrcValue uint16 = PRESET_VALUE
  17. for ucI := uint8(0); ucI < ucX; ucI++ {
  18. uiCrcValue = uiCrcValue ^ uint16(pucY[ucI])
  19. for ucJ := uint8(0); ucJ < 8; ucJ++ {
  20. if uiCrcValue&0x0001 != 0 {
  21. uiCrcValue = (uiCrcValue >> 1) ^ POLYNOMIAL
  22. } else {
  23. uiCrcValue = uiCrcValue >> 1
  24. }
  25. }
  26. }
  27. return (uiCrcValue << 8) | (uiCrcValue >> 8)
  28. }
  29. func parseResp(resp []byte) (bool, error) {
  30. if len(resp) < 7 {
  31. return false, fmt.Errorf("返回报文长度不足,实际长度: %d", len(resp))
  32. }
  33. if resp[0] != 0xCF {
  34. return false, fmt.Errorf("帧头错误,期望0xCF,实际0x%02X", resp[0])
  35. }
  36. dataWithoutCRC := resp[:len(resp)-2]
  37. recvCRC := binary.LittleEndian.Uint16(resp[len(resp)-2:])
  38. calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC)))
  39. if recvCRC != calcCRC {
  40. return false, fmt.Errorf("CRC校验失败 | 接收:0x%04X | 计算:0x%04X", recvCRC, calcCRC)
  41. }
  42. fmt.Println("✅ 设备上报原始Hex:", hex.EncodeToString(resp))
  43. fmt.Println("=====================================")
  44. cmd := binary.BigEndian.Uint16(resp[2:4])
  45. switch cmd {
  46. case 0x0072:
  47. fmt.Println("📦 解析【设备版本信息】")
  48. if len(resp) >= 8 {
  49. fmt.Printf("硬件版本: %d.%d\n", resp[6], resp[7])
  50. if len(resp) >= 9 {
  51. fmt.Printf("软件版本: %d.%d\n", resp[8], resp[9])
  52. }
  53. }
  54. case 0x0064:
  55. fmt.Println("📦 解析【设备远程网口信息】")
  56. if len(resp) >= 16 {
  57. status := resp[6]
  58. option := resp[7]
  59. enable := resp[8]
  60. ip := fmt.Sprintf("%d.%d.%d.%d", resp[9], resp[10], resp[11], resp[12])
  61. port := binary.BigEndian.Uint16([]byte{resp[13], resp[14]})
  62. heartTime := resp[15]
  63. fmt.Printf("执行状态: 0x%02X (%s)\n", status, map[byte]string{0x00: "成功", 0xFF: "失败"}[status])
  64. fmt.Printf("操作类型: 0x%02X (0x02=读取)\n", option)
  65. fmt.Printf("上报使能: %d (1=开启, 0=关闭)\n", enable)
  66. fmt.Printf("远程IP: %s\n", ip)
  67. fmt.Printf("远程端口: %d\n", port)
  68. fmt.Printf("心跳时间: %d*5S = %dS\n", heartTime, heartTime*5)
  69. }
  70. default:
  71. fmt.Println("📦 解析【设备主动上报数据】")
  72. // 这里根据设备上报的自定义指令格式扩展解析逻辑
  73. }
  74. fmt.Println("=====================================")
  75. return true, nil
  76. }
  77. // ===================== 串口持续读取核心逻辑 =====================
  78. // SerialReaderWrapper 封装串口读取器,增加协程控制
  79. type SerialReaderWrapper struct {
  80. reader *uhf.SerialReader
  81. ctx context.Context
  82. cancel context.CancelFunc
  83. dataChan chan []byte // 上报数据通道
  84. errChan chan error // 错误通道
  85. wg sync.WaitGroup // 等待组,确保协程优雅退出
  86. isRunning bool // 读取协程运行状态
  87. comPort string // 串口号
  88. baudRate int // 波特率
  89. reconnect bool // 是否自动重连
  90. }
  91. // NewSerialReaderWrapper 创建串口读取包装器
  92. func NewSerialReaderWrapper(comPort string, baudRate int, reconnect bool) *SerialReaderWrapper {
  93. ctx, cancel := context.WithCancel(context.Background())
  94. return &SerialReaderWrapper{
  95. reader: uhf.NewSerialReader(comPort, baudRate),
  96. ctx: ctx,
  97. cancel: cancel,
  98. dataChan: make(chan []byte, 100), // 缓冲通道,避免阻塞(根据上报频率调整大小)
  99. errChan: make(chan error, 10),
  100. comPort: comPort,
  101. baudRate: baudRate,
  102. reconnect: reconnect,
  103. }
  104. }
  105. // StartReadLoop 启动串口持续读取协程
  106. func (w *SerialReaderWrapper) StartReadLoop() error {
  107. if w.isRunning {
  108. return fmt.Errorf("读取协程已在运行")
  109. }
  110. // 先连接串口
  111. if err := w.reader.Connect(); err != nil {
  112. return fmt.Errorf("串口初始连接失败: %v", err)
  113. }
  114. w.isRunning = true
  115. w.wg.Add(1)
  116. // 启动读取协程
  117. go func() {
  118. defer func() {
  119. // 协程退出时清理
  120. w.isRunning = false
  121. w.reader.Disconnect()
  122. w.wg.Done()
  123. close(w.dataChan)
  124. close(w.errChan)
  125. fmt.Println("🔌 串口读取协程已退出")
  126. }()
  127. fmt.Println("🔄 串口持续读取协程已启动,等待设备上报数据...")
  128. for {
  129. select {
  130. // 检测退出信号
  131. case <-w.ctx.Done():
  132. fmt.Println("📢 收到退出信号,停止读取串口")
  133. return
  134. default:
  135. // 循环读取串口数据
  136. buf, err := w.reader.ReadData()
  137. if err != nil {
  138. w.errChan <- fmt.Errorf("读取串口失败: %v", err)
  139. // 自动重连逻辑
  140. if w.reconnect {
  141. fmt.Println("🔁 尝试重新连接串口...")
  142. w.reader.Disconnect()
  143. // 重连间隔2秒,避免频繁重试
  144. time.Sleep(2 * time.Second)
  145. if reconnectErr := w.reader.Connect(); reconnectErr != nil {
  146. w.errChan <- fmt.Errorf("串口重连失败: %v", reconnectErr)
  147. } else {
  148. w.errChan <- fmt.Errorf("串口重连成功")
  149. }
  150. } else {
  151. // 不重连则退出协程
  152. return
  153. }
  154. continue
  155. }
  156. // 读取到有效数据,发送到数据通道
  157. if len(buf) > 0 {
  158. select {
  159. case w.dataChan <- buf:
  160. // 数据成功写入通道
  161. default:
  162. // 通道满时丢弃(或根据业务调整,比如阻塞/报错)
  163. w.errChan <- fmt.Errorf("数据通道已满,丢弃上报数据: %x", buf)
  164. }
  165. }
  166. }
  167. }
  168. }()
  169. return nil
  170. }
  171. // StopReadLoop 停止读取协程(优雅退出)
  172. func (w *SerialReaderWrapper) StopReadLoop() {
  173. if !w.isRunning {
  174. return
  175. }
  176. // 发送退出信号
  177. w.cancel()
  178. // 等待协程退出
  179. w.wg.Wait()
  180. fmt.Println("✅ 串口读取协程已优雅停止")
  181. }
  182. // ProcessData 处理上报数据(可在主线程/另一个协程中运行)
  183. func (w *SerialReaderWrapper) ProcessData() {
  184. for {
  185. select {
  186. case <-w.ctx.Done():
  187. return
  188. // 处理上报数据
  189. case data, ok := <-w.dataChan:
  190. if !ok {
  191. return
  192. }
  193. fmt.Println("\n📥 收到设备上报数据:")
  194. // 解析数据(复用你已有的parseResp函数)
  195. if _, err := parseResp(data); err != nil {
  196. fmt.Printf("❌ 解析上报数据失败: %v | 原始数据: %x\n", err, data)
  197. }
  198. // 处理错误
  199. case err, ok := <-w.errChan:
  200. if !ok {
  201. return
  202. }
  203. fmt.Printf("⚠️ 串口读取错误: %v\n", err)
  204. }
  205. }
  206. }
  207. // ===================== 测试主函数 =====================
  208. func TestContinuousSerialRead(t *testing.T) {
  209. // 1. 创建串口包装器(COM9,115200波特率,开启自动重连)
  210. serialWrapper := NewSerialReaderWrapper("COM9", 115200, true)
  211. defer serialWrapper.StopReadLoop() // 测试结束时停止协程
  212. // 2. 启动持续读取协程
  213. if err := serialWrapper.StartReadLoop(); err != nil {
  214. t.Fatalf("❌ 启动串口读取协程失败: %v", err)
  215. }
  216. // 3. 启动数据处理协程(主线程也可以直接调用ProcessData,这里用协程避免阻塞)
  217. go serialWrapper.ProcessData()
  218. // 4. 主线程保持运行(模拟业务逻辑,比如运行30秒后退出)
  219. fmt.Println("🚀 主线程运行中,按Ctrl+C退出...")
  220. time.Sleep(30 * time.Second) // 运行30秒,可替换为无限循环(for {})
  221. // 5. 主动停止(可选,defer也会自动停止)
  222. serialWrapper.StopReadLoop()
  223. }