tcp_reader.go 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229
  1. package uhf
  2. import (
  3. "errors"
  4. "fmt"
  5. "net"
  6. "sync"
  7. "sync/atomic"
  8. "time"
  9. "go.uber.org/zap"
  10. "wails-app/internal/global"
  11. )
  12. type TCPReader struct {
  13. IP string
  14. Port int
  15. conn net.Conn
  16. connected atomic.Bool
  17. buffer []byte // 内部拼包缓冲区(关键修复)
  18. ioMu sync.Mutex
  19. }
  20. // 实现 Reader 接口的继电器方法(不再panic)
  21. func (t *TCPReader) CloseRelay1(validTime byte) error {
  22. frame := buildRelayFrame(0x0077, RELAY_OP_CLOSE, validTime)
  23. return t.sendRelayCommand("open", 1, validTime, frame)
  24. }
  25. func (t *TCPReader) ReleaseRelay1() error {
  26. frame := buildRelayFrame(0x0077, RELAY_OP_RELEASE, 0)
  27. return t.sendRelayCommand("release", 1, 0, frame)
  28. }
  29. func (t *TCPReader) CloseRelay2(validTime byte) error {
  30. frame := buildRelayFrame(0x0078, RELAY_OP_CLOSE, validTime)
  31. return t.sendRelayCommand("close", 2, validTime, frame)
  32. }
  33. func (t *TCPReader) ReleaseRelay2() error {
  34. frame := buildRelayFrame(0x0078, RELAY_OP_RELEASE, 0)
  35. return t.sendRelayCommand("release", 2, 0, frame)
  36. }
  37. func NewTCPReader(ip string, port int) *TCPReader {
  38. return &TCPReader{
  39. IP: ip,
  40. Port: port,
  41. }
  42. }
  43. func (t *TCPReader) Connect() error {
  44. t.ioMu.Lock()
  45. defer t.ioMu.Unlock()
  46. addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port))
  47. conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
  48. if err != nil {
  49. t.connected.Store(false)
  50. t.logWarn("UHF TCP 连接失败", zap.Error(err))
  51. return err
  52. }
  53. t.conn = conn
  54. t.connected.Store(true)
  55. t.logInfo("UHF TCP 连接成功")
  56. return nil
  57. }
  58. func (t *TCPReader) Disconnect() error {
  59. t.ioMu.Lock()
  60. defer t.ioMu.Unlock()
  61. if t.conn != nil {
  62. _ = t.conn.Close()
  63. }
  64. t.connected.Store(false)
  65. t.buffer = nil
  66. t.logInfo("UHF TCP 连接已关闭")
  67. return nil
  68. }
  69. func (t *TCPReader) IsConnected() bool {
  70. return t.connected.Load()
  71. }
  72. func (t *TCPReader) SendData(data []byte) error {
  73. t.ioMu.Lock()
  74. defer t.ioMu.Unlock()
  75. return t.sendDataLocked(data)
  76. }
  77. func (t *TCPReader) sendDataLocked(data []byte) error {
  78. if !t.connected.Load() {
  79. return errors.New("TCP未连接")
  80. }
  81. _, err := t.conn.Write(data)
  82. if err != nil {
  83. t.connected.Store(false)
  84. t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err))
  85. }
  86. return err
  87. }
  88. // ==============================
  89. // 🔥 关键修复:TCP 自动拼包,按协议变长帧提取(标签帧25字节,继电器应答9字节等)
  90. // ==============================
  91. func (t *TCPReader) ReadData() ([]byte, error) {
  92. t.ioMu.Lock()
  93. defer t.ioMu.Unlock()
  94. return t.readDataLocked()
  95. }
  96. func (t *TCPReader) readDataLocked() ([]byte, error) {
  97. if !t.connected.Load() {
  98. return nil, errors.New("TCP未连接")
  99. }
  100. // 读超时,防止永久阻塞
  101. t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
  102. buf := make([]byte, 1024)
  103. n, err := t.conn.Read(buf)
  104. if err != nil {
  105. if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
  106. return nil, ErrReadTimeout
  107. }
  108. t.connected.Store(false)
  109. t.logWarn("UHF TCP 读取失败", zap.Error(err))
  110. return nil, err
  111. }
  112. // 加入缓冲区
  113. t.buffer = append(t.buffer, buf[:n]...)
  114. t.logDebug("UHF TCP 收到数据块",
  115. zap.Int("packet_bytes", n),
  116. zap.Int("buffered_bytes", len(t.buffer)),
  117. zap.String("packet_hex", fmt.Sprintf("%x", buf[:n])),
  118. )
  119. frame, rest, ok := extractUHFFrame(t.buffer)
  120. t.buffer = rest
  121. if !ok {
  122. if len(t.buffer) > 0 {
  123. return nil, ErrIncompleteFrame
  124. }
  125. return nil, ErrReadTimeout
  126. }
  127. return frame, nil
  128. }
  129. func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
  130. t.ioMu.Lock()
  131. defer t.ioMu.Unlock()
  132. if err := t.sendDataLocked(data); err != nil {
  133. return nil, err
  134. }
  135. t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
  136. return t.readDataLocked()
  137. }
  138. // isRelayAckFor 判断响应帧是否为对应指令的应答(设备回显命令字)。
  139. func isRelayAckFor(response, frame []byte) bool {
  140. return len(response) >= 7 && response[0] == 0xCF &&
  141. response[2] == frame[2] && response[3] == frame[3]
  142. }
  143. func (t *TCPReader) sendRelayCommand(action string, relay int, validTime byte, frame []byte) error {
  144. t.logInfo("UHF TCP 道闸指令下发",
  145. zap.String("action", action),
  146. zap.Int("relay", relay),
  147. zap.Uint8("valid_time_seconds", validTime),
  148. zap.String("packet_hex", fmt.Sprintf("%x", frame)),
  149. )
  150. t.ioMu.Lock()
  151. defer t.ioMu.Unlock()
  152. if err := t.sendDataLocked(frame); err != nil {
  153. t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err))
  154. return err
  155. }
  156. // 等待设备应答(回显命令字)。期间收到的标签上报帧放回缓冲区交还读循环。
  157. for attempt := 0; attempt < 3; attempt++ {
  158. response, err := t.readDataLocked()
  159. if err != nil {
  160. break
  161. }
  162. if isRelayAckFor(response, frame) {
  163. status := "unknown"
  164. if len(response) >= 7 {
  165. status = fmt.Sprintf("0x%02x", response[5])
  166. }
  167. t.logInfo("UHF TCP 道闸指令收到响应",
  168. zap.String("action", action),
  169. zap.Int("relay", relay),
  170. zap.String("status", status),
  171. zap.String("response_hex", fmt.Sprintf("%x", response)),
  172. )
  173. if len(response) >= 7 && response[5] != 0x00 {
  174. return fmt.Errorf("道闸指令被设备拒绝: 状态码 %s", status)
  175. }
  176. return nil
  177. }
  178. t.buffer = append(append(make([]byte, 0, len(response)+len(t.buffer)), response...), t.buffer...)
  179. }
  180. err := fmt.Errorf("道闸指令未收到设备应答(帧已发出)")
  181. t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err))
  182. return err
  183. }
  184. func (t *TCPReader) logInfo(message string, fields ...zap.Field) {
  185. if global.GVA_LOG == nil {
  186. return
  187. }
  188. fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...)
  189. global.GVA_LOG.Info(message, fields...)
  190. }
  191. func (t *TCPReader) logDebug(message string, fields ...zap.Field) {
  192. if global.GVA_LOG == nil {
  193. return
  194. }
  195. fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...)
  196. global.GVA_LOG.Debug(message, fields...)
  197. }
  198. func (t *TCPReader) logWarn(message string, fields ...zap.Field) {
  199. if global.GVA_LOG == nil {
  200. return
  201. }
  202. fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...)
  203. global.GVA_LOG.Warn(message, fields...)
  204. }