| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229 |
- package uhf
- import (
- "errors"
- "fmt"
- "net"
- "sync"
- "sync/atomic"
- "time"
- "go.uber.org/zap"
- "wails-app/internal/global"
- )
- type TCPReader struct {
- IP string
- Port int
- conn net.Conn
- connected atomic.Bool
- buffer []byte // 内部拼包缓冲区(关键修复)
- ioMu sync.Mutex
- }
- // 实现 Reader 接口的继电器方法(不再panic)
- func (t *TCPReader) CloseRelay1(validTime byte) error {
- frame := buildRelayFrame(0x0077, RELAY_OP_CLOSE, validTime)
- return t.sendRelayCommand("open", 1, validTime, frame)
- }
- func (t *TCPReader) ReleaseRelay1() error {
- frame := buildRelayFrame(0x0077, RELAY_OP_RELEASE, 0)
- return t.sendRelayCommand("release", 1, 0, frame)
- }
- func (t *TCPReader) CloseRelay2(validTime byte) error {
- frame := buildRelayFrame(0x0078, RELAY_OP_CLOSE, validTime)
- return t.sendRelayCommand("close", 2, validTime, frame)
- }
- func (t *TCPReader) ReleaseRelay2() error {
- frame := buildRelayFrame(0x0078, RELAY_OP_RELEASE, 0)
- return t.sendRelayCommand("release", 2, 0, frame)
- }
- func NewTCPReader(ip string, port int) *TCPReader {
- return &TCPReader{
- IP: ip,
- Port: port,
- }
- }
- func (t *TCPReader) Connect() error {
- t.ioMu.Lock()
- defer t.ioMu.Unlock()
- addr := net.JoinHostPort(t.IP, fmt.Sprintf("%d", t.Port))
- conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
- if err != nil {
- t.connected.Store(false)
- t.logWarn("UHF TCP 连接失败", zap.Error(err))
- return err
- }
- t.conn = conn
- t.connected.Store(true)
- t.logInfo("UHF TCP 连接成功")
- return nil
- }
- func (t *TCPReader) Disconnect() error {
- t.ioMu.Lock()
- defer t.ioMu.Unlock()
- if t.conn != nil {
- _ = t.conn.Close()
- }
- t.connected.Store(false)
- t.buffer = nil
- t.logInfo("UHF TCP 连接已关闭")
- return nil
- }
- func (t *TCPReader) IsConnected() bool {
- return t.connected.Load()
- }
- func (t *TCPReader) SendData(data []byte) error {
- t.ioMu.Lock()
- defer t.ioMu.Unlock()
- return t.sendDataLocked(data)
- }
- func (t *TCPReader) sendDataLocked(data []byte) error {
- if !t.connected.Load() {
- return errors.New("TCP未连接")
- }
- _, err := t.conn.Write(data)
- if err != nil {
- t.connected.Store(false)
- t.logWarn("UHF TCP 写入失败", zap.String("packet_hex", fmt.Sprintf("%x", data)), zap.Error(err))
- }
- return err
- }
- // ==============================
- // 🔥 关键修复:TCP 自动拼包,按协议变长帧提取(标签帧25字节,继电器应答9字节等)
- // ==============================
- func (t *TCPReader) ReadData() ([]byte, error) {
- t.ioMu.Lock()
- defer t.ioMu.Unlock()
- return t.readDataLocked()
- }
- func (t *TCPReader) readDataLocked() ([]byte, error) {
- if !t.connected.Load() {
- return nil, errors.New("TCP未连接")
- }
- // 读超时,防止永久阻塞
- t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
- buf := make([]byte, 1024)
- n, err := t.conn.Read(buf)
- if err != nil {
- if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
- return nil, ErrReadTimeout
- }
- t.connected.Store(false)
- t.logWarn("UHF TCP 读取失败", zap.Error(err))
- return nil, err
- }
- // 加入缓冲区
- t.buffer = append(t.buffer, buf[:n]...)
- t.logDebug("UHF TCP 收到数据块",
- zap.Int("packet_bytes", n),
- zap.Int("buffered_bytes", len(t.buffer)),
- zap.String("packet_hex", fmt.Sprintf("%x", buf[:n])),
- )
- frame, rest, ok := extractUHFFrame(t.buffer)
- t.buffer = rest
- if !ok {
- if len(t.buffer) > 0 {
- return nil, ErrIncompleteFrame
- }
- return nil, ErrReadTimeout
- }
- return frame, nil
- }
- func (t *TCPReader) SendAndRecv(data []byte) ([]byte, error) {
- t.ioMu.Lock()
- defer t.ioMu.Unlock()
- if err := t.sendDataLocked(data); err != nil {
- return nil, err
- }
- t.conn.SetReadDeadline(time.Now().Add(3 * time.Second))
- 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 {
- t.logInfo("UHF TCP 道闸指令下发",
- zap.String("action", action),
- zap.Int("relay", relay),
- zap.Uint8("valid_time_seconds", validTime),
- zap.String("packet_hex", fmt.Sprintf("%x", frame)),
- )
- 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))
- return err
- }
- // 等待设备应答(回显命令字)。期间收到的标签上报帧放回缓冲区交还读循环。
- 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) {
- 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.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) {
- 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.Warn(message, fields...)
- }
|