tcp_reader.go 2.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127
  1. package uhf
  2. import (
  3. "errors"
  4. "fmt"
  5. "net"
  6. "time"
  7. )
  8. type TCPReader struct {
  9. IP string
  10. Port int
  11. conn net.Conn
  12. isConnected bool
  13. buffer []byte // 内部拼包缓冲区(关键修复)
  14. }
  15. // 实现 Reader 接口的继电器方法(不再panic)
  16. func (t *TCPReader) CloseRelay1(validTime byte) error {
  17. frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
  18. _, err := t.SendAndRecv(frame)
  19. return err
  20. }
  21. func (t *TCPReader) ReleaseRelay1() error {
  22. frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
  23. _, err := t.SendAndRecv(frame)
  24. return err
  25. }
  26. func (t *TCPReader) CloseRelay2(validTime byte) error {
  27. frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
  28. _, err := t.SendAndRecv(frame)
  29. return err
  30. }
  31. func (t *TCPReader) ReleaseRelay2() error {
  32. frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
  33. _, err := t.SendAndRecv(frame)
  34. return err
  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. addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port))
  44. conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
  45. if err != nil {
  46. t.isConnected = false
  47. return err
  48. }
  49. t.conn = conn
  50. t.isConnected = true
  51. return nil
  52. }
  53. func (t *TCPReader) Disconnect() error {
  54. if t.conn != nil {
  55. _ = t.conn.Close()
  56. }
  57. t.isConnected = false
  58. t.buffer = nil
  59. return nil
  60. }
  61. func (t *TCPReader) IsConnected() bool {
  62. return t.isConnected
  63. }
  64. func (t *TCPReader) SendData(data []byte) error {
  65. if !t.isConnected {
  66. return errors.New("TCP未连接")
  67. }
  68. _, err := t.conn.Write(data)
  69. if err != nil {
  70. t.isConnected = false
  71. }
  72. return err
  73. }
  74. // ==============================
  75. // 🔥 关键修复:TCP 自动拼包,返回完整25字节帧
  76. // ==============================
  77. func (t *TCPReader) ReadData() ([]byte, error) {
  78. if !t.isConnected {
  79. return nil, errors.New("TCP未连接")
  80. }
  81. // 读超时,防止永久阻塞
  82. t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
  83. buf := make([]byte, 1024)
  84. n, err := t.conn.Read(buf)
  85. if err != nil {
  86. t.isConnected = false
  87. return nil, err
  88. }
  89. // 加入缓冲区
  90. t.buffer = append(t.buffer, buf[:n]...)
  91. // 找完整帧:0xCF 开头 + 25字节
  92. for len(t.buffer) >= 25 {
  93. if t.buffer[0] == 0xCF {
  94. frame := t.buffer[:25]
  95. t.buffer = t.buffer[25:]
  96. return frame, nil
  97. } else {
  98. t.buffer = t.buffer[1:]
  99. }
  100. }
  101. return nil, errors.New("等待完整帧")
  102. }
  103. func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
  104. if err := t.SendData(data); err != nil {
  105. return nil, err
  106. }
  107. t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
  108. return t.ReadData()
  109. }