tcp_reader.go 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197
  1. package uhf
  2. import (
  3. "errors"
  4. "fmt"
  5. "net"
  6. "sync"
  7. "time"
  8. "go.uber.org/zap"
  9. "wails-app/internal/global"
  10. )
  11. type TCPReader struct {
  12. IP string
  13. Port int
  14. conn net.Conn
  15. isConnected bool
  16. buffer []byte // 内部拼包缓冲区(关键修复)
  17. ioMu sync.Mutex
  18. }
  19. // 实现 Reader 接口的继电器方法(不再panic)
  20. func (t *TCPReader) CloseRelay1(validTime byte) error {
  21. frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
  22. return t.sendRelayCommand("open", 1, validTime, frame)
  23. }
  24. func (t *TCPReader) ReleaseRelay1() error {
  25. frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
  26. return t.sendRelayCommand("release", 1, 0, frame)
  27. }
  28. func (t *TCPReader) CloseRelay2(validTime byte) error {
  29. frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
  30. return t.sendRelayCommand("close", 2, validTime, frame)
  31. }
  32. func (t *TCPReader) ReleaseRelay2() error {
  33. frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
  34. return t.sendRelayCommand("release", 2, 0, frame)
  35. }
  36. func NewTCPReader(ip string, port int) *TCPReader {
  37. return &TCPReader{
  38. IP: ip,
  39. Port: port,
  40. }
  41. }
  42. func (t *TCPReader) Connect() error {
  43. t.ioMu.Lock()
  44. defer t.ioMu.Unlock()
  45. addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port))
  46. conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
  47. if err != nil {
  48. t.isConnected = false
  49. t.logWarn("UHF TCP 连接失败", zap.Error(err))
  50. return err
  51. }
  52. t.conn = conn
  53. t.isConnected = true
  54. t.logInfo("UHF TCP 连接成功")
  55. return nil
  56. }
  57. func (t *TCPReader) Disconnect() error {
  58. t.ioMu.Lock()
  59. defer t.ioMu.Unlock()
  60. if t.conn != nil {
  61. _ = t.conn.Close()
  62. }
  63. t.isConnected = false
  64. t.buffer = nil
  65. t.logInfo("UHF TCP 连接已关闭")
  66. return nil
  67. }
  68. func (t *TCPReader) IsConnected() bool {
  69. t.ioMu.Lock()
  70. defer t.ioMu.Unlock()
  71. return t.isConnected
  72. }
  73. func (t *TCPReader) SendData(data []byte) error {
  74. t.ioMu.Lock()
  75. defer t.ioMu.Unlock()
  76. return t.sendDataLocked(data)
  77. }
  78. func (t *TCPReader) sendDataLocked(data []byte) error {
  79. if !t.isConnected {
  80. return errors.New("TCP未连接")
  81. }
  82. _, err := t.conn.Write(data)
  83. if err != nil {
  84. t.isConnected = false
  85. t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err))
  86. }
  87. return err
  88. }
  89. // ==============================
  90. // 🔥 关键修复:TCP 自动拼包,返回完整25字节帧
  91. // ==============================
  92. func (t *TCPReader) ReadData() ([]byte, error) {
  93. t.ioMu.Lock()
  94. defer t.ioMu.Unlock()
  95. return t.readDataLocked()
  96. }
  97. func (t *TCPReader) readDataLocked() ([]byte, error) {
  98. if !t.isConnected {
  99. return nil, errors.New("TCP未连接")
  100. }
  101. // 读超时,防止永久阻塞
  102. t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
  103. buf := make([]byte, 1024)
  104. n, err := t.conn.Read(buf)
  105. if err != nil {
  106. if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
  107. return nil, ErrReadTimeout
  108. }
  109. t.isConnected = false
  110. t.logWarn("UHF TCP 读取失败", zap.Error(err))
  111. return nil, err
  112. }
  113. // 加入缓冲区
  114. t.buffer = append(t.buffer, buf[:n]...)
  115. t.logInfo("UHF TCP 收到数据块",
  116. zap.Int("packet_bytes", n),
  117. zap.Int("buffered_bytes", len(t.buffer)),
  118. zap.String("packet_hex", fmt.Sprintf("%x", buf[:n])),
  119. )
  120. // 找完整帧:0xCF 开头 + 25字节
  121. for len(t.buffer) >= 25 {
  122. if t.buffer[0] == 0xCF {
  123. frame := t.buffer[:25]
  124. t.buffer = t.buffer[25:]
  125. return frame, nil
  126. } else {
  127. t.buffer = t.buffer[1:]
  128. }
  129. }
  130. return nil, ErrIncompleteFrame
  131. }
  132. func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
  133. t.ioMu.Lock()
  134. defer t.ioMu.Unlock()
  135. if err := t.sendDataLocked(data); err != nil {
  136. return nil, err
  137. }
  138. t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
  139. return t.readDataLocked()
  140. }
  141. func (t *TCPReader) sendRelayCommand(action string, relay int, validTime byte, frame []byte) error {
  142. t.logInfo("UHF TCP 道闸指令下发",
  143. zap.String("action", action),
  144. zap.Int("relay", relay),
  145. zap.Uint8("valid_time_seconds", validTime),
  146. zap.String("packet_hex", fmt.Sprintf("%x", frame)),
  147. )
  148. response, err := t.SendAndRecv(frame)
  149. if err != nil {
  150. t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err))
  151. return err
  152. }
  153. t.logInfo("UHF TCP 道闸指令收到响应",
  154. zap.String("action", action),
  155. zap.Int("response_bytes", len(response)),
  156. zap.String("response_hex", fmt.Sprintf("%x", response)),
  157. )
  158. return nil
  159. }
  160. func (t *TCPReader) logInfo(message string, fields ...zap.Field) {
  161. if global.GVA_LOG == nil {
  162. return
  163. }
  164. fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...)
  165. global.GVA_LOG.Info(message, fields...)
  166. }
  167. func (t *TCPReader) logWarn(message string, fields ...zap.Field) {
  168. if global.GVA_LOG == nil {
  169. return
  170. }
  171. fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...)
  172. global.GVA_LOG.Warn(message, fields...)
  173. }