Files
membank/auto-check/pkg/stress/runner.go
12600k-rog-d4 b673ed62f9 docker
2026-08-24 00:35:30 +08:00

487 lines
14 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 运行命令(通过环境变量传入参数)
containerName := fmt.Sprintf("auto-check-%d", time.Now().UnixNano())
dockerRunCmd := fmt.Sprintf("docker run --rm --name %s "+
"-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 <image>
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
}