|
|
@@ -11,39 +11,41 @@ type TCPReader struct {
|
|
|
Port int
|
|
|
conn net.Conn
|
|
|
isConnected bool
|
|
|
- readBufSize int // 读取缓冲区大小
|
|
|
+ buffer []byte // 内部拼包缓冲区(关键修复)
|
|
|
}
|
|
|
|
|
|
+// 实现 Reader 接口的继电器方法(不再panic)
|
|
|
func (t *TCPReader) CloseRelay1(validTime byte) error {
|
|
|
- //TODO implement me
|
|
|
- panic("implement me")
|
|
|
+ frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
|
|
|
+ _, err := t.SendAndRecv(frame)
|
|
|
+ return err
|
|
|
}
|
|
|
|
|
|
func (t *TCPReader) ReleaseRelay1() error {
|
|
|
- //TODO implement me
|
|
|
- panic("implement me")
|
|
|
+ frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
|
|
|
+ _, err := t.SendAndRecv(frame)
|
|
|
+ return err
|
|
|
}
|
|
|
|
|
|
func (t *TCPReader) CloseRelay2(validTime byte) error {
|
|
|
- //TODO implement me
|
|
|
- panic("implement me")
|
|
|
+ frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
|
|
|
+ _, err := t.SendAndRecv(frame)
|
|
|
+ return err
|
|
|
}
|
|
|
|
|
|
func (t *TCPReader) ReleaseRelay2() error {
|
|
|
- //TODO implement me
|
|
|
- panic("implement me")
|
|
|
+ frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
|
|
|
+ _, err := t.SendAndRecv(frame)
|
|
|
+ return err
|
|
|
}
|
|
|
|
|
|
-// 初始化时设置默认缓冲区
|
|
|
func NewTCPReader(ip string, port int) *TCPReader {
|
|
|
return &TCPReader{
|
|
|
- IP: ip,
|
|
|
- Port: port,
|
|
|
- readBufSize: 1024, // 默认1K缓冲区
|
|
|
+ IP: ip,
|
|
|
+ Port: port,
|
|
|
}
|
|
|
}
|
|
|
|
|
|
-// Connect 建立TCP连接(原有逻辑不变)
|
|
|
func (t *TCPReader) Connect() error {
|
|
|
addr := net.JoinHostPort(t.IP, string(rune(t.Port)))
|
|
|
conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
|
|
|
@@ -56,51 +58,69 @@ func (t *TCPReader) Connect() error {
|
|
|
return nil
|
|
|
}
|
|
|
|
|
|
-// Disconnect 断开连接(原有逻辑不变)
|
|
|
func (t *TCPReader) Disconnect() error {
|
|
|
if t.conn != nil {
|
|
|
- err := t.conn.Close()
|
|
|
- t.isConnected = false
|
|
|
- return err
|
|
|
+ _ = t.conn.Close()
|
|
|
}
|
|
|
+ t.isConnected = false
|
|
|
+ t.buffer = nil
|
|
|
return nil
|
|
|
}
|
|
|
|
|
|
-// IsConnected 检查状态(原有逻辑不变)
|
|
|
func (t *TCPReader) IsConnected() bool {
|
|
|
return t.isConnected
|
|
|
}
|
|
|
|
|
|
-// SendData 仅发送数据(多设备批量发送用)
|
|
|
func (t *TCPReader) SendData(data []byte) error {
|
|
|
if !t.isConnected {
|
|
|
- return errors.New("TCP连接未建立")
|
|
|
+ return errors.New("TCP未连接")
|
|
|
}
|
|
|
_, err := t.conn.Write(data)
|
|
|
+ if err != nil {
|
|
|
+ t.isConnected = false
|
|
|
+ }
|
|
|
return err
|
|
|
}
|
|
|
|
|
|
-// ReadData 阻塞读取TCP数据(多设备监听用)
|
|
|
+// ==============================
|
|
|
+// 🔥 关键修复:TCP 自动拼包,返回完整25字节帧
|
|
|
+// ==============================
|
|
|
func (t *TCPReader) ReadData() ([]byte, error) {
|
|
|
if !t.isConnected {
|
|
|
- return nil, errors.New("TCP连接未建立")
|
|
|
+ return nil, errors.New("TCP未连接")
|
|
|
}
|
|
|
- buf := make([]byte, t.readBufSize)
|
|
|
+
|
|
|
+ // 读超时,防止永久阻塞
|
|
|
+ t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
|
|
|
+
|
|
|
+ buf := make([]byte, 1024)
|
|
|
n, err := t.conn.Read(buf)
|
|
|
if err != nil {
|
|
|
- t.isConnected = false // 读失败标记为断开
|
|
|
+ t.isConnected = false
|
|
|
return nil, err
|
|
|
}
|
|
|
- return buf[:n], nil
|
|
|
+
|
|
|
+ // 加入缓冲区
|
|
|
+ t.buffer = append(t.buffer, buf[:n]...)
|
|
|
+
|
|
|
+ // 找完整帧:0xCF 开头 + 25字节
|
|
|
+ for len(t.buffer) >= 25 {
|
|
|
+ if t.buffer[0] == 0xCF {
|
|
|
+ frame := t.buffer[:25]
|
|
|
+ t.buffer = t.buffer[25:]
|
|
|
+ return frame, nil
|
|
|
+ } else {
|
|
|
+ t.buffer = t.buffer[1:]
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ return nil, errors.New("等待完整帧")
|
|
|
}
|
|
|
|
|
|
-// SendAndRecv 发送并等待响应(一问一答,多设备指令交互用)
|
|
|
func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
|
|
|
if err := t.SendData(data); err != nil {
|
|
|
return nil, err
|
|
|
}
|
|
|
- // 设置读超时,避免永久阻塞
|
|
|
t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
|
|
|
- defer t.conn.SetReadDeadline(time.Time{}) // 恢复默认
|
|
|
return t.ReadData()
|
|
|
}
|