diff --git a/agent/app/service/alert.go b/agent/app/service/alert.go index 1ccf62e1a..44f9a43ee 100644 --- a/agent/app/service/alert.go +++ b/agent/app/service/alert.go @@ -4,11 +4,8 @@ import ( "encoding/json" "fmt" "mime" - "sort" "strconv" "strings" - "sync" - "time" "github.com/1Panel-dev/1Panel/agent/app/dto" "github.com/1Panel-dev/1Panel/agent/app/model" @@ -20,12 +17,10 @@ import ( alertUtil "github.com/1Panel-dev/1Panel/agent/utils/alert" alertconfig "github.com/1Panel-dev/1Panel/agent/utils/alert_config" alertwebhook "github.com/1Panel-dev/1Panel/agent/utils/alert_webhook" - "github.com/1Panel-dev/1Panel/agent/utils/cmd" "github.com/1Panel-dev/1Panel/agent/utils/copier" "github.com/1Panel-dev/1Panel/agent/utils/email" "github.com/1Panel-dev/1Panel/agent/utils/xpack" "github.com/1Panel-dev/1Panel/agent/utils/xpack/providers" - "github.com/shirou/gopsutil/v4/disk" ) type AlertService struct{} @@ -358,122 +353,14 @@ func (a AlertService) UpdateStatus(id uint, status string) error { } func (a AlertService) GetDisks() ([]dto.DiskDTO, error) { - var disks []dto.DiskDTO - excludes := map[string]struct{}{ - "/mnt/cdrom": {}, "/boot": {}, "/boot/efi": {}, "/dev": {}, "/dev/shm": {}, - "/run/lock": {}, "/run": {}, "/run/shm": {}, "/run/user": {}, + infos := loadDiskInfo(true) + disks := make([]dto.DiskDTO, 0, len(infos)) + for _, item := range infos { + disks = append(disks, dto.DiskDTO(item)) } - stdout, err := executeDiskCommand() - if err != nil { - return disks, nil - } - - lines := strings.Split(stdout, "\n") - var mounts []dto.AlertDiskInfo - - for _, line := range lines { - fields := strings.Fields(line) - if len(fields) < 7 { - continue - } - mountPoint := strings.Join(fields[6:], " ") - if shouldExclude(fields, mountPoint, excludes) { - continue - } - mounts = append(mounts, dto.AlertDiskInfo{Type: fields[1], Device: fields[0], Mount: mountPoint}) - - } - - var ( - wg sync.WaitGroup - mu sync.Mutex - ) - wg.Add(len(mounts)) - for i := 0; i < len(mounts); i++ { - go func(timeoutCh <-chan time.Time, mount dto.AlertDiskInfo) { - defer wg.Done() - - var itemData dto.DiskDTO - itemData.Path = mount.Mount - itemData.Type = mount.Type - itemData.Device = mount.Device - select { - case <-timeoutCh: - mu.Lock() - disks = append(disks, itemData) - mu.Unlock() - global.LOG.Errorf("load disk info from %s failed, err: timeout", mount.Mount) - default: - state, err := disk.Usage(mount.Mount) - if err != nil { - mu.Lock() - disks = append(disks, itemData) - mu.Unlock() - global.LOG.Errorf("load disk info from %s failed, err: %v", mount.Mount, err) - return - } - itemData.Total = state.Total - itemData.Free = state.Free - itemData.Used = state.Used - itemData.UsedPercent = state.UsedPercent - itemData.InodesTotal = state.InodesTotal - itemData.InodesUsed = state.InodesUsed - itemData.InodesFree = state.InodesFree - itemData.InodesUsedPercent = state.InodesUsedPercent - mu.Lock() - disks = append(disks, itemData) - mu.Unlock() - } - }(time.After(5*time.Second), mounts[i]) - } - wg.Wait() - - sort.Slice(disks, func(i, j int) bool { - return disks[i].Path < disks[j].Path - }) return disks, nil } -func executeDiskCommand() (string, error) { - cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(2 * time.Second)) - stdout, err := cmdMgr.RunWithStdout("df", "-hT", "-P") - if err != nil { - cmdMgr2 := cmd.NewCommandMgr(cmd.WithTimeout(1 * time.Second)) - stdout, err = cmdMgr2.RunWithStdout("df", "-lhT", "-P") - } - if err != nil { - return stdout, err - } - var lines []string - for _, line := range strings.Split(stdout, "\n") { - if !strings.Contains(line, "/") || strings.Contains(line, "tmpfs") || strings.Contains(line, "snap/core") || strings.Contains(line, "udev") { - continue - } - lines = append(lines, line) - } - if len(lines) == 0 { - return "", nil - } - return strings.Join(lines, "\n"), nil -} - -func shouldExclude(fields []string, mountPoint string, excludes map[string]struct{}) bool { - if strings.HasPrefix(mountPoint, "/snap") || len(strings.Split(mountPoint, "/")) > 10 { - return true - } - if strings.TrimSpace(fields[1]) == "tmpfs" { - return true - } - if strings.Contains(fields[2], "K") { - return true - } - if strings.Contains(mountPoint, "docker") { - return true - } - _, excluded := excludes[mountPoint] - return excluded -} - func (a AlertService) PageAlertLogs(search dto.AlertLogSearch) (int64, []dto.AlertLogDTO, error) { var ( opts []repo.DBOption diff --git a/agent/app/service/alert_helper.go b/agent/app/service/alert_helper.go index f9ca2e791..1fdcd0ed0 100644 --- a/agent/app/service/alert_helper.go +++ b/agent/app/service/alert_helper.go @@ -1066,50 +1066,37 @@ func processAllDisks(alert dto.AlertDTO) error { global.LOG.Errorf("error getting disk list, err: %v", err) return err } - var errMsgs []string for _, item := range diskList { - err := checkAndCreateDiskAlert(alert, item.Path) - if err != nil { - errMsg := fmt.Sprintf("disk path %s process failed: %v", item.Path, err) - errMsgs = append(errMsgs, errMsg) - global.LOG.Errorf("%s", errMsg) + if item.Total == 0 { continue } - } - if len(errMsgs) > 0 { - return fmt.Errorf("batch process disks failed, error count: %d, details: %s", len(errMsgs), strings.Join(errMsgs, "; ")) + checkAndCreateDiskAlert(alert, item.Path, &disk.UsageStat{Used: item.Used, UsedPercent: item.UsedPercent}) } return nil } func processSingleDisk(alert dto.AlertDTO) error { - err := checkAndCreateDiskAlert(alert, alert.Project) + usageStat, err := loadDiskUsageWithTimeout(alert.Project, true) if err != nil { - global.LOG.Errorf("%s", err.Error()) + global.LOG.Errorf("error getting disk usage for %s, err: %v", alert.Project, err) return err } + checkAndCreateDiskAlert(alert, alert.Project, usageStat) return nil } -func checkAndCreateDiskAlert(alert dto.AlertDTO, path string) error { - usageStat, err := psutil.DISK.GetUsage(path, false) - if err != nil { - global.LOG.Errorf("error getting disk usage for %s, err: %v", path, err) - return err - } - +func checkAndCreateDiskAlert(alert dto.AlertDTO, path string, usageStat *disk.UsageStat) { usedTotal, usedStr := calculateUsedTotal(alert.Cycle, usageStat) commonTotal := float64(alert.Count) if alert.Cycle == 1 { commonTotal *= 1024 * 1024 * 1024 } if usedTotal < commonTotal { - return nil + return } params := createAlertDiskParams(path, usedStr) sender := NewAlertSender(alert, alert.Project) sender.ResourceSend(path, params) - return nil } func calculateUsedTotal(cycle uint, usageStat *disk.UsageStat) (float64, string) { diff --git a/agent/app/service/dashboard.go b/agent/app/service/dashboard.go index 19268b1b1..49bf3e2d6 100644 --- a/agent/app/service/dashboard.go +++ b/agent/app/service/dashboard.go @@ -242,7 +242,7 @@ func (u *DashboardService) LoadCurrentInfo(ioOption string, netOption string) *d currentInfo.SwapMemoryUsed = swapInfo.Used currentInfo.SwapMemoryUsedPercent = swapInfo.UsedPercent - currentInfo.DiskData = loadDiskInfo() + currentInfo.DiskData = loadDiskInfo(false) currentInfo.GPUData, currentInfo.NPUData, currentInfo.XPUData = loadAcceleratorInfo() if ioOption == "all" { @@ -455,26 +455,9 @@ type diskInfo struct { Device string } -func loadDiskInfo() []dto.DiskInfo { +func loadDiskInfo(forceRefresh bool) []dto.DiskInfo { var datas []dto.DiskInfo - cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(2 * time.Second)) - format := `NR>1 && !/tmpfs|snap\/core|udev/ {printf "%s\t%s\t%s\t%s\t%s\t%s\t%s\n", $1, $2, $3, $4, $5, $6, $7}` - stdout, err := cmdMgr.RunPipe( - cmd.PipeCommand{Name: "df", Args: []string{"-hT", "-P"}}, - cmd.PipeCommand{Name: "awk", Args: []string{format}}, - ) - if err != nil { - global.LOG.Errorf("load disk info with df -hT -P failed, err: %v", err) - cmdMgr2 := cmd.NewCommandMgr(cmd.WithTimeout(1 * time.Second)) - stdout, err = cmdMgr2.RunPipe( - cmd.PipeCommand{Name: "df", Args: []string{"-lhT", "-P"}}, - cmd.PipeCommand{Name: "awk", Args: []string{format}}, - ) - if err != nil { - global.LOG.Errorf("load disk info with df -lhT -P failed, err: %v", err) - return datas - } - } + stdout := loadDiskMounts() lines := strings.Split(stdout, "\n") var mounts []diskInfo @@ -522,43 +505,22 @@ func loadDiskInfo() []dto.DiskInfo { itemData.Type = mount.Type itemData.Device = mount.Device - type diskResult struct { - state *disk.UsageStat - err error - } - resultCh := make(chan diskResult, 1) - - go func() { - state, err := psutil.DISK.GetUsage(mount.Mount, false) - resultCh <- diskResult{state: state, err: err} - }() - - select { - case <-time.After(5 * time.Second): - mu.Lock() - datas = append(datas, itemData) - mu.Unlock() - global.LOG.Errorf("load disk info from %s failed, err: timeout", mount.Mount) - case result := <-resultCh: - if result.err != nil { - mu.Lock() - datas = append(datas, itemData) - mu.Unlock() - global.LOG.Errorf("load disk info from %s failed, err: %v", mount.Mount, result.err) - return - } - itemData.Total = result.state.Total - itemData.Free = result.state.Free - itemData.Used = result.state.Used - itemData.UsedPercent = result.state.UsedPercent - itemData.InodesTotal = result.state.InodesTotal - itemData.InodesUsed = result.state.InodesUsed - itemData.InodesFree = result.state.InodesFree - itemData.InodesUsedPercent = result.state.InodesUsedPercent - mu.Lock() - datas = append(datas, itemData) - mu.Unlock() + state, err := loadDiskUsageWithTimeout(mount.Mount, forceRefresh) + if err != nil { + global.LOG.Errorf("load disk info from %s failed, err: %v", mount.Mount, err) + } else { + itemData.Total = state.Total + itemData.Free = state.Free + itemData.Used = state.Used + itemData.UsedPercent = state.UsedPercent + itemData.InodesTotal = state.InodesTotal + itemData.InodesUsed = state.InodesUsed + itemData.InodesFree = state.InodesFree + itemData.InodesUsedPercent = state.InodesUsedPercent } + mu.Lock() + datas = append(datas, itemData) + mu.Unlock() }(mounts[i]) } wg.Wait() @@ -569,6 +531,69 @@ func loadDiskInfo() []dto.DiskInfo { return datas } +var diskMountsMu sync.Mutex + +func loadDiskMounts() string { + if !diskMountsMu.TryLock() { + return "" + } + resultCh := make(chan string, 1) + go func() { + var stdout string + defer func() { + diskMountsMu.Unlock() + resultCh <- stdout + }() + cmdMgr := cmd.NewCommandMgr(cmd.WithTimeout(2 * time.Second)) + format := `NR>1 && !/tmpfs|snap\/core|udev/ {printf "%s\t%s\t%s\t%s\t%s\t%s\t%s\n", $1, $2, $3, $4, $5, $6, $7}` + output, err := cmdMgr.RunPipe( + cmd.PipeCommand{Name: "df", Args: []string{"-hT", "-P"}}, + cmd.PipeCommand{Name: "awk", Args: []string{format}}, + ) + if err != nil { + global.LOG.Errorf("load disk info with df -hT -P failed, err: %v", err) + cmdMgr2 := cmd.NewCommandMgr(cmd.WithTimeout(1 * time.Second)) + output, err = cmdMgr2.RunPipe( + cmd.PipeCommand{Name: "df", Args: []string{"-lhT", "-P"}}, + cmd.PipeCommand{Name: "awk", Args: []string{format}}, + ) + if err != nil { + global.LOG.Errorf("load disk info with df -lhT -P failed, err: %v", err) + return + } + } + + stdout = output + }() + timer := time.NewTimer(3 * time.Second) + defer timer.Stop() + select { + case stdout := <-resultCh: + return stdout + case <-timer.C: + global.LOG.Error("load disk mounts timed out; df collection is still running") + return "" + } +} + +func loadDiskUsageWithTimeout(path string, forceRefresh bool) (*disk.UsageStat, error) { + type diskResult struct { + state *disk.UsageStat + err error + } + resultCh := make(chan diskResult, 1) + go func() { + state, err := psutil.DISK.GetUsage(path, forceRefresh) + resultCh <- diskResult{state: state, err: err} + }() + select { + case <-time.After(5 * time.Second): + return nil, fmt.Errorf("load disk usage from %s: timeout", path) + case result := <-resultCh: + return result.state, result.err + } +} + func loadAcceleratorInfo() ([]dto.GPUInfo, []dto.NPUInfo, []dto.XPUInfo) { ok, client := accelerator.New() if !ok { diff --git a/agent/app/service/device.go b/agent/app/service/device.go index ed2f363ec..5c7f917de 100644 --- a/agent/app/service/device.go +++ b/agent/app/service/device.go @@ -73,7 +73,7 @@ func (u *DeviceService) LoadBaseInfo() (dto.DeviceBaseInfo, error) { if baseInfo.SwapMemoryTotal != 0 { baseInfo.SwapDetails = loadSwap() } - disks := loadDiskInfo() + disks := loadDiskInfo(false) for _, item := range disks { baseInfo.MaxSize += item.Free } diff --git a/agent/app/service/file.go b/agent/app/service/file.go index bf2f82fd8..4eada4eb2 100644 --- a/agent/app/service/file.go +++ b/agent/app/service/file.go @@ -1310,7 +1310,7 @@ func (f *FileService) BatchCheckFiles(req request.FilePathsCheck) []response.Exi } func (f *FileService) GetHostMount() []dto.DiskInfo { - return loadDiskInfo() + return loadDiskInfo(false) } func (f *FileService) GetUsersAndGroups() (*response.UserGroupResponse, error) {