package stress import ( "fmt" "net" "os" "strconv" "strings" "time" "auto-check/pkg/config" "auto-check/pkg/discovery" "auto-check/pkg/sshclient" "golang.org/x/crypto/ssh" ) // Results 测试结果表(IP → 报告),包级公开 var Results = make(map[string]*Report) // ============================ // 执行器 // ============================ // Runner 压力测试执行器(统一在 Docker 容器中执行,工具预装在镜像里) type Runner struct { cfg Config client *ssh.Client metrics []Metric sudoPwd string // 非 root 用户(如 fnOS admin)执行 docker 命令时的 sudo 密码 dockerOK bool // 目标机 Docker 是否可用 } // NewRunner 创建执行器(检测 Docker → 确保镜像 → 启动容器执行测试) func NewRunner(cfg Config, types []TestType, client *ssh.Client, sudoPwd string) *Runner { dockerInfo := IsDockerAvailable(client, sudoPwd) if dockerInfo.Available { fmt.Printf(" [Docker] 检测到 Docker 版本 %s,压测将在容器中执行\n", dockerInfo.Version) } else { fmt.Printf(" [Docker] 目标机未检测到 Docker,无法执行压测(飞牛默认已安装,请检查 Docker 服务状态)\n") } return &Runner{ cfg: cfg, client: client, metrics: BuildMetrics(types), sudoPwd: sudoPwd, dockerOK: dockerInfo.Available, } } // Run 执行压力测试(启动容器 → 解析结果 → 生成报告) func (r *Runner) Run(ip string) *Report { report := &Report{IP: ip, StartTime: time.Now()} fmt.Printf("\n════════ [%s] 压力测试开始 ════════\n", ip) // Docker 是唯一执行模式,不可用时直接判失败 if !r.dockerOK { report.AddResult(Result{Type: "docker", Status: StatusFail, Error: "目标机 Docker 不可用"}) report.EndTime = time.Now() report.Duration = report.EndTime.Sub(report.StartTime) return report } fmt.Printf(" [模式] Docker 容器 (%s)\n", r.cfg.DockerImage) // 显示将执行的指标及被跳过的原因(未知类型) fmt.Printf(" 指标: ") for _, m := range r.metrics { if m.Enabled { fmt.Printf("%s ", m.Name) } } fmt.Println() for _, m := range r.metrics { if !m.Enabled && !m.IsMonitor { fmt.Printf(" [跳过] %s: %s\n", m.Name, m.DisableReason) } } // 启动容器执行测试(脚本已打包在镜像中,环境变量传参) output, err := r.runInDocker() if err != nil { fmt.Printf(" [执行] 容器执行出错: %v\n", err) } // 解析容器输出(与宿主机模式输出格式相同) if output != "" { r.parseResults(report, output) } report.EndTime = time.Now() report.Duration = report.EndTime.Sub(report.StartTime) // 生成综合 HTML 报告(跑分卡片 + 曲线图) if len(report.Results) > 0 { reportPath, err := WriteReportHTML(ip, report, "reports") if err != nil { fmt.Printf(" [报告] 生成失败: %v\n", err) } else { fmt.Printf(" [报告] %s\n", reportPath) } } fmt.Printf("\n════════ [%s] 压力测试完成 ════════\n", ip) fmt.Println(report.ToText()) return report } // parseResults 解析脚本的结构化输出 func (r *Runner) parseResults(report *Report, output string) { lines := strings.Split(output, "\n") var currentTest string var currentStatus string var currentDuration string var currentOutput []string var currentSamples []Sample var currentScore float64 var currentMetrics map[string]float64 var sampleHeader []string inSamples := false inSysInfo := false flush := func() { if currentTest == "" { return } status := StatusPass switch currentStatus { case "fail": status = StatusFail case "skip": status = StatusSkip case "error": status = StatusError } // 温度超限判定(仅对 temp 监控) if TestType(currentTest) == MonitorTemp && status == StatusPass && r.cfg.TempLimit > 0 { if maxT, ok := currentMetrics["max_temp"]; ok && maxT > r.cfg.TempLimit { status = StatusFail currentOutput = append(currentOutput, fmt.Sprintf("温度超限: 最高 %.1f°C > 阈值 %.1f°C", maxT, r.cfg.TempLimit)) } } duration, _ := time.ParseDuration(currentDuration) report.AddResult(Result{ Type: TestType(currentTest), Status: status, Output: strings.Join(currentOutput, "\n"), Duration: duration, Samples: currentSamples, Score: currentScore, Metrics: currentMetrics, }) currentTest = "" currentStatus = "" currentDuration = "" currentOutput = nil currentSamples = nil currentScore = 0 currentMetrics = nil sampleHeader = nil } for _, line := range lines { line = strings.TrimSpace(line) // 采样数据块 if line == "===SAMPLES===" { inSamples = true sampleHeader = nil continue } if line == "===END_SAMPLES===" { inSamples = false continue } if inSamples { parts := strings.Split(line, ",") if sampleHeader == nil { // 首行为表头:time,指标1,指标2,... for _, p := range parts { sampleHeader = append(sampleHeader, strings.TrimSpace(p)) } continue } if len(parts) < 1 || strings.TrimSpace(parts[0]) == "" { continue } s := Sample{Time: strings.TrimSpace(parts[0]), Values: map[string]float64{}} for i := 1; i < len(parts) && i < len(sampleHeader); i++ { name := sampleHeader[i] if name == "" || name == "time" { continue } if v, err := strconv.ParseFloat(strings.TrimSpace(parts[i]), 64); err == nil { s.Values[name] = v } } currentSamples = append(currentSamples, s) continue } if line == "===SYSTEM_INFO===" { inSysInfo = true continue } if inSysInfo { if line == "" { inSysInfo = false continue } if i := strings.Index(line, ":"); i > 0 { k := strings.TrimSpace(line[:i]) v := strings.TrimSpace(line[i+1:]) if report.SysInfo == nil { report.SysInfo = map[string]string{} } report.SysInfo[k] = v } continue } if strings.HasPrefix(line, "===TEST:") { flush() currentTest = strings.TrimSuffix(strings.TrimPrefix(line, "===TEST:"), "===") continue } if strings.HasPrefix(line, "===MONITOR:") { flush() currentTest = strings.TrimSuffix(strings.TrimPrefix(line, "===MONITOR:"), "===") continue } if line == "===END===" { flush() continue } // 解析 key:value if strings.HasPrefix(line, "status:") { currentStatus = strings.TrimPrefix(line, "status:") continue } if strings.HasPrefix(line, "duration:") { durStr := strings.TrimPrefix(line, "duration:") ms, _ := strconv.ParseInt(strings.TrimSuffix(durStr, "ms"), 10, 64) currentDuration = fmt.Sprintf("%dms", ms) continue } if strings.HasPrefix(line, "score:") { currentScore, _ = strconv.ParseFloat(strings.TrimPrefix(line, "score:"), 64) continue } if strings.HasPrefix(line, "metric:") { kv := strings.TrimPrefix(line, "metric:") if i := strings.Index(kv, "="); i > 0 { name := strings.TrimSpace(kv[:i]) val, _ := strconv.ParseFloat(strings.TrimSpace(kv[i+1:]), 64) if currentMetrics == nil { currentMetrics = map[string]float64{} } currentMetrics[name] = val } continue } if strings.HasPrefix(line, "output:") { currentOutput = append(currentOutput, strings.TrimPrefix(line, "output:")) continue } // 普通输出行 if currentTest != "" && line != "" { currentOutput = append(currentOutput, line) } } flush() } // ============================ // Docker 容器执行(工具与脚本均预装在镜像中) // ============================ // runInDocker 在 Docker 容器中执行压测(工具与脚本预装在镜像中) // 固定分发模式:确保目标机有镜像(无则上传本地 tar + docker load)→ 启动容器 → 解析结果 func (r *Runner) runInDocker() (string, error) { // 1. 确保目标机存在压测镜像(不存在则上传本地导出的 tar 并 load) if err := r.ensureImage(); err != nil { return "", err } // 2. 构建 Docker 运行命令(通过环境变量传入参数) // 固定挂载:/sys 只读(容器内是独立 sysfs,sensors 需读宿主机硬件传感器节点才能采集温度); // /etc/localtime 只读(容器时区与宿主机一致,dmesg -T 与采样时间戳才准确)。 containerName := fmt.Sprintf("auto-check-%d", time.Now().UnixNano()) dockerRunCmd := fmt.Sprintf("docker run --rm --name %s "+ "-v /sys:/sys:ro "+ "-v /etc/localtime:/etc/localtime:ro "+ "-e AUTOCHECK_TYPES=%s "+ "-e AUTOCHECK_DURATION=%d "+ "-e AUTOCHECK_THREADS=%d "+ "-e AUTOCHECK_MEM_SIZE_MB=%d "+ "-e AUTOCHECK_DISK_SIZE_MB=%d "+ "-e AUTOCHECK_DISK_DIR=%s "+ "-e AUTOCHECK_TEMP_INTERVAL=%d "+ "-e AUTOCHECK_SAMPLE_INTERVAL=2 "+ "%s %s", containerName, r.cfg.Types, int(r.cfg.Duration.Seconds()), r.cfg.Threads, r.cfg.MemSizeMB, r.cfg.DiskSizeMB, r.cfg.DiskDir, int(r.cfg.TempLogInt.Seconds()), r.cfg.DockerArgs, r.cfg.DockerImage) fmt.Printf(" [Docker] 启动容器: %s\n", containerName) output, err := runPrivileged(r.client, dockerRunCmd, r.sudoPwd) if err != nil { return "", err } return output, nil } // ensureImage 确保目标机存在压测镜像(固定分发模式,不依赖 registry): // 镜像已存在则直接使用;否则把工程目录中导出的镜像 tar 通过 SFTP // 上传到目标机,再 docker load 载入。 func (r *Runner) ensureImage() error { imageExists, _ := runPrivileged(r.client, fmt.Sprintf("docker images %s --format '{{.Repository}}:{{.Tag}}'", r.cfg.DockerImage), r.sudoPwd) if strings.TrimSpace(imageExists) != "" { fmt.Printf(" [Docker] 镜像已存在: %s\n", r.cfg.DockerImage) return nil } // 本地导出的镜像文件(docker save -o xxx.tar ) if r.cfg.DockerTar == "" { return fmt.Errorf("目标机不存在镜像 %s 且未配置 docker_tar", r.cfg.DockerImage) } if _, err := os.Stat(r.cfg.DockerTar); err != nil { return fmt.Errorf("本地镜像文件不存在: %s(请先执行 docker save -o %s %s)", r.cfg.DockerTar, r.cfg.DockerTar, r.cfg.DockerImage) } fmt.Printf(" [Docker] 目标机镜像不存在,开始上传...\n") remoteTar := "/tmp/auto-check-image.tar" if err := sshclient.UploadFile(r.client, r.cfg.DockerTar, remoteTar); err != nil { return fmt.Errorf("镜像上传失败: %w", err) } loadOut, err := runPrivileged(r.client, fmt.Sprintf("docker load -i %s 2>&1", remoteTar), r.sudoPwd) sshclient.RunCommand(r.client, "rm -f "+remoteTar) // 清理远端临时 tar(忽略失败) if err != nil || strings.Contains(loadOut, "Error") { return fmt.Errorf("docker load 失败: %s", truncate(loadOut, 300)) } fmt.Printf(" [Docker] 镜像载入成功: %s\n", r.cfg.DockerImage) return nil } // ============================ // 业务入口(workflow 调用) // ============================ // TestDevice 对单台设备执行 SSH 登录 + 压力测试(业务入口) func TestDevice(ip string, cfg config.Config) *Report { sshCfg := sshclient.NewConfig(cfg.SSH.User, cfg.SSH.Password, cfg.SSH.KeyFile, cfg.SSH.Port, cfg.Scan.Timeout) sshCfg.SudoPwd = cfg.SSH.SudoPwd conn, err := sshclient.Connect(ip, sshCfg) if err != nil { return NewSSHFailReport(ip, err) } defer conn.Close() // fnOS 的 admin 用户 HOME 目录可能不存在,登录会输出 // "Could not chdir to home directory" 横幅污染命令输出,先建好 HOME sshclient.EnsureHome(conn) // sudo 密码:优先用专属 sudo_password;为空时退回 SSH 登录密码(飞牛 admin 同密码) sudoPwd := cfg.SSH.SudoPwd if sudoPwd == "" { sudoPwd = cfg.SSH.Password } // SSH 登录成功后回填远端 MAC(发现阶段只做了 ping,暂无 MAC) if mac := queryRemoteMAC(conn); mac != "" { if dev, ok := discovery.Devices[ip]; ok { dev.MAC = mac discovery.Devices[ip] = dev } fmt.Printf(" [MAC] %s -> %s\n", ip, mac) } // 由默认配置 + yaml 映射构造(避免硬编码,便于后续扩展字段) stressCfg := DefaultConfig() stressCfg.Types = cfg.Stress.Types stressCfg.Duration = cfg.Stress.Duration stressCfg.Threads = cfg.Stress.Threads stressCfg.MemSizeMB = cfg.Stress.MemSizeMB stressCfg.DiskSizeMB = cfg.Stress.DiskSizeMB stressCfg.DiskDir = cfg.Stress.DiskDir stressCfg.TempLogInt = cfg.Stress.TempInterval stressCfg.TempLimit = cfg.Stress.TempLimit // Docker 镜像配置(工具与脚本预装在镜像中,本地导出的 tar 上传到目标机后 load) stressCfg.DockerImage = cfg.Stress.DockerImage stressCfg.DockerTar = cfg.Stress.DockerTar stressCfg.DockerArgs = cfg.Stress.DockerArgs return NewRunner(stressCfg, ParseTypes(cfg.Stress.Types), conn, sudoPwd).Run(ip) } // NewSSHFailReport 构造 SSH 连接失败报告 func NewSSHFailReport(ip string, err error) *Report { return &Report{ IP: ip, StartTime: time.Now(), Results: []Result{{Type: "ssh", Status: StatusFail, Error: err.Error()}}, Failed: 1, } } // queryRemoteMAC 查询远端主机的 MAC 地址(取第一个有效的单播地址)。 // 失败或无有效地址时返回空字符串,不阻断主流程。 func queryRemoteMAC(client *ssh.Client) string { out, err := sshclient.RunCommand(client, "cat /sys/class/net/*/address 2>/dev/null") if err != nil { return "" } for _, line := range strings.Split(out, "\n") { mac := strings.TrimSpace(line) if mac == "" { continue } hw, e := net.ParseMAC(mac) if e != nil || len(hw) != 6 { continue } // 排除零地址、广播、组播 if hw[0] == 0 && hw[1] == 0 && hw[2] == 0 && hw[3] == 0 && hw[4] == 0 && hw[5] == 0 { continue } if hw[0]&0x01 != 0 { continue } return hw.String() } return "" } // ============================ // 工具函数 // ============================ // IsPassed 报告是否通过 func IsPassed(rpt *Report) bool { return rpt != nil && rpt.Failed == 0 && rpt.Errors == 0 } // IsSSHFail 报告是否为 SSH 连接失败(连不上,无任何有效测试项) func IsSSHFail(rpt *Report) bool { if rpt == nil || len(rpt.Results) != 1 { return false } r := rpt.Results[0] return r.Type == "ssh" && (r.Status == StatusFail || r.Status == StatusError) } // Status 报告状态字符串 func Status(rpt *Report) string { if rpt == nil { return StatusUnknown } if IsPassed(rpt) { return StatusPass } return StatusFail } // ParseTypes 逗号分隔字符串 → TestType 列表 func ParseTypes(s string) []TestType { var types []TestType for _, t := range strings.Split(s, ",") { if t = strings.TrimSpace(t); t != "" { types = append(types, TestType(t)) } } return types }