64 lines
1.6 KiB
Go
64 lines
1.6 KiB
Go
package agent
|
|
|
|
import (
|
|
"io"
|
|
"net"
|
|
"time"
|
|
)
|
|
|
|
// handleData 处理节点主动打开的数据流:
|
|
// 首帧 CmdTarget 携带目标信息 → 拨号本地 TCP 目标 → 回执 → 双向管道。
|
|
func (a *Agent) handleData(stream net.Conn) {
|
|
defer stream.Close()
|
|
f, err := ReadFrame(stream)
|
|
if err != nil {
|
|
return
|
|
}
|
|
if f.Cmd() != CmdTarget {
|
|
a.logger().Warnf("agent 数据流首帧命令错误 cmd=%d", f.Cmd())
|
|
return
|
|
}
|
|
t := Tunnel{}
|
|
if err := f.Unmarshal(&t); err != nil {
|
|
return
|
|
}
|
|
|
|
dialTimeout := a.DialTimeout
|
|
if dialTimeout <= 0 { // 零值 Agent 兜底
|
|
dialTimeout = defaultDialTimeout
|
|
}
|
|
conn, err := net.DialTimeout("tcp", t.TargetAddr(), dialTimeout)
|
|
if err != nil {
|
|
a.logger().Warnf("agent 拨号本地目标失败 target=%s err=%s", t.TargetAddr(), err.Error())
|
|
rf := NewFrame(FrameVersion, CmdTarget, f.StreamID())
|
|
_ = rf.Marshal(Ret{Code: 1, Msg: err.Error()})
|
|
_ = WriteFrame(stream, rf)
|
|
return
|
|
}
|
|
defer conn.Close()
|
|
// agent↔目标段 TCP 保活: 目标假死(进程在但不响应/不发 FIN)时是
|
|
// 全链路唯一没有存活检测的段, 不开 keepalive 观看端会无限挂死画面
|
|
if tc, ok := conn.(*net.TCPConn); ok {
|
|
tc.SetKeepAlive(true)
|
|
tc.SetKeepAlivePeriod(5 * time.Second)
|
|
}
|
|
|
|
rf := NewFrame(FrameVersion, CmdTarget, f.StreamID())
|
|
_ = rf.Marshal(Ret{Code: 0, Msg: "ok"})
|
|
if err := WriteFrame(stream, rf); err != nil {
|
|
return
|
|
}
|
|
|
|
pipe(conn, stream)
|
|
}
|
|
|
|
// pipe 双向拷贝,任一方向结束即关闭两端
|
|
func pipe(a, b net.Conn) {
|
|
go func() {
|
|
io.Copy(b, a)
|
|
b.Close()
|
|
}()
|
|
io.Copy(a, b)
|
|
a.Close()
|
|
}
|