Files

115 lines
2.7 KiB
Go

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()
}