重构 auto-check:拆分包架构,精简 workflow 为纯调度循环
各包职责: - config: 分类 YAML 配置 - model: Device/IP/MAC 公共类型 - discovery: 扫描 + ARP MAC 解析 + Devices 变量 - sshclient: SSH 连接 - stress: stress-ng/stressapptest 压测 + Results 变量 - report: 80mm 热敏模板/sparkline/打印 + Printed 变量 - workflow: 永久循环调度 + Start/Stop workflow 已精简为 100 行纯编排,所有业务逻辑下沉各包
This commit is contained in:
@@ -6,161 +6,96 @@ import (
|
||||
"time"
|
||||
|
||||
"auto-check/pkg/config"
|
||||
"auto-check/pkg/discovery"
|
||||
"auto-check/pkg/report"
|
||||
"auto-check/pkg/stress"
|
||||
)
|
||||
|
||||
// ============================
|
||||
// Context — 阶段间共享数据
|
||||
// ============================
|
||||
|
||||
// Context 工作流上下文,保存阶段间传递的数据
|
||||
type Context struct {
|
||||
base context.Context // 可取消上下文
|
||||
|
||||
// 配置(由入口设置)
|
||||
Config config.Config
|
||||
|
||||
// 循环模式下的状态数据(key 均为设备 MAC)
|
||||
Devices map[string]*Device // 设备列表(MAC → 设备)
|
||||
Results map[string]*stress.Report // 测试结果(nil/失败=需测试)
|
||||
Printed map[string]string // 上次打印的测试状态
|
||||
}
|
||||
|
||||
// NewContext 创建上下文
|
||||
func NewContext() *Context {
|
||||
return &Context{
|
||||
base: context.Background(),
|
||||
Devices: make(map[string]*Device),
|
||||
Results: make(map[string]*stress.Report),
|
||||
Printed: make(map[string]string),
|
||||
}
|
||||
}
|
||||
|
||||
// WithCancel 支持取消
|
||||
func (c *Context) WithCancel() (context.CancelFunc, error) {
|
||||
ctx, cancel := context.WithCancel(c.base)
|
||||
c.base = ctx
|
||||
return cancel, nil
|
||||
}
|
||||
|
||||
// Done 返回取消信号
|
||||
func (c *Context) Done() <-chan struct{} {
|
||||
return c.base.Done()
|
||||
}
|
||||
|
||||
// Err 返回取消原因
|
||||
func (c *Context) Err() error {
|
||||
return c.base.Err()
|
||||
}
|
||||
|
||||
// ============================
|
||||
// Stage — 阶段接口
|
||||
// ============================
|
||||
|
||||
// Stage 工作流阶段
|
||||
type Stage interface {
|
||||
Name() string
|
||||
Execute(ctx *Context) error
|
||||
}
|
||||
|
||||
// ============================
|
||||
// Workflow — 流水线调度器
|
||||
// ============================
|
||||
|
||||
// OnStageError 阶段失败处理策略
|
||||
type OnStageError int
|
||||
|
||||
const (
|
||||
// StopOnError 失败即中断整个工作流
|
||||
StopOnError OnStageError = iota
|
||||
// ContinueOnError 失败继续执行后续阶段
|
||||
ContinueOnError
|
||||
)
|
||||
|
||||
// StageResult 阶段执行结果
|
||||
type StageResult struct {
|
||||
Stage string
|
||||
Start time.Time
|
||||
End time.Time
|
||||
Duration time.Duration
|
||||
Err error
|
||||
}
|
||||
|
||||
// Workflow Pipeline 工作流
|
||||
// Workflow 固定工作流:扫描 → 压测 → 打印状态,按间隔循环
|
||||
// 状态数据由各包自持(discovery/stress/report),workflow 直接调用其包级函数
|
||||
type Workflow struct {
|
||||
stages []Stage
|
||||
onError OnStageError
|
||||
results []StageResult
|
||||
startTime time.Time
|
||||
cfg config.Config
|
||||
base context.Context
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
// New 创建空工作流
|
||||
func New() *Workflow {
|
||||
return &Workflow{
|
||||
onError: StopOnError,
|
||||
// New 创建工作流
|
||||
func New(cfg config.Config) *Workflow {
|
||||
base, cancel := context.WithCancel(context.Background())
|
||||
return &Workflow{cfg: cfg, base: base, cancel: cancel}
|
||||
}
|
||||
|
||||
// Start 启动固定工作流循环(永久运行,直至取消)
|
||||
func (w *Workflow) Start() error {
|
||||
interval := w.cfg.Workflow.Interval
|
||||
if interval <= 0 {
|
||||
interval = 10 * time.Second
|
||||
}
|
||||
}
|
||||
|
||||
// SetOnError 设置失败策略
|
||||
func (w *Workflow) SetOnError(policy OnStageError) *Workflow {
|
||||
w.onError = policy
|
||||
return w
|
||||
}
|
||||
fmt.Printf("══ 工作流启动(间隔 %v,设备以 IP 为 key)══\n\n", interval)
|
||||
w.round() // 首轮立即执行
|
||||
|
||||
// AddStage 追加阶段(串行)
|
||||
func (w *Workflow) AddStage(stage Stage) *Workflow {
|
||||
w.stages = append(w.stages, stage)
|
||||
return w
|
||||
}
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
// Run 依次执行所有阶段
|
||||
func (w *Workflow) Run(ctx *Context) error {
|
||||
w.startTime = time.Now()
|
||||
fmt.Printf("══ 工作流启动 (%d 个阶段) ══\n\n", len(w.stages))
|
||||
|
||||
for i, stage := range w.stages {
|
||||
// 检查是否被取消
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
fmt.Printf(" [!] 工作流已取消,停止于阶段 %d/%d\n", i+1, len(w.stages))
|
||||
return ctx.Err()
|
||||
default:
|
||||
}
|
||||
|
||||
fmt.Printf("════ 阶段 %d/%d: %s ════\n", i+1, len(w.stages), stage.Name())
|
||||
start := time.Now()
|
||||
err := stage.Execute(ctx)
|
||||
duration := time.Since(start)
|
||||
|
||||
result := StageResult{
|
||||
Stage: stage.Name(),
|
||||
Start: start,
|
||||
End: time.Now(),
|
||||
Duration: duration,
|
||||
Err: err,
|
||||
}
|
||||
w.results = append(w.results, result)
|
||||
|
||||
if err != nil {
|
||||
fmt.Printf(" [✗] 阶段失败: %v\n", err)
|
||||
switch w.onError {
|
||||
case StopOnError:
|
||||
fmt.Printf("══ 工作流中止(策略: 失败即停止)══\n")
|
||||
return err
|
||||
case ContinueOnError:
|
||||
fmt.Printf(" [!] 继续执行后续阶段(策略: 失败继续)\n\n")
|
||||
continue
|
||||
}
|
||||
} else {
|
||||
fmt.Printf(" [✓] 阶段完成 (%v)\n\n", duration.Round(time.Millisecond))
|
||||
case <-w.base.Done():
|
||||
fmt.Println("══ 工作流已停止 ══")
|
||||
return nil
|
||||
case <-ticker.C:
|
||||
w.round()
|
||||
}
|
||||
}
|
||||
|
||||
fmt.Printf("══ 工作流完成,总耗时 %v ══\n", time.Since(w.startTime).Round(time.Millisecond))
|
||||
return nil
|
||||
}
|
||||
|
||||
// Results 阶段执行结果列表
|
||||
func (w *Workflow) Results() []StageResult {
|
||||
return w.results
|
||||
// Stop 停止工作流
|
||||
func (w *Workflow) Stop() {
|
||||
w.cancel()
|
||||
}
|
||||
|
||||
// round 单轮:扫描 → 压测 → 打印状态变化
|
||||
func (w *Workflow) round() {
|
||||
fmt.Printf("════ 轮次开始 [%s] ════\n", time.Now().Format("15:04:05"))
|
||||
|
||||
// 1. 扫描(discovery 包),更新设备表
|
||||
cfg := w.cfg
|
||||
discovery.NewScanner(cfg.Scan.CIDR, cfg.Scan.Timeout, cfg.Scan.Concurrency).Discover()
|
||||
fmt.Printf(" [扫描] 设备 %d 台\n", len(discovery.Devices))
|
||||
|
||||
// 2. 压测未通过设备(stress.TestDevice 执行,report 保存+打印)
|
||||
for ip, dev := range discovery.Devices {
|
||||
if stress.IsPassed(stress.Results[ip]) {
|
||||
continue // 已通过,跳过
|
||||
}
|
||||
fmt.Printf(" [测试] %s (%s) 开始压测...\n", ip, dev.MAC)
|
||||
rpt := stress.TestDevice(ip, w.cfg)
|
||||
stress.Results[ip] = rpt
|
||||
report.SaveAndPrintDeviceReport(&dev, rpt, cfg.Report.Path, cfg.Report.Print)
|
||||
}
|
||||
|
||||
// 3. 打印状态变化(首次或与上次不一致)
|
||||
for ip, rpt := range stress.Results {
|
||||
dev, ok := discovery.Devices[ip]
|
||||
if !ok {
|
||||
continue // 设备本轮缺席,状态保留等它回来
|
||||
}
|
||||
status := stress.Status(rpt)
|
||||
if prev, ok := report.Printed[ip]; !ok || prev != status {
|
||||
fmt.Printf(" [状态] %s %s → %s\n", dev.Label(), statusIcon(status), status)
|
||||
report.Printed[ip] = status
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// statusIcon 图标
|
||||
func statusIcon(s string) string {
|
||||
switch s {
|
||||
case "pass":
|
||||
return "✓"
|
||||
case "fail":
|
||||
return "✗"
|
||||
default:
|
||||
return "○"
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user