You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Golang获取远程SSH执行命令的PID为空,求解决方法

Golang SSH会话获取远程命令PID为空的解决方法及Ansible异步更新逻辑迁移

我尝试通过Golang的SSH会话获取远程服务器上执行命令的PID,但返回的PID为空,同时正在将Ansible的集群异步更新逻辑迁移到Golang。

原Golang代码

func (c *SSHClient) RunCommandWithContext(ctx context.Context, command string) (string, error) {
    if c.Client == nil {
        return "", errors.New("SSH client is not connected")
    }

    session, err := c.Client.NewSession()
    if err != nil {
        return "", err
    }
    defer session.Close()

    var stdoutBuf bytes.Buffer
    session.Stdout = &stdoutBuf

    // Channel to signal completion
    done := make(chan error, 1)

    // Run the command and get the PID
    var stdout, stderr bytes.Buffer
    session.Stdout = &stdout
    session.Stderr = &stderr

    go func() {
        done <- session.Run(command)
    }()

    // Extract the PID from the output
    output := strings.TrimSpace(stdout.String())
    fmt.Println("Output:", output) // 这里输出为空
    pid, err := strconv.Atoi(output)
    fmt.Println("PID:", pid) // 需要获取远程执行命令的PID
    fmt.Println("Error:", err)
    if err != nil {
        log.Fatalf("Failed to parse PID: %v", err)
    }
    log.Printf("Started command with PID %d\n", pid)

    ticker := time.NewTicker(2 * time.Second)
    defer ticker.Stop()
    for {
        select {
        case <-ctx.Done():
            // Context timeout
            if err := session.Signal(ssh.SIGKILL); err != nil {
                return "", fmt.Errorf("failed to kill process: %w", err)
            }
            return "", ctx.Err()
        case <-ticker.C:
            log.Println("Checking command status...")
            if !isProcessRunning(session, pid) {
                log.Println("Command finished successfully")
                printLogs(stdout.String(), stderr.String())
                return "", nil
            }
        case err := <-done:
            // Command completed
            if err != nil {
                return "", err
            }
            return stdoutBuf.String(), nil
        }
    }
    return "", nil
}

// This program asynchronously checks if the command has completed or not
func isProcessRunning(session *ssh.Session, pid int) bool { 
    command := fmt.Sprintf("ps -p %d", pid)
    if err := session.Run(command); err != nil {
        return false
    }
    // Run the command and get the PID
    var stdout, stderr bytes.Buffer
    session.Stdout = &stdout
    session.Stderr = &stderr
    output := strings.TrimSpace(stdout.String())
    return strings.Contains(output, strconv.Itoa(pid))
}

原Ansible异步更新逻辑

- name: Update cluster
  shell:
    cmd: "{{ bmctl_tool }} update cluster -c {{ cluster }} --workspace-dir {{ abm_path }} --kubeconfig {{ cluster_kubeconfig }} --bootstrap-cluster-pod-cidr {{ bootstrap_cluster_pod_cidr }} --bootstrap-cluster-service-cidr {{ bootstrap_cluster_service_cidr }}"
  environment:
    GOOGLE_APPLICATION_CREDENTIALS: "{{ gcp_service_account }}"
  async: 2000
  poll: 0
  register: cluster_update

- name: Check on cluster update async task
  async_status:
    jid: "{{ cluster_update.ansible_job_id }}"
  register: update_job_result
  until: update_job_result.finished
  no_log: true
  retries: 100
  delay: 20

问题分析

  1. PID为空的核心原因:
    • 启动session.Run的goroutine后立刻读取stdout,此时命令还未执行或未输出任何内容,导致stdout为空。
    • 原执行的命令本身不会输出PID,即使等待命令执行,也无法获取PID。
  2. isProcessRunning函数的问题:复用同一个SSH Session执行ps命令,且未正确设置输出缓冲,导致无法正确判断进程状态。
  3. Ansible异步逻辑的对应关系:Ansible的async:2000 poll:0是让命令后台运行并返回任务ID,后续定期检查状态;Golang需要模拟这一逻辑,让命令后台运行、记录PID、定期检查进程状态。

解决方案

1. 修正Golang代码,获取远程命令PID并实现异步检查

修改远程执行的命令,让它后台运行并输出自身PID;同时重构进程状态检查逻辑,避免复用Session导致的问题:

import (
    "bytes"
    "errors"
    "fmt"
    "log"
    "strconv"
    "strings"
    "time"

    "golang.org/x/crypto/ssh"
)

func (c *SSHClient) RunCommandWithContext(ctx context.Context, command string) (string, error) {
    if c.Client == nil {
        return "", errors.New("SSH client is not connected")
    }

    // 1. 包装命令:后台运行,输出PID到stdout,同时重定向日志到文件
    logFile := "/tmp/cluster_update.log"
    wrappedCmd := fmt.Sprintf("nohup %s > %s 2>&1 & echo $!", command, logFile)

    // 2. 创建Session执行包装后的命令,获取PID
    pidSession, err := c.Client.NewSession()
    if err != nil {
        return "", fmt.Errorf("failed to create PID session: %w", err)
    }
    defer pidSession.Close()

    var stdout, stderr bytes.Buffer
    pidSession.Stdout = &stdout
    pidSession.Stderr = &stderr

    if err := pidSession.Run(wrappedCmd); err != nil {
        return "", fmt.Errorf("failed to start command: %w, stderr: %s", err, stderr.String())
    }

    // 3. 解析PID
    pidStr := strings.TrimSpace(stdout.String())
    pid, err := strconv.Atoi(pidStr)
    if err != nil {
        return "", fmt.Errorf("invalid PID output: %s, error: %w", pidStr, err)
    }
    log.Printf("Started cluster update with PID %d\n", pid)

    // 4. 定期检查进程状态(模拟Ansible的delay:20)
    ticker := time.NewTicker(20 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            // 超时或取消,强制杀死进程
            killSession, err := c.Client.NewSession()
            if err != nil {
                return "", fmt.Errorf("failed to create kill session: %w", err)
            }
            defer killSession.Close()

            if err := killSession.Run(fmt.Sprintf("kill -9 %d", pid)); err != nil {
                return "", fmt.Errorf("failed to kill process %d: %w", pid, err)
            }
            return "", ctx.Err()
        case <-ticker.C:
            log.Println("Checking cluster update status...")
            running, err := c.isProcessRunning(pid)
            if err != nil {
                log.Printf("Warning: failed to check process status: %v", err)
                continue
            }
            if !running {
                // 进程结束,读取日志并返回
                logSession, err := c.Client.NewSession()
                if err != nil {
                    return "", fmt.Errorf("failed to create log session: %w", err)
                }
                defer logSession.Close()

                var logOutput bytes.Buffer
                logSession.Stdout = &logOutput
                if err := logSession.Run(fmt.Sprintf("cat %s", logFile)); err != nil {
                    return "", fmt.Errorf("failed to read command log: %w", err)
                }
                log.Println("Cluster update completed successfully")
                return logOutput.String(), nil
            }
        }
    }
}

// 单独创建Session检查进程状态,避免复用Session的问题
func (c *SSHClient) isProcessRunning(pid int) (bool, error) {
    session, err := c.Client.NewSession()
    if err != nil {
        return false, err
    }
    defer session.Close()

    // ps命令退出码为0表示进程存在,非0表示不存在
    cmd := fmt.Sprintf("ps -p %d > /dev/null 2>&1", pid)
    err = session.Run(cmd)
    return err == nil, nil
}

2. Ansible异步逻辑迁移说明

  • 后台执行命令:用nohup ... & echo $!实现命令后台运行,同时输出PID(对应Ansible的async:2000 poll:0)。
  • 状态检查:每隔20秒检查一次PID对应的进程是否存在(对应Ansible的delay:20和until: update_job_result.finished)。
  • 超时处理:用context.WithTimeout实现2000秒的超时(对应Ansible的async:2000)。
  • 日志记录:将命令输出重定向到远程服务器的日志文件,方便进程结束后读取完整输出(对应Ansible任务的输出记录)。

内容的提问来源于stack exchange,提问作者Rad4

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 18:14:55