package node import ( "io" "log/slog" "net" "sync/atomic" "time" "git.zeroonesoft.cn/golib/zonat/agent" ) // TunnelServer 一条已注册隧道的公网侧监听。 // 纯 TCP 字节管道:无协议嗅探、无 HTTP 代理、无证书,任何 TCP 协议原样通过。 type TunnelServer struct { tunnel agent.Tunnel session *Session listener net.Listener connCount atomic.Int64 lastActive atomic.Int64 // unix nano closed atomic.Bool } // Port 公网监听端口 func (t *TunnelServer) Port() int { return t.tunnel.ListenPort } func (t *TunnelServer) touch() { t.lastActive.Store(time.Now().UnixNano()) } func (t *TunnelServer) lastActiveTime() time.Time { return time.Unix(0, t.lastActive.Load()) } func (t *TunnelServer) acceptLoop() { for { conn, err := t.listener.Accept() if err != nil { return } if t.closed.Load() { conn.Close() return } go t.handleConn(conn) } } // handleConn 公网连接 → smux 新流 → 首帧告知目标 → agent 回执后双向管道。 // 与 zonat 多路复用路径一致,仅去掉协议嗅探与限流。 func (t *TunnelServer) handleConn(conn net.Conn) { defer conn.Close() if t.session.isDead() { return } stream, err := t.session.OpenStream() if err != nil { return } defer stream.Close() // 转发握手(告知目标+等 agent 回执)限时: agent 会话半死/目标拨号不回时 // 观看端不再无限挂起, 到点立即断开让观看端快速失败可重试 if hs := t.session.node.HandshakeTimeout; hs > 0 { _ = stream.SetDeadline(time.Now().Add(hs)) } tf := agent.NewFrame(agent.FrameVersion, agent.CmdTarget, 0) if err := tf.Marshal(t.tunnel); err != nil { return } if err := agent.WriteFrame(stream, tf); err != nil { return } rf, err := agent.ReadFrame(stream) if err != nil { slog.Warn("节点 agent拨号应答失败", "id", t.tunnel.Id, "err", err.Error()) return } _ = stream.SetDeadline(time.Time{}) // 数据阶段不受握手超时约束 ret := agent.Ret{} if err := rf.Unmarshal(&ret); err != nil || ret.Code != 0 { slog.Warn("节点 agent拨号失败", "id", t.tunnel.Id, "msg", ret.Msg) return } // 公网侧 TCP 保活(smux 流自身由会话级 keepalive 保护) if tc, ok := conn.(*net.TCPConn); ok { tc.SetKeepAlive(true) tc.SetKeepAlivePeriod(5 * time.Second) } t.connCount.Add(1) t.touch() defer func() { t.connCount.Add(-1) t.touch() }() pipe(conn, stream) } func (t *TunnelServer) close() { if t.closed.CompareAndSwap(false, true) && t.listener != nil { t.listener.Close() } } // pipe 双向拷贝,任一方向结束即关闭两端 func pipe(a, b net.Conn) { go func() { io.Copy(b, a) b.Close() }() io.Copy(a, b) a.Close() }