|
@@ -5,6 +5,7 @@ import (
|
|
|
"fmt"
|
|
"fmt"
|
|
|
"net"
|
|
"net"
|
|
|
"sync"
|
|
"sync"
|
|
|
|
|
+ "sync/atomic"
|
|
|
"time"
|
|
"time"
|
|
|
|
|
|
|
|
"go.uber.org/zap"
|
|
"go.uber.org/zap"
|
|
@@ -12,32 +13,32 @@ import (
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
type TCPReader struct {
|
|
type TCPReader struct {
|
|
|
- IP string
|
|
|
|
|
- Port int
|
|
|
|
|
- conn net.Conn
|
|
|
|
|
- isConnected bool
|
|
|
|
|
- buffer []byte // 内部拼包缓冲区(关键修复)
|
|
|
|
|
- ioMu sync.Mutex
|
|
|
|
|
|
|
+ IP string
|
|
|
|
|
+ Port int
|
|
|
|
|
+ conn net.Conn
|
|
|
|
|
+ connected atomic.Bool
|
|
|
|
|
+ buffer []byte // 内部拼包缓冲区(关键修复)
|
|
|
|
|
+ ioMu sync.Mutex
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// 实现 Reader 接口的继电器方法(不再panic)
|
|
// 实现 Reader 接口的继电器方法(不再panic)
|
|
|
func (t *TCPReader) CloseRelay1(validTime byte) error {
|
|
func (t *TCPReader) CloseRelay1(validTime byte) error {
|
|
|
- frame := buildRelayFrame(0x0077, 1, RELAY_OP_CLOSE, validTime)
|
|
|
|
|
|
|
+ frame := buildRelayFrame(0x0077, RELAY_OP_CLOSE, validTime)
|
|
|
return t.sendRelayCommand("open", 1, validTime, frame)
|
|
return t.sendRelayCommand("open", 1, validTime, frame)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) ReleaseRelay1() error {
|
|
func (t *TCPReader) ReleaseRelay1() error {
|
|
|
- frame := buildRelayFrame(0x0077, 1, RELAY_OP_RELEASE, 0)
|
|
|
|
|
|
|
+ frame := buildRelayFrame(0x0077, RELAY_OP_RELEASE, 0)
|
|
|
return t.sendRelayCommand("release", 1, 0, frame)
|
|
return t.sendRelayCommand("release", 1, 0, frame)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) CloseRelay2(validTime byte) error {
|
|
func (t *TCPReader) CloseRelay2(validTime byte) error {
|
|
|
- frame := buildRelayFrame(0x0078, 2, RELAY_OP_CLOSE, validTime)
|
|
|
|
|
|
|
+ frame := buildRelayFrame(0x0078, RELAY_OP_CLOSE, validTime)
|
|
|
return t.sendRelayCommand("close", 2, validTime, frame)
|
|
return t.sendRelayCommand("close", 2, validTime, frame)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) ReleaseRelay2() error {
|
|
func (t *TCPReader) ReleaseRelay2() error {
|
|
|
- frame := buildRelayFrame(0x0078, 2, RELAY_OP_RELEASE, 0)
|
|
|
|
|
|
|
+ frame := buildRelayFrame(0x0078, RELAY_OP_RELEASE, 0)
|
|
|
return t.sendRelayCommand("release", 2, 0, frame)
|
|
return t.sendRelayCommand("release", 2, 0, frame)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -54,12 +55,12 @@ func (t *TCPReader) Connect() error {
|
|
|
addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port))
|
|
addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port))
|
|
|
conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
|
|
conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
|
|
|
if err != nil {
|
|
if err != nil {
|
|
|
- t.isConnected = false
|
|
|
|
|
|
|
+ t.connected.Store(false)
|
|
|
t.logWarn("UHF TCP 连接失败", zap.Error(err))
|
|
t.logWarn("UHF TCP 连接失败", zap.Error(err))
|
|
|
return err
|
|
return err
|
|
|
}
|
|
}
|
|
|
t.conn = conn
|
|
t.conn = conn
|
|
|
- t.isConnected = true
|
|
|
|
|
|
|
+ t.connected.Store(true)
|
|
|
t.logInfo("UHF TCP 连接成功")
|
|
t.logInfo("UHF TCP 连接成功")
|
|
|
return nil
|
|
return nil
|
|
|
}
|
|
}
|
|
@@ -70,16 +71,14 @@ func (t *TCPReader) Disconnect() error {
|
|
|
if t.conn != nil {
|
|
if t.conn != nil {
|
|
|
_ = t.conn.Close()
|
|
_ = t.conn.Close()
|
|
|
}
|
|
}
|
|
|
- t.isConnected = false
|
|
|
|
|
|
|
+ t.connected.Store(false)
|
|
|
t.buffer = nil
|
|
t.buffer = nil
|
|
|
t.logInfo("UHF TCP 连接已关闭")
|
|
t.logInfo("UHF TCP 连接已关闭")
|
|
|
return nil
|
|
return nil
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) IsConnected() bool {
|
|
func (t *TCPReader) IsConnected() bool {
|
|
|
- t.ioMu.Lock()
|
|
|
|
|
- defer t.ioMu.Unlock()
|
|
|
|
|
- return t.isConnected
|
|
|
|
|
|
|
+ return t.connected.Load()
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) SendData(data []byte) error {
|
|
func (t *TCPReader) SendData(data []byte) error {
|
|
@@ -89,19 +88,19 @@ func (t *TCPReader) SendData(data []byte) error {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) sendDataLocked(data []byte) error {
|
|
func (t *TCPReader) sendDataLocked(data []byte) error {
|
|
|
- if !t.isConnected {
|
|
|
|
|
|
|
+ if !t.connected.Load() {
|
|
|
return errors.New("TCP未连接")
|
|
return errors.New("TCP未连接")
|
|
|
}
|
|
}
|
|
|
_, err := t.conn.Write(data)
|
|
_, err := t.conn.Write(data)
|
|
|
if err != nil {
|
|
if err != nil {
|
|
|
- t.isConnected = false
|
|
|
|
|
|
|
+ t.connected.Store(false)
|
|
|
t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err))
|
|
t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err))
|
|
|
}
|
|
}
|
|
|
return err
|
|
return err
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// ==============================
|
|
// ==============================
|
|
|
-// 🔥 关键修复:TCP 自动拼包,返回完整25字节帧
|
|
|
|
|
|
|
+// 🔥 关键修复:TCP 自动拼包,按协议变长帧提取(标签帧25字节,继电器应答9字节等)
|
|
|
// ==============================
|
|
// ==============================
|
|
|
func (t *TCPReader) ReadData() ([]byte, error) {
|
|
func (t *TCPReader) ReadData() ([]byte, error) {
|
|
|
t.ioMu.Lock()
|
|
t.ioMu.Lock()
|
|
@@ -110,7 +109,7 @@ func (t *TCPReader) ReadData() ([]byte, error) {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) readDataLocked() ([]byte, error) {
|
|
func (t *TCPReader) readDataLocked() ([]byte, error) {
|
|
|
- if !t.isConnected {
|
|
|
|
|
|
|
+ if !t.connected.Load() {
|
|
|
return nil, errors.New("TCP未连接")
|
|
return nil, errors.New("TCP未连接")
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -123,31 +122,28 @@ func (t *TCPReader) readDataLocked() ([]byte, error) {
|
|
|
if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
|
|
if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
|
|
|
return nil, ErrReadTimeout
|
|
return nil, ErrReadTimeout
|
|
|
}
|
|
}
|
|
|
- t.isConnected = false
|
|
|
|
|
|
|
+ t.connected.Store(false)
|
|
|
t.logWarn("UHF TCP 读取失败", zap.Error(err))
|
|
t.logWarn("UHF TCP 读取失败", zap.Error(err))
|
|
|
return nil, err
|
|
return nil, err
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// 加入缓冲区
|
|
// 加入缓冲区
|
|
|
t.buffer = append(t.buffer, buf[:n]...)
|
|
t.buffer = append(t.buffer, buf[:n]...)
|
|
|
- t.logInfo("UHF TCP 收到数据块",
|
|
|
|
|
|
|
+ t.logDebug("UHF TCP 收到数据块",
|
|
|
zap.Int("packet_bytes", n),
|
|
zap.Int("packet_bytes", n),
|
|
|
zap.Int("buffered_bytes", len(t.buffer)),
|
|
zap.Int("buffered_bytes", len(t.buffer)),
|
|
|
zap.String("packet_hex", fmt.Sprintf("%x", buf[:n])),
|
|
zap.String("packet_hex", fmt.Sprintf("%x", 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:]
|
|
|
|
|
|
|
+ frame, rest, ok := extractUHFFrame(t.buffer)
|
|
|
|
|
+ t.buffer = rest
|
|
|
|
|
+ if !ok {
|
|
|
|
|
+ if len(t.buffer) > 0 {
|
|
|
|
|
+ return nil, ErrIncompleteFrame
|
|
|
}
|
|
}
|
|
|
|
|
+ return nil, ErrReadTimeout
|
|
|
}
|
|
}
|
|
|
-
|
|
|
|
|
- return nil, ErrIncompleteFrame
|
|
|
|
|
|
|
+ return frame, nil
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
|
|
func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
|
|
@@ -160,6 +156,12 @@ func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
|
|
|
return t.readDataLocked()
|
|
return t.readDataLocked()
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+// isRelayAckFor 判断响应帧是否为对应指令的应答(设备回显命令字)。
|
|
|
|
|
+func isRelayAckFor(response, frame []byte) bool {
|
|
|
|
|
+ return len(response) >= 7 && response[0] == 0xCF &&
|
|
|
|
|
+ response[2] == frame[2] && response[3] == frame[3]
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
func (t *TCPReader) sendRelayCommand(action string, relay int, validTime byte, frame []byte) error {
|
|
func (t *TCPReader) sendRelayCommand(action string, relay int, validTime byte, frame []byte) error {
|
|
|
t.logInfo("UHF TCP 道闸指令下发",
|
|
t.logInfo("UHF TCP 道闸指令下发",
|
|
|
zap.String("action", action),
|
|
zap.String("action", action),
|
|
@@ -167,17 +169,39 @@ func (t *TCPReader) sendRelayCommand(action string, relay int, validTime byte, f
|
|
|
zap.Uint8("valid_time_seconds", validTime),
|
|
zap.Uint8("valid_time_seconds", validTime),
|
|
|
zap.String("packet_hex", fmt.Sprintf("%x", frame)),
|
|
zap.String("packet_hex", fmt.Sprintf("%x", frame)),
|
|
|
)
|
|
)
|
|
|
- response, err := t.SendAndRecv(frame)
|
|
|
|
|
- if err != nil {
|
|
|
|
|
|
|
+ t.ioMu.Lock()
|
|
|
|
|
+ defer t.ioMu.Unlock()
|
|
|
|
|
+ if err := t.sendDataLocked(frame); err != nil {
|
|
|
t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err))
|
|
t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err))
|
|
|
return err
|
|
return err
|
|
|
}
|
|
}
|
|
|
- t.logInfo("UHF TCP 道闸指令收到响应",
|
|
|
|
|
- zap.String("action", action),
|
|
|
|
|
- zap.Int("response_bytes", len(response)),
|
|
|
|
|
- zap.String("response_hex", fmt.Sprintf("%x", response)),
|
|
|
|
|
- )
|
|
|
|
|
- return nil
|
|
|
|
|
|
|
+ // 等待设备应答(回显命令字)。期间收到的标签上报帧放回缓冲区交还读循环。
|
|
|
|
|
+ for attempt := 0; attempt < 3; attempt++ {
|
|
|
|
|
+ response, err := t.readDataLocked()
|
|
|
|
|
+ if err != nil {
|
|
|
|
|
+ break
|
|
|
|
|
+ }
|
|
|
|
|
+ if isRelayAckFor(response, frame) {
|
|
|
|
|
+ status := "unknown"
|
|
|
|
|
+ if len(response) >= 7 {
|
|
|
|
|
+ status = fmt.Sprintf("0x%02x", response[5])
|
|
|
|
|
+ }
|
|
|
|
|
+ t.logInfo("UHF TCP 道闸指令收到响应",
|
|
|
|
|
+ zap.String("action", action),
|
|
|
|
|
+ zap.Int("relay", relay),
|
|
|
|
|
+ zap.String("status", status),
|
|
|
|
|
+ zap.String("response_hex", fmt.Sprintf("%x", response)),
|
|
|
|
|
+ )
|
|
|
|
|
+ if len(response) >= 7 && response[5] != 0x00 {
|
|
|
|
|
+ return fmt.Errorf("道闸指令被设备拒绝: 状态码 %s", status)
|
|
|
|
|
+ }
|
|
|
|
|
+ return nil
|
|
|
|
|
+ }
|
|
|
|
|
+ t.buffer = append(append(make([]byte, 0, len(response)+len(t.buffer)), response...), t.buffer...)
|
|
|
|
|
+ }
|
|
|
|
|
+ err := fmt.Errorf("道闸指令未收到设备应答(帧已发出)")
|
|
|
|
|
+ t.logWarn("UHF TCP 道闸指令失败", zap.String("action", action), zap.Error(err))
|
|
|
|
|
+ return err
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
func (t *TCPReader) logInfo(message string, fields ...zap.Field) {
|
|
func (t *TCPReader) logInfo(message string, fields ...zap.Field) {
|
|
@@ -188,6 +212,14 @@ func (t *TCPReader) logInfo(message string, fields ...zap.Field) {
|
|
|
global.GVA_LOG.Info(message, fields...)
|
|
global.GVA_LOG.Info(message, fields...)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+func (t *TCPReader) logDebug(message string, fields ...zap.Field) {
|
|
|
|
|
+ if global.GVA_LOG == nil {
|
|
|
|
|
+ return
|
|
|
|
|
+ }
|
|
|
|
|
+ fields = append([]zap.Field{zap.String("tcp_address", net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port)))}, fields...)
|
|
|
|
|
+ global.GVA_LOG.Debug(message, fields...)
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
func (t *TCPReader) logWarn(message string, fields ...zap.Field) {
|
|
func (t *TCPReader) logWarn(message string, fields ...zap.Field) {
|
|
|
if global.GVA_LOG == nil {
|
|
if global.GVA_LOG == nil {
|
|
|
return
|
|
return
|