reader.go 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230
  1. package uhf
  2. import (
  3. "context"
  4. "encoding/binary"
  5. "encoding/hex"
  6. "errors"
  7. "fmt"
  8. "server/dao"
  9. "server/global"
  10. "sync"
  11. "time"
  12. )
  13. // Reader 通用读写接口(串口和TCP都实现它)
  14. type Reader interface {
  15. Connect() error
  16. Disconnect() error
  17. SendData([]byte) error
  18. ReadData() ([]byte, error)
  19. IsConnected() bool
  20. }
  21. // 上报数据模型
  22. type ReportData struct {
  23. DeviceCode string `json:"device_code"` // 设备编码
  24. Hex string `json:"hex"` // 原始报文
  25. Epcs []string `json:"epcs"` // 解析出的EPC列表
  26. RSSI int `json:"rssi"` // 信号强度
  27. Antenna int `json:"antenna"` // 天线号
  28. Timestamp int64 `json:"timestamp"` // 上报时间
  29. }
  30. // 全局设备管理器(管理所有已连接的设备)
  31. var DeviceManager = &Manager{
  32. devices: make(map[string]*DeviceHandler),
  33. }
  34. type Manager struct {
  35. mu sync.RWMutex
  36. devices map[string]*DeviceHandler
  37. }
  38. // DeviceHandler 每个设备的处理实例(包含读写器和协程)
  39. type DeviceHandler struct {
  40. reader Reader
  41. device *dao.UHFReader
  42. ctx context.Context
  43. cancel context.CancelFunc
  44. dataChan chan *ReportData
  45. isRunning bool
  46. }
  47. // Register 注册设备
  48. func (m *Manager) Register(code string, h *DeviceHandler) {
  49. m.mu.Lock()
  50. defer m.mu.Unlock()
  51. m.devices[code] = h
  52. }
  53. // Unregister 注销设备
  54. func (m *Manager) Unregister(code string) {
  55. m.mu.Lock()
  56. defer m.mu.Unlock()
  57. delete(m.devices, code)
  58. }
  59. // Get 获取设备
  60. func (m *Manager) Get(code string) (*DeviceHandler, bool) {
  61. m.mu.RLock()
  62. defer m.mu.RUnlock()
  63. h, ok := m.devices[code]
  64. return h, ok
  65. }
  66. // ===================== 复用原有CRC和解析函数(无需修改)=====================
  67. func uiCrc16Cal(pucY []byte, ucX uint8) uint16 {
  68. const PRESET_VALUE = 0xFFFF
  69. const POLYNOMIAL = 0x8408
  70. var uiCrcValue uint16 = PRESET_VALUE
  71. for ucI := uint8(0); ucI < ucX; ucI++ {
  72. uiCrcValue = uiCrcValue ^ uint16(pucY[ucI])
  73. for ucJ := uint8(0); ucJ < 8; ucJ++ {
  74. if uiCrcValue&0x0001 != 0 {
  75. uiCrcValue = (uiCrcValue >> 1) ^ POLYNOMIAL
  76. } else {
  77. uiCrcValue = uiCrcValue >> 1
  78. }
  79. }
  80. }
  81. return (uiCrcValue << 8) | (uiCrcValue >> 8)
  82. }
  83. // StartDeviceHandler 启动单个设备的读写协程
  84. func StartDeviceHandler(device *dao.UHFReader) error {
  85. // 1. 根据连接类型创建读写器
  86. var r Reader
  87. switch device.ConnectType {
  88. case "serial":
  89. r = NewSerialReader(device.COMPort, device.BaudRate)
  90. case "tcp":
  91. r = NewTCPReader(device.IPAddress, device.Port)
  92. default:
  93. return fmt.Errorf("不支持的连接类型: %s", device.ConnectType)
  94. }
  95. // 2. 连接设备
  96. if err := r.Connect(); err != nil {
  97. return err
  98. }
  99. // 3. 创建上下文和handler
  100. ctx, cancel := context.WithCancel(context.Background())
  101. h := &DeviceHandler{
  102. reader: r,
  103. device: device,
  104. ctx: ctx,
  105. cancel: cancel,
  106. dataChan: make(chan *ReportData, 100),
  107. isRunning: true,
  108. }
  109. // 4. 注册到管理器
  110. DeviceManager.Register(device.DeviceCode, h)
  111. // 5. 启动读写协程
  112. go func() {
  113. defer func() {
  114. h.isRunning = false
  115. r.Disconnect()
  116. close(h.dataChan)
  117. DeviceManager.Unregister(device.DeviceCode)
  118. }()
  119. for {
  120. select {
  121. case <-ctx.Done():
  122. return
  123. default:
  124. buf, err := r.ReadData()
  125. if err != nil {
  126. // 读取出错,更新设备状态
  127. global.GVA_DB.Model(device).Update("status", "offline")
  128. // 自动重连
  129. time.Sleep(2 * time.Second)
  130. if err := r.Connect(); err != nil {
  131. continue
  132. }
  133. global.GVA_DB.Model(device).Update("status", "online")
  134. continue
  135. }
  136. if len(buf) == 0 {
  137. continue
  138. }
  139. // 解析上报数据
  140. report, err := parseReportData(device.DeviceCode, buf)
  141. if err != nil {
  142. continue
  143. }
  144. // 推送到数据通道
  145. select {
  146. case h.dataChan <- report:
  147. default:
  148. }
  149. }
  150. }
  151. }()
  152. // 6. 启动业务处理协程(你后续在这里写业务逻辑)
  153. go func() {
  154. for {
  155. select {
  156. case <-ctx.Done():
  157. return
  158. case report, ok := <-h.dataChan:
  159. if !ok {
  160. return
  161. }
  162. // 这里写你的业务逻辑,比如:
  163. // 1. 去重过滤标签
  164. // 2. 写入数据库
  165. // 3. 推送到前端websocket
  166. handleReportData(report)
  167. }
  168. }
  169. }()
  170. // 更新设备状态
  171. global.GVA_DB.Model(device).Update("status", "online")
  172. return nil
  173. }
  174. // parseReportData 解析上报数据(把你之前的逻辑搬过来)
  175. func parseReportData(deviceCode string, buf []byte) (*ReportData, error) {
  176. if len(buf) < 25 || buf[0] != 0xCF {
  177. return nil, errors.New("无效帧")
  178. }
  179. // CRC校验(复用你之前的函数)
  180. dataWithoutCRC := buf[:len(buf)-2]
  181. recvCRC := binary.LittleEndian.Uint16(buf[len(buf)-2:])
  182. calcCRC := uiCrc16Cal(dataWithoutCRC, uint8(len(dataWithoutCRC)))
  183. if recvCRC != calcCRC {
  184. return nil, errors.New("CRC校验失败")
  185. }
  186. // 解析字段
  187. rssi := int(buf[6])
  188. antenna := int(buf[22])
  189. epc := hex.EncodeToString(buf[9:21])
  190. return &ReportData{
  191. DeviceCode: deviceCode,
  192. Hex: hex.EncodeToString(buf),
  193. Epcs: []string{epc},
  194. RSSI: rssi,
  195. Antenna: antenna,
  196. Timestamp: time.Now().Unix(),
  197. }, nil
  198. }
  199. // handleReportData
  200. func handleReportData(report *ReportData) {
  201. // 示例:打印上报数据
  202. fmt.Printf("[业务处理] 设备:%s EPC:%s RSSI:%d 天线:%d\n",
  203. report.DeviceCode, report.Epcs[0], report.RSSI, report.Antenna)
  204. }