package main import ( "bufio" "bytes" "context" "errors" "fmt" "io/fs" "log" "os" "os/exec" "os/signal" "path" "path/filepath" "sort" "strconv" "strings" "syscall" "time" ) const ( usersRoot = "/var/spanel/userdata" logRoot = "/var/log/spanel/loadmonitor" psPath = "/usr/bin/ps" psTimeout = 15 * time.Second interval = 5 * time.Minute userColWidth = 32 dirPerm = 0755 filePerm = 0644 tmpSuffix = ".tmp" ) type agg struct { cpuSum float64 memSum float64 mysqlCnt int peakCPU float64 peakCmd string } type dayAgg struct { cpuSum float64 memSum float64 mysqlSum float64 samples int peakCPU float64 peakCmd string } func main() { log.SetFlags(log.LstdFlags | log.LUTC) log.Println("spmonitor starting") if err := os.MkdirAll(logRoot, dirPerm); err != nil { log.Fatalf("ERROR creating log root %s: %v", logRoot, err) } if err := os.MkdirAll(usersRoot, dirPerm); err != nil { log.Printf("WARN: created missing users root %s", usersRoot) } ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer cancel() sleepUntilNextTick(ctx, interval) ticker := time.NewTicker(interval) defer ticker.Stop() for { runOnce(ctx) select { case <-ctx.Done(): log.Println("shutdown requested, exiting") return case <-ticker.C: } } } func sleepUntilNextTick(ctx context.Context, step time.Duration) { now := time.Now() next := now.Truncate(step).Add(step) d := time.Until(next) timer := time.NewTimer(d) defer timer.Stop() select { case <-ctx.Done(): case <-timer.C: } } func runOnce(ctx context.Context) { start := time.Now() defer func() { log.Printf("snapshot finished in %s", time.Since(start).Round(time.Millisecond)) }() if err := compactOldDays(); err != nil { log.Printf("ERROR compacting old days: %v", err) } valid, err := loadValidUsers(usersRoot) if err != nil { log.Printf("ERROR load users: %v", err) return } if len(valid) == 0 { log.Printf("WARN: no valid users under %s; writing empty snapshot", usersRoot) } out, err := runPS(ctx) if err != nil { log.Printf("ERROR ps: %v", err) } stats := parsePS(out, valid) destPath, err := buildDestPath(time.Now()) if err != nil { log.Printf("ERROR build path: %v", err) return } if err := writeSnapshot(destPath, stats); err != nil { log.Printf("ERROR write snapshot: %v", err) return } } func loadValidUsers(root string) (map[string]struct{}, error) { m := make(map[string]struct{}) entries, err := os.ReadDir(root) if err != nil { if errors.Is(err, fs.ErrNotExist) { return m, nil } return nil, err } for _, e := range entries { if e.IsDir() { name := e.Name() if name == "" || strings.HasPrefix(name, ".") { continue } m[name] = struct{}{} } } return m, nil } func runPS(parent context.Context) ([]byte, error) { ctx, cancel := context.WithTimeout(parent, psTimeout) defer cancel() cmd := exec.CommandContext(ctx, psPath, "-eo", "user:32,pcpu,pmem,pid,args", "--no-headers") cmd.Env = append(os.Environ(), "LC_ALL=C", "LANG=C") var stdout, stderr bytes.Buffer cmd.Stdout = &stdout cmd.Stderr = &stderr if err := cmd.Run(); err != nil { return nil, fmt.Errorf("%w: %s", err, strings.TrimSpace(stderr.String())) } return stdout.Bytes(), nil } func parsePS(out []byte, valid map[string]struct{}) map[string]*agg { res := make(map[string]*agg) sc := bufio.NewScanner(bytes.NewReader(out)) const maxLine = 1024 * 1024 buf := make([]byte, 64*1024) sc.Buffer(buf, maxLine) for sc.Scan() { line := sc.Text() if len(line) < userColWidth { continue } userField := strings.TrimSpace(line[:userColWidth]) rest := strings.TrimSpace(line[userColWidth:]) if userField == "" || rest == "" { continue } if _, ok := valid[userField]; !ok { continue } parts := strings.Fields(rest) if len(parts) < 4 { continue } pcpu, err1 := strconv.ParseFloat(parts[0], 64) pmem, err2 := strconv.ParseFloat(parts[1], 64) _, err3 := strconv.Atoi(parts[2]) if err1 != nil || err2 != nil || err3 != nil { continue } args := strings.Join(parts[3:], " ") args = sanitizeArgs(args) a := res[userField] if a == nil { a = &agg{} res[userField] = a } a.cpuSum += pcpu a.memSum += pmem if pcpu > a.peakCPU { a.peakCPU = pcpu a.peakCmd = args } if looksLikeMysqlClient(args) { a.mysqlCnt++ } } return res } func sanitizeArgs(s string) string { s = strings.ReplaceAll(s, "\n", " ") s = strings.ReplaceAll(s, "\t", " ") return strings.TrimSpace(s) } func looksLikeMysqlClient(args string) bool { l := strings.ToLower(strings.TrimSpace(args)) if strings.Contains(l, "mysqld") || strings.Contains(l, "mariadbd") { return false } first := l if i := strings.IndexByte(l, ' '); i > 0 { first = l[:i] } base := path.Base(first) return base == "mysql" || base == "mysqladmin" } func buildDestPath(t time.Time) (string, error) { year := t.Year() mon := t.Format("Jan") day := t.Day() dir := filepath.Join( logRoot, fmt.Sprintf("%04d", year), mon, fmt.Sprintf("%d", day), ) if err := os.MkdirAll(dir, dirPerm); err != nil { return "", err } return filepath.Join(dir, fmt.Sprintf("%d", t.Unix())), nil } func writeSnapshot(dest string, stats map[string]*agg) error { tmp := dest + tmpSuffix users := make([]string, 0, len(stats)) for u := range stats { users = append(users, u) } sort.Strings(users) f, err := os.OpenFile(tmp, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, filePerm) if err != nil { return err } w := bufio.NewWriterSize(f, 64*1024) for _, u := range users { a := stats[u] // user:cpu_avg:mem_avg:mysql_conn_count:cpu_peak_cmd line := fmt.Sprintf("%s:%.2f:%.2f:%d:%s\r\n", u, a.cpuSum, a.memSum, a.mysqlCnt, a.peakCmd) if _, err := w.WriteString(line); err != nil { _ = f.Close() _ = os.Remove(tmp) return err } } if err := w.Flush(); err != nil { _ = f.Close() _ = os.Remove(tmp) return err } if err := f.Sync(); err != nil { _ = f.Close() _ = os.Remove(tmp) return err } if err := f.Close(); err != nil { _ = os.Remove(tmp) return err } return os.Rename(tmp, dest) } // ---------- Daily compaction ---------- // compactOldDays scans logRoot/YYYY/Mon/D dirs older than "today" // If a day contains timestamp files and lacks a "daily" file it aggregates // them into one summary and deletes the raw snapshots func compactOldDays() error { today := time.Now() todayY := fmt.Sprintf("%04d", today.Year()) todayM := today.Format("Jan") todayD := fmt.Sprintf("%d", today.Day()) yearDirs, err := os.ReadDir(logRoot) if err != nil { if errors.Is(err, fs.ErrNotExist) { return nil } return err } for _, y := range yearDirs { if !y.IsDir() { continue } yName := y.Name() monDir := filepath.Join(logRoot, yName) monEntries, err := os.ReadDir(monDir) if err != nil { continue } for _, m := range monEntries { if !m.IsDir() { continue } mName := m.Name() dayDir := filepath.Join(monDir, mName) dayEntries, err := os.ReadDir(dayDir) if err != nil { continue } for _, d := range dayEntries { if !d.IsDir() { continue } dName := d.Name() if yName == todayY && mName == todayM && dName == todayD { continue } dir := filepath.Join(dayDir, dName) if err := compactOneDay(dir); err != nil { log.Printf("WARN: compact %s: %v", dir, err) } } } } return nil } func compactOneDay(dayPath string) error { dailyPath := filepath.Join(dayPath, "daily") if _, err := os.Stat(dailyPath); err == nil { return nil } entries, err := os.ReadDir(dayPath) if err != nil { return err } var snapFiles []string for _, e := range entries { if e.IsDir() { continue } name := e.Name() if name == "daily" { continue } if allDigits(name) { snapFiles = append(snapFiles, filepath.Join(dayPath, name)) } } if len(snapFiles) == 0 { return nil } perUser := make(map[string]*dayAgg) for _, f := range snapFiles { if err := aggregateSnapshotFile(f, perUser); err != nil { log.Printf("WARN: reading %s: %v", f, err) } } tmp := dailyPath + tmpSuffix out, err := os.OpenFile(tmp, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, filePerm) if err != nil { return err } w := bufio.NewWriterSize(out, 64*1024) users := make([]string, 0, len(perUser)) for u := range perUser { users = append(users, u) } sort.Strings(users) for _, u := range users { a := perUser[u] if a.samples == 0 { continue } avgCPU := a.cpuSum / float64(a.samples) avgMEM := a.memSum / float64(a.samples) avgSQL := a.mysqlSum / float64(a.samples) line := fmt.Sprintf("%s:%.2f:%.2f:%.2f:%s\r\n", u, avgCPU, avgMEM, avgSQL, a.peakCmd) if _, err := w.WriteString(line); err != nil { _ = out.Close() _ = os.Remove(tmp) return err } } if err := w.Flush(); err != nil { _ = out.Close() _ = os.Remove(tmp) return err } if err := out.Sync(); err != nil { _ = out.Close() _ = os.Remove(tmp) return err } if err := out.Close(); err != nil { _ = os.Remove(tmp) return err } if err := os.Rename(tmp, dailyPath); err != nil { _ = os.Remove(tmp) return err } for _, f := range snapFiles { _ = os.Remove(f) } return nil } func allDigits(s string) bool { if s == "" { return false } for i := 0; i < len(s); i++ { if s[i] < '0' || s[i] > '9' { return false } } return true } func aggregateSnapshotFile(path string, perUser map[string]*dayAgg) error { f, err := os.Open(path) if err != nil { return err } defer f.Close() sc := bufio.NewScanner(f) const maxLine = 1024 * 1024 buf := make([]byte, 64*1024) sc.Buffer(buf, maxLine) for sc.Scan() { line := strings.TrimSpace(sc.Text()) if line == "" { continue } content := line // user:cpu:mem:mysql:cmd... parts := strings.SplitN(content, ":", 5) if len(parts) < 5 { continue } user := parts[0] cpu, err1 := strconv.ParseFloat(parts[1], 64) mem, err2 := strconv.ParseFloat(parts[2], 64) mysql, err3 := strconv.ParseFloat(parts[3], 64) cmd := parts[4] if err1 != nil || err2 != nil || err3 != nil { continue } a := perUser[user] if a == nil { a = &dayAgg{} perUser[user] = a } a.cpuSum += cpu a.memSum += mem a.mysqlSum += mysql a.samples++ if cpu > a.peakCPU { a.peakCPU = cpu a.peakCmd = cmd } } return sc.Err() }