修复: SFTP连接首拉不显示与传输挂起
- watcher握手期伪触发致幻影本地列表,根因修复 - 连接探测2.5s超时竞速,握手10s预算 - 30s保活自动掐死半开连接,取连接时剔除已断开 - 副本命名与进度计数抽公共,下载缓存清理去重
This commit is contained in:
+100
-4
@@ -11,11 +11,22 @@ import (
|
||||
"golang.org/x/crypto/ssh"
|
||||
)
|
||||
|
||||
// 连接建立全流程(拨号+认证+SFTP 子系统)的总超时预算
|
||||
const connectTimeout = 10 * time.Second
|
||||
|
||||
// 保活参数:周期发送 keepalive 并等待回复,超过 keepaliveReplyWait 未回复判定连接已死
|
||||
const (
|
||||
keepaliveInterval = 30 * time.Second
|
||||
keepaliveReplyWait = 30 * time.Second
|
||||
)
|
||||
|
||||
// Client SFTP 客户端封装(单连接)
|
||||
type Client struct {
|
||||
config *Config
|
||||
client *sftp.Client
|
||||
sshClient *ssh.Client
|
||||
stopKeep chan struct{} // 当前 SSH 连接保活循环的停止信号
|
||||
closed bool // 保活失败或显式关闭后置位,供取连接时剔除
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
@@ -73,6 +84,11 @@ func (m *Manager) Disconnect(host string, port int) {
|
||||
}
|
||||
}
|
||||
|
||||
// Evict 从连接池剔除连接(底层已断开,仅移除不重复关闭)
|
||||
func (m *Manager) Evict(connID string) {
|
||||
m.clients.Delete(connID)
|
||||
}
|
||||
|
||||
// Shutdown 关闭所有连接
|
||||
func (m *Manager) Shutdown() {
|
||||
m.clients.Range(func(key, value any) bool {
|
||||
@@ -84,7 +100,18 @@ func (m *Manager) Shutdown() {
|
||||
|
||||
// --- 内部 ---
|
||||
|
||||
// newClient 建立连接并启动保活循环
|
||||
func newClient(config *Config) (*Client, error) {
|
||||
c, err := buildClient(config)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
c.startKeepalive()
|
||||
return c, nil
|
||||
}
|
||||
|
||||
// buildClient 拨号+认证+SFTP 子系统建立(全程受 connectTimeout 约束),不启动保活
|
||||
func buildClient(config *Config) (*Client, error) {
|
||||
sshConfig := &ssh.ClientConfig{
|
||||
Config: ssh.Config{
|
||||
KeyExchanges: []string{
|
||||
@@ -115,8 +142,8 @@ func newClient(config *Config) (*Client, error) {
|
||||
}
|
||||
sshConfig.Auth = []ssh.AuthMethod{ssh.PublicKeys(signer)}
|
||||
} else if config.Password != "" {
|
||||
pw := config.Password
|
||||
sshConfig.Auth = []ssh.AuthMethod{
|
||||
pw := config.Password
|
||||
sshConfig.Auth = []ssh.AuthMethod{
|
||||
ssh.Password(pw),
|
||||
ssh.KeyboardInteractive(func(user, instruction string, questions []string, echos []bool) ([]string, error) {
|
||||
answers := make([]string, len(questions))
|
||||
@@ -131,10 +158,13 @@ func newClient(config *Config) (*Client, error) {
|
||||
}
|
||||
|
||||
addr := fmt.Sprintf("%s:%d", config.Host, config.Port)
|
||||
sshConn, err := net.DialTimeout("tcp", addr, config.Timeout)
|
||||
// 拨号与认证、SFTP 子系统握手共享同一总超时预算,认证阶段无界挂起由 deadline 掐断
|
||||
deadline := time.Now().Add(connectTimeout)
|
||||
sshConn, err := net.DialTimeout("tcp", addr, connectTimeout)
|
||||
if err != nil {
|
||||
return nil, &ConnectionError{Op: "dial", Err: err}
|
||||
}
|
||||
sshConn.SetDeadline(deadline)
|
||||
|
||||
sshConnConn, chans, reqs, err := ssh.NewClientConn(sshConn, addr, sshConfig)
|
||||
if err != nil {
|
||||
@@ -143,7 +173,9 @@ func newClient(config *Config) (*Client, error) {
|
||||
}
|
||||
|
||||
sshClient := ssh.NewClient(sshConnConn, chans, reqs)
|
||||
|
||||
sftpClient, err := sftp.NewClient(sshClient)
|
||||
sshConn.SetDeadline(time.Time{})
|
||||
if err != nil {
|
||||
sshClient.Close()
|
||||
return nil, &ConnectionError{Op: "sftp_init", Err: err}
|
||||
@@ -156,6 +188,62 @@ func newClient(config *Config) (*Client, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
// startKeepalive 启动当前 SSH 连接的保活循环(须持有 c.mu 或独占 Client 时调用)
|
||||
// 防止空闲连接被 NAT/防火墙杀掉;探测失败时关闭连接并置位 closed,唤醒所有阻塞中的操作
|
||||
func (c *Client) startKeepalive() {
|
||||
stop := make(chan struct{})
|
||||
sshClient := c.sshClient
|
||||
c.stopKeep = stop
|
||||
go func() {
|
||||
t := time.NewTicker(keepaliveInterval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-stop:
|
||||
return
|
||||
case <-t.C:
|
||||
}
|
||||
if !sshKeepalive(sshClient, keepaliveReplyWait) {
|
||||
c.markDead(sshClient)
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// sshKeepalive 发送一次 keepalive 请求并限时等待回复,超时或出错视作连接已死
|
||||
func sshKeepalive(sshClient *ssh.Client, wait time.Duration) bool {
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
_, _, err := sshClient.SendRequest("keepalive@openssh.com", true, nil)
|
||||
done <- err
|
||||
}()
|
||||
select {
|
||||
case err := <-done:
|
||||
return err == nil
|
||||
case <-time.After(wait):
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// markDead 标记连接死亡并关闭底层资源
|
||||
// 仅当容器仍持有该 SSH 连接时生效,避免误杀 reconnect 换上的新连接
|
||||
func (c *Client) markDead(sshClient *ssh.Client) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
if c.sshClient != sshClient {
|
||||
return
|
||||
}
|
||||
c.closeLocked()
|
||||
}
|
||||
|
||||
// IsClosed 连接是否已断开(保活失败或显式关闭)
|
||||
func (c *Client) IsClosed() bool {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
return c.closed
|
||||
}
|
||||
|
||||
// IsHealthy 检查连接是否健康(先取引用再解锁,避免持锁做 I/O)
|
||||
func (c *Client) IsHealthy() bool {
|
||||
c.mu.Lock()
|
||||
@@ -204,7 +292,8 @@ func (c *Client) WithRetry(fn func(*sftp.Client) error) error {
|
||||
}
|
||||
|
||||
func (c *Client) reconnect() error {
|
||||
nc, err := newClient(c.config)
|
||||
// 用 buildClient 而非 newClient,避免保活循环绑定到即将丢弃的临时容器
|
||||
nc, err := buildClient(c.config)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -213,6 +302,8 @@ func (c *Client) reconnect() error {
|
||||
c.closeLocked()
|
||||
c.client = nc.client
|
||||
c.sshClient = nc.sshClient
|
||||
c.closed = false
|
||||
c.startKeepalive()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -223,6 +314,10 @@ func (c *Client) Close() {
|
||||
}
|
||||
|
||||
func (c *Client) closeLocked() {
|
||||
if c.stopKeep != nil {
|
||||
close(c.stopKeep)
|
||||
c.stopKeep = nil
|
||||
}
|
||||
if c.client != nil {
|
||||
c.client.Close()
|
||||
c.client = nil
|
||||
@@ -231,6 +326,7 @@ func (c *Client) closeLocked() {
|
||||
c.sshClient.Close()
|
||||
c.sshClient = nil
|
||||
}
|
||||
c.closed = true
|
||||
}
|
||||
|
||||
// RunCommand 通过 SSH Session 执行远程命令,返回 stdout
|
||||
|
||||
Reference in New Issue
Block a user