From 4cd520b18baab3e0f31d551476ef91bb6d69d482 Mon Sep 17 00:00:00 2001 From: nabil1440 <52530910+nabil1440@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:51:41 +0600 Subject: [PATCH 1/4] feat(agent): send the CPU, memory and disk use of each site, and conform to contract v0.5.0 Contract v0.5.0, "sites": one item for each Docker Compose project in the home folder of the server user. The CPU comes from the cgroup v2 usage_usec of the containers that both ticks saw, as a share of all the CPUs. The memory is memory.current minus inactive_file. A background walk measures the disk use of each folder at most each hour, at the idle I/O priority, and the next sample carries it one time. sites is null when the agent cannot read Docker or on cgroup v1, and [] when Docker runs and no project matches. The README pins contract v0.5.0. Closes #45 --- README.md | 4 +- internal/agent/clean.go | 27 +++ internal/agent/config.go | 2 +- internal/agent/outbox_test.go | 30 +++ internal/agent/wire/wire.go | 17 +- internal/metrics/diskusage_other.go | 13 ++ internal/metrics/diskusage_unix.go | 68 ++++++ internal/metrics/diskusage_unix_test.go | 97 ++++++++ internal/metrics/metrics.go | 20 +- internal/metrics/sites.go | 292 +++++++++++++++++++++++ internal/metrics/sites_test.go | 297 ++++++++++++++++++++++++ internal/metrics/sys_linux.go | 17 ++ internal/metrics/sys_other.go | 2 + 13 files changed, 880 insertions(+), 6 deletions(-) create mode 100644 internal/metrics/diskusage_other.go create mode 100644 internal/metrics/diskusage_unix.go create mode 100644 internal/metrics/diskusage_unix_test.go create mode 100644 internal/metrics/sites.go create mode 100644 internal/metrics/sites_test.go diff --git a/README.md b/README.md index 5c2fdaa..91d5260 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ Easy CLI tool for servers managed by FlyWP. -Conforms to the FlyWP monitoring agent contract v0.4.0. +Conforms to the FlyWP monitoring agent contract v0.5.0. ## Installation @@ -103,7 +103,7 @@ All arguments after the WP-CLI command (or after the command for `fly exec`) go ### Monitoring agent -`fly agent run` is the FlyWP monitoring agent. It runs all the time under systemd (`fly-agent.service`, as the server user, not root), and FlyWP installs it. Each minute it measures CPU, load, memory, swap, disk and network traffic, the pressure (PSI) and the disk activity. It reads the server each 10 seconds, so each minute also has its peaks. It sends the values and the server status (restart needed, waiting updates, OS, kernel, uptime, CPU count, Docker state and version) to FlyWP. It keeps unsent data on disk for up to 24 hours. FlyWP can update and restart the agent through it, without SSH. The agent does not need Docker. When Docker runs, the agent reads its socket with two requests only: `GET /version` and `GET /containers/json`. +`fly agent run` is the FlyWP monitoring agent. It runs all the time under systemd (`fly-agent.service`, as the server user, not root), and FlyWP installs it. Each minute it measures CPU, load, memory, swap, disk and network traffic, the pressure (PSI) and the disk activity. It reads the server each 10 seconds, so each minute also has its peaks. It also measures the CPU, the memory and the disk use of each site: each Docker Compose project in the home folder of the server user. It measures the disk use at most one time each hour, at the lowest I/O priority. It sends the values and the server status (restart needed, waiting updates, OS, kernel, uptime, CPU count, Docker state and version) to FlyWP. It keeps unsent data on disk for up to 24 hours. FlyWP can update and restart the agent through it, without SSH. The agent does not need Docker. When Docker runs, the agent reads its socket with two requests only: `GET /version` and `GET /containers/json`. It reads `FLY_AGENT_URL` (https), `FLY_AGENT_TOKEN` and `FLY_AGENT_SERVER_ID` from `/etc/fly/agent.env`, and keeps its state in `STATE_DIRECTORY` (`/var/lib/fly-agent`). diff --git a/internal/agent/clean.go b/internal/agent/clean.go index 46f5bd6..fdfe535 100644 --- a/internal/agent/clean.go +++ b/internal/agent/clean.go @@ -19,6 +19,8 @@ const ( maxStatusTextLen = 255 maxArchLen = 16 maxCPUCount = 4096 + maxSites = 1000 + maxDirectoryLen = 255 maxEventNameLen = 64 maxErrorLen = 2000 ) @@ -46,9 +48,34 @@ func cleanSample(s wire.Sample) wire.Sample { } { *v = clampIntPtr(*v) } + s.Sites = cleanSites(s.Sites) return s } +// cleanSites keeps at most maxSites items. It drops an item whose directory +// is empty or too long: a cut directory would match an other site. A nil +// slice stays nil, and an empty slice stays empty: they mean different things. +func cleanSites(sites []wire.Site) []wire.Site { + if sites == nil { + return nil + } + + out := make([]wire.Site, 0, min(len(sites), maxSites)) + for _, site := range sites { + if len(out) == maxSites { + break + } + if n := utf8.RuneCountInString(site.Directory); n == 0 || n > maxDirectoryLen { + continue + } + site.CPUPercent = clampPtr(site.CPUPercent, 0, 100) + site.MemoryUsedBytes = clampIntPtr(site.MemoryUsedBytes) + site.DiskUsedBytes = clampIntPtr(site.DiskUsedBytes) + out = append(out, site) + } + return out +} + func cleanStatus(s wire.Status) wire.Status { s.OS = truncate(s.OS, maxStatusTextLen) s.Kernel = truncate(s.Kernel, maxStatusTextLen) diff --git a/internal/agent/config.go b/internal/agent/config.go index 705a14e..afcbddf 100644 --- a/internal/agent/config.go +++ b/internal/agent/config.go @@ -1,6 +1,6 @@ // Package agent is the FlyWP monitoring agent: the long-running mode of fly // that "fly agent run" starts. It follows the FlyWP monitoring agent -// contract v0.4.0. +// contract v0.5.0. package agent import ( diff --git a/internal/agent/outbox_test.go b/internal/agent/outbox_test.go index 8fdcecd..ded606e 100644 --- a/internal/agent/outbox_test.go +++ b/internal/agent/outbox_test.go @@ -242,6 +242,36 @@ func TestCleanDockerStatus(t *testing.T) { } } +func TestCleanSites(t *testing.T) { + if s := cleanSample(wire.Sample{}); s.Sites != nil { + t.Errorf("sites = %v, want nil to stay nil (not known)", s.Sites) + } + s := cleanSample(wire.Sample{Sites: []wire.Site{}}) + data, err := json.Marshal(s) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(data), `"sites":[]`) { + t.Errorf("sample JSON = %s, want an empty list to stay empty (no project)", data) + } + + cpu, big := 250.0, uint64(math.MaxUint64) + many := []wire.Site{{Directory: ""}, {Directory: strings.Repeat("d", 256)}, {Directory: "example.com", CPUPercent: &cpu, MemoryUsedBytes: &big}} + for i := range 1200 { + many = append(many, wire.Site{Directory: fmt.Sprintf("site%d.com", i)}) + } + s = cleanSample(wire.Sample{Sites: many}) + if len(s.Sites) != maxSites { + t.Errorf("%d sites, want at most %d", len(s.Sites), maxSites) + } + if got := s.Sites[0]; got.Directory != "example.com" || *got.CPUPercent != 100 || *got.MemoryUsedBytes != math.MaxInt64 { + t.Errorf("first site = %+v, want example.com with its values in range, after the empty and the long directory", got) + } + if cpu != 250 { + t.Error("cleanSample() changed the value of the caller") + } +} + func TestCleanEventDropsACommandIDThatIsNotAULID(t *testing.T) { if e := cleanEvent(wire.Event{CommandID: "not-a-ulid"}); e.CommandID != "" { t.Errorf("command_id = %q, want it removed", e.CommandID) diff --git a/internal/agent/wire/wire.go b/internal/agent/wire/wire.go index 0247254..952f3cb 100644 --- a/internal/agent/wire/wire.go +++ b/internal/agent/wire/wire.go @@ -1,5 +1,5 @@ // Package wire holds the JSON bodies of the FlyWP monitoring agent contract -// v0.4.0: the requests that the agent sends and the replies that it reads. +// v0.5.0: the requests that the agent sends and the replies that it reads. package wire import ( @@ -88,6 +88,21 @@ type Sample struct { DiskWriteMaxBytesPerSecond *uint64 `json:"disk_write_max_bytes_per_second"` DiskReadMaxOpsPerSecond *uint64 `json:"disk_read_max_ops_per_second"` DiskWriteMaxOpsPerSecond *uint64 `json:"disk_write_max_ops_per_second"` + + // Sites holds one item for each Docker Compose project in the home + // folder of the server user (contract v0.5.0). nil (JSON null) means that + // the agent cannot read Docker; an empty, non-nil slice ([]) means that + // Docker runs and no project matches. + Sites []Site `json:"sites"` +} + +// Site is the use of one Docker Compose project in the minute. nil (JSON +// null) means "not known". DiskUsedBytes is set in one sample each hour. +type Site struct { + Directory string `json:"directory"` + CPUPercent *float64 `json:"cpu_percent"` + MemoryUsedBytes *uint64 `json:"memory_used_bytes"` + DiskUsedBytes *uint64 `json:"disk_used_bytes"` } // MetricsReply is the reply to POST /agent/v1/metrics. diff --git a/internal/metrics/diskusage_other.go b/internal/metrics/diskusage_other.go new file mode 100644 index 0000000..a903918 --- /dev/null +++ b/internal/metrics/diskusage_other.go @@ -0,0 +1,13 @@ +//go:build !unix + +package metrics + +import ( + "context" + "errors" + "runtime" +) + +func diskUsage(context.Context, string) (uint64, int, error) { + return 0, 0, errors.New("disk use is not measured on " + runtime.GOOS) +} diff --git a/internal/metrics/diskusage_unix.go b/internal/metrics/diskusage_unix.go new file mode 100644 index 0000000..4bcb7b5 --- /dev/null +++ b/internal/metrics/diskusage_unix.go @@ -0,0 +1,68 @@ +//go:build unix + +package metrics + +import ( + "context" + "io/fs" + "path/filepath" + "syscall" +) + +// diskUsage returns the space that the files under root take on the disk: the +// allocated blocks, as du shows them (st_blocks × 512). It does not follow a +// symbolic link, does not go into an other file system, and counts a file +// with more than one hard link one time. It skips what it cannot read, and +// returns the number of those skips. It stops when ctx is done. +func diskUsage(ctx context.Context, root string) (used uint64, skipped int, err error) { + var dev uint64 + seen := map[[2]uint64]bool{} + + err = filepath.WalkDir(root, func(path string, d fs.DirEntry, err error) error { + if ctxErr := ctx.Err(); ctxErr != nil { + return ctxErr + } + if err != nil { + if path == root { + return err + } + // A folder that cannot be read is skipped. Its own blocks + // were counted before. + skipped++ + return nil + } + + info, err := d.Info() + if err != nil { + skipped++ + return nil + } + st, ok := info.Sys().(*syscall.Stat_t) + if !ok { + skipped++ + return nil + } + + if path == root { + dev = uint64(st.Dev) + } else if uint64(st.Dev) != dev { + // A mount point of an other file system. + if d.IsDir() { + return filepath.SkipDir + } + return nil + } + if !d.IsDir() && st.Nlink > 1 { + key := [2]uint64{uint64(st.Dev), uint64(st.Ino)} + if seen[key] { + return nil + } + seen[key] = true + } + + used += uint64(st.Blocks) * 512 + return nil + }) + + return used, skipped, err +} diff --git a/internal/metrics/diskusage_unix_test.go b/internal/metrics/diskusage_unix_test.go new file mode 100644 index 0000000..115af84 --- /dev/null +++ b/internal/metrics/diskusage_unix_test.go @@ -0,0 +1,97 @@ +//go:build unix + +package metrics + +import ( + "context" + "os" + "path/filepath" + "strings" + "syscall" + "testing" +) + +// blocks returns the allocated bytes of a path, as du counts them. +func blocks(t *testing.T, path string) uint64 { + t.Helper() + var st syscall.Stat_t + if err := syscall.Lstat(path, &st); err != nil { + t.Fatal(err) + } + return uint64(st.Blocks) * 512 +} + +func TestDiskUsage(t *testing.T) { + root := filepath.Join(t.TempDir(), "example.com") + outside := t.TempDir() + for name, size := range map[string]int{"wp-config.php": 3000, "app/index.php": 100, "app/uploads/big.jpg": 200000} { + path := filepath.Join(root, name) + if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, []byte(strings.Repeat("x", size)), 0o644); err != nil { + t.Fatal(err) + } + } + // A hard link counts one time. A symbolic link to a big file outside + // the folder is not followed. + if err := os.Link(filepath.Join(root, "app/uploads/big.jpg"), filepath.Join(root, "big-link.jpg")); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(outside, "huge"), []byte(strings.Repeat("y", 1<<20)), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Symlink(filepath.Join(outside, "huge"), filepath.Join(root, "huge-link")); err != nil { + t.Fatal(err) + } + + var want uint64 + for _, p := range []string{"", "app", "app/uploads", "wp-config.php", "app/index.php", "app/uploads/big.jpg", "huge-link"} { + want += blocks(t, filepath.Join(root, p)) + } + + got, skipped, err := diskUsage(context.Background(), root) + if err != nil || skipped != 0 { + t.Fatalf("diskUsage() = %d, %d skipped, %v", got, skipped, err) + } + if got != want { + t.Errorf("diskUsage() = %d, want %d: the blocks of each file one time, without the target of the symbolic link", got, want) + } +} + +func TestDiskUsageSkipsWhatItCannotRead(t *testing.T) { + if os.Getuid() == 0 { + t.Skip("root can read each folder") + } + root := t.TempDir() + locked := filepath.Join(root, "locked") + if err := os.MkdirAll(locked, 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(locked, "secret"), []byte(strings.Repeat("s", 100000)), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Chmod(locked, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(locked, 0o755) }) + + got, skipped, err := diskUsage(context.Background(), root) + if err != nil { + t.Fatal(err) + } + if skipped != 1 || got != blocks(t, root)+blocks(t, locked) { + t.Errorf("diskUsage() = %d with %d skipped, want the two folders and 1 skip", got, skipped) + } +} + +func TestDiskUsageStops(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if _, _, err := diskUsage(ctx, t.TempDir()); err == nil { + t.Error("diskUsage() = nil error, want the error of the context") + } + if _, _, err := diskUsage(context.Background(), filepath.Join(t.TempDir(), "gone")); err == nil { + t.Error("diskUsage() = nil error, want an error for a folder that does not exist") + } +} diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index c8c1313..1febb53 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -1,7 +1,8 @@ // Package metrics measures a Linux server for the monitoring agent: CPU, // load, memory, swap, disk, network, pressure (PSI) and disk activity each -// minute, with the peaks of the minute from a reading each 10 seconds, and -// the status of the server (contract v0.4.0). It needs no root. +// minute, with the peaks of the minute from a reading each 10 seconds, the +// use of each site, and the status of the server (contract v0.5.0). It needs +// no root. package metrics import ( @@ -73,6 +74,9 @@ type Collector struct { release func() string aptCheck func(ctx context.Context) ([]byte, error) docker *dockerapi.Client + // home is the home folder of the server user, where the Docker Compose + // projects of the sites are. + home string // prev is the reading of the last tick, and readings are the readings // after it, oldest first. They give the windows of the next sample. @@ -91,6 +95,11 @@ type Collector struct { // not known, until it is known again. dockerWarned bool + // prevContainers are the CPU times of the containers at the last tick, + // and walker measures the disk use of the sites. + prevContainers *containerReadings + walker *diskWalker + updatesAt time.Time updatesKnown bool updatesTotal uint64 @@ -110,6 +119,8 @@ func New(root, stateDir string, log *slog.Logger) *Collector { minFirst: minFirstMinute, } c.docker = dockerapi.New(c.file("var/run/docker.sock")) + c.home = homeDir() + c.walker = newDiskWalker() // Take a reading now, so that the first sample has a CPU value for the // time since the start. The saved reading comes before it only when it is @@ -191,6 +202,10 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { } } + // The containers come right after the reading of the tick: their CPU + // times must be of the same moment. + sites := c.sites(context.Background(), cur) + c.refreshUpdates(context.Background()) var s wire.Sample @@ -223,6 +238,7 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { } setPeaks(&s, c.prev, c.readings, cur) + s.Sites = sites c.prev = &cur c.readings = nil diff --git a/internal/metrics/sites.go b/internal/metrics/sites.go new file mode 100644 index 0000000..722bf1b --- /dev/null +++ b/internal/metrics/sites.go @@ -0,0 +1,292 @@ +package metrics + +import ( + "bufio" + "bytes" + "context" + "fmt" + "os" + "os/user" + "path/filepath" + "slices" + "strconv" + "strings" + "sync" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" +) + +const ( + // workingDirLabel is the label that Docker Compose puts on each container + // of a project: the folder of the project. + workingDirLabel = "com.docker.compose.project.working_dir" + + // The disk use of each project folder is measured at most each hour, in + // the background, and a walk of one folder stops after 5 minutes. + sitesDiskEvery = time.Hour + sitesDiskTimeout = 5 * time.Minute +) + +// containerUsage is the CPU time of one container, in microseconds, at a tick. +type containerUsage struct { + project string + usec uint64 +} + +// containerReadings are the CPU times of the containers at one tick. +type containerReadings struct { + bootID string + at time.Time + usage map[string]containerUsage // by container id +} + +// homeDir returns the home folder of the user that runs the agent: the +// server user. It is "" when it is not known. +func homeDir() string { + if u, err := user.Current(); err == nil && u.HomeDir != "" { + return filepath.Clean(u.HomeDir) + } + if h, err := os.UserHomeDir(); err == nil { + return filepath.Clean(h) + } + return "" +} + +// sites measures each Docker Compose project in the home folder at the tick +// cur (contract v0.5.0). It returns nil when the agent cannot read Docker or +// the cgroups: Docker does not run, the socket refuses the agent, or the +// server has cgroup v1. It returns an empty, non-nil slice when no project +// matches. +func (c *Collector) sites(ctx context.Context, cur reading) []wire.Site { + if c.home == "" { + return nil + } + if _, err := os.Stat(c.file("sys/fs/cgroup/cgroup.controllers")); err != nil { + // cgroup v1, or no cgroups: the agent reads only cgroup v2. + return nil + } + containers, err := c.docker.Containers(ctx) + if err != nil { + c.log.Debug("listing the containers", "error", err) + return nil + } + + type project struct { + cpuUsec uint64 + cpuOK bool + mem uint64 + memOK bool + } + projects := map[string]*project{} + now := containerReadings{bootID: cur.BootID, at: cur.At, usage: map[string]containerUsage{}} + prev := c.prevContainers + usable := prev != nil && prev.bootID == cur.BootID && cur.At.After(prev.at) && cur.At.Sub(prev.at) <= maxAge + + for _, ct := range containers { + dir := filepath.Clean(ct.Labels[workingDirLabel]) + if ct.Labels[workingDirLabel] == "" || filepath.Dir(dir) != c.home { + continue + } + name := filepath.Base(dir) + p := projects[name] + if p == nil { + p = &project{} + projects[name] = p + } + + cg := c.cgroupDir(ct.ID) + if cg == "" { + continue + } + if usec, err := parseFile(c, filepath.Join(cg, "cpu.stat"), parseCPUStat); err == nil { + now.usage[ct.ID] = containerUsage{project: name, usec: usec} + // Only a container that both ticks saw, with the same id: a + // container that started or restarted in the minute is left out. + if old, ok := prev.lookup(ct.ID); usable && ok && old.project == name && usec >= old.usec { + p.cpuUsec += usec - old.usec + p.cpuOK = true + } + } + if mem, err := c.containerMemory(cg); err == nil { + p.mem += mem + p.memOK = true + } + } + c.prevContainers = &now + + cpus, cpuErr := parseFile(c, "proc/stat", parseCPUCount) + disk := c.walker.take() + out := make([]wire.Site, 0, len(projects)) + for _, name := range slices.Sorted(func(yield func(string) bool) { + for name := range projects { + if !yield(name) { + return + } + } + }) { + p := projects[name] + site := wire.Site{Directory: name} + if p.cpuOK && cpuErr == nil { + us := float64(cur.At.Sub(prev.at).Microseconds()) * float64(cpus) + v := min(float64(p.cpuUsec)/us*100, 100) + site.CPUPercent = &v + } + if p.memOK { + site.MemoryUsedBytes = &p.mem + } + if d, ok := disk[name]; ok { + site.DiskUsedBytes = &d + } + out = append(out, site) + } + + // Measure the disk use of the folders in the background, at most each + // hour. The result goes into the next sample. + folders := map[string]string{} + for name := range projects { + folders[name] = filepath.Join(c.home, name) + } + c.walker.start(folders, c.log) + + return out +} + +// lookup returns the reading of a container at the tick before. +func (r *containerReadings) lookup(id string) (containerUsage, bool) { + if r == nil { + return containerUsage{}, false + } + u, ok := r.usage[id] + return u, ok +} + +// cgroupDir returns the cgroup v2 folder of a container, relative to the +// root, or "". The systemd cgroup driver (the default of Ubuntu) and the +// cgroupfs driver use different folders. +func (c *Collector) cgroupDir(id string) string { + if id == "" || strings.ContainsAny(id, "/.") { + return "" + } + for _, dir := range []string{ + filepath.Join("sys/fs/cgroup/system.slice", "docker-"+id+".scope"), + filepath.Join("sys/fs/cgroup/docker", id), + } { + if _, err := os.Stat(c.file(dir, "cpu.stat")); err == nil { + return dir + } + } + return "" +} + +// containerMemory returns memory.current − inactive_file of a cgroup, as +// docker stats shows it on cgroup v2. The page cache that the kernel can drop +// is not used memory. A negative value counts as 0. +func (c *Collector) containerMemory(cg string) (uint64, error) { + current, err := parseFile(c, filepath.Join(cg, "memory.current"), func(data []byte) (uint64, error) { + return strconv.ParseUint(string(bytes.TrimSpace(data)), 10, 64) + }) + if err != nil { + return 0, err + } + inactive, err := parseFile(c, filepath.Join(cg, "memory.stat"), statValue("inactive_file")) + if err != nil { + return 0, err + } + return current - min(inactive, current), nil +} + +// parseCPUStat reads usage_usec from a cgroup v2 cpu.stat. +func parseCPUStat(data []byte) (uint64, error) { + return statValue("usage_usec")(data) +} + +// statValue returns a parser for the value of key in a cgroup file of +// "key value" lines. +func statValue(key string) func([]byte) (uint64, error) { + return func(data []byte) (uint64, error) { + s := bufio.NewScanner(bytes.NewReader(data)) + for s.Scan() { + k, v, ok := strings.Cut(s.Text(), " ") + if ok && k == key { + return strconv.ParseUint(strings.TrimSpace(v), 10, 64) + } + } + return 0, fmt.Errorf("no %s", key) + } +} + +// diskWalker measures the disk use of the project folders in a goroutine of +// its own, so that a large folder does not delay the samples. The collector +// takes each result one time. +type diskWalker struct { + mu sync.Mutex + running bool + last time.Time // the start of the last walk + results map[string]uint64 // by directory, not yet sent + // done is closed when the running walk ends. Tests wait for it. + done chan struct{} + // walk measures one folder. Tests replace it. + walk func(ctx context.Context, root string) (uint64, int, error) +} + +func newDiskWalker() *diskWalker { + return &diskWalker{walk: diskUsage} +} + +// start begins a walk of the folders, by directory, when no walk runs and the +// last walk started one hour ago or more. +func (w *diskWalker) start(folders map[string]string, log interface { + Warn(msg string, args ...any) +}) { + w.mu.Lock() + defer w.mu.Unlock() + if w.running || len(folders) == 0 || (!w.last.IsZero() && time.Since(w.last) < sitesDiskEvery) { + return + } + w.running, w.last = true, time.Now() + done := make(chan struct{}) + w.done = done + + go func() { + defer close(done) + // The walk runs at the lowest I/O priority, so that it does not + // slow the sites. The priority is a property of the thread: the + // thread stays locked, and ends with the goroutine. + lowerIOPriority() + + results := map[string]uint64{} + for _, name := range slices.Sorted(func(yield func(string) bool) { + for name := range folders { + if !yield(name) { + return + } + } + }) { + ctx, cancel := context.WithTimeout(context.Background(), sitesDiskTimeout) + used, skipped, err := w.walk(ctx, folders[name]) + cancel() + if err != nil { + log.Warn("cannot measure the disk use of a site; sending null", "directory", name, "error", err) + continue + } + if skipped > 0 { + log.Warn("some files of a site cannot be read; its disk use is lower than the real use", "directory", name, "skipped", skipped) + } + results[name] = used + } + + w.mu.Lock() + defer w.mu.Unlock() + w.results, w.running = results, false + }() +} + +// take returns the results of the last walk that ended, one time. +func (w *diskWalker) take() map[string]uint64 { + w.mu.Lock() + defer w.mu.Unlock() + r := w.results + w.results = nil + return r +} diff --git a/internal/metrics/sites_test.go b/internal/metrics/sites_test.go new file mode 100644 index 0000000..cbaa40b --- /dev/null +++ b/internal/metrics/sites_test.go @@ -0,0 +1,297 @@ +package metrics + +import ( + "context" + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/flywp/server-cli/internal/agent/wire" + "github.com/flywp/server-cli/internal/dockerapi" + "github.com/flywp/server-cli/internal/testutil" +) + +// fakeContainer is a running container of the fake Docker Engine. +type fakeContainer struct { + id, dir string +} + +// containerJSON is the answer of GET /containers/json for the containers. +func containerJSON(cs []fakeContainer) string { + var items []string + for _, c := range cs { + labels := "{}" + if c.dir != "" { + labels = fmt.Sprintf(`{"com.docker.compose.project.working_dir":%q,"com.docker.compose.service":"php"}`, c.dir) + } + items = append(items, fmt.Sprintf(`{"Id":%q,"Names":["/%s"],"Labels":%s}`, c.id, c.id, labels)) + } + return "[" + strings.Join(items, ",") + "]" +} + +// cgroup writes the cgroup v2 files of a container, with the systemd driver. +func (s *server) cgroup(id string, usageUsec, current, inactiveFile uint64) { + s.cgroupAt(filepath.Join("sys/fs/cgroup/system.slice", "docker-"+id+".scope"), usageUsec, current, inactiveFile) +} + +func (s *server) cgroupAt(dir string, usageUsec, current, inactiveFile uint64) { + s.write("sys/fs/cgroup/cgroup.controllers", "cpuset cpu io memory pids\n") + s.write(filepath.Join(dir, "cpu.stat"), fmt.Sprintf("usage_usec %d\nuser_usec 1\nsystem_usec 2\n", usageUsec)) + s.write(filepath.Join(dir, "memory.current"), fmt.Sprintf("%d\n", current)) + s.write(filepath.Join(dir, "memory.stat"), fmt.Sprintf("anon 1\nfile 2\nactive_file 3\ninactive_file %d\n", inactiveFile)) +} + +// sitesCollector returns a collector whose home is /home/fly under the root, +// and whose Docker Engine lists the containers. The disk walk returns 4096 +// bytes for each folder. +func sitesCollector(t *testing.T, srv *server, cs []fakeContainer) *Collector { + t.Helper() + engine := testutil.NewEngine(t, testutil.EngineHandler("29.7.1", containerJSON(cs))) + c := srv.collector(t.TempDir()) + c.docker = dockerapi.New(engine.Socket) + c.home = filepath.Join(srv.root, "home/fly") + c.walker.walk = func(context.Context, string) (uint64, int, error) { return 4096, 0, nil } + return c +} + +// home is the folder of a project in the home folder of the fake server. +func (s *server) home(name string) string { + return filepath.Join(s.root, "home/fly", name) +} + +func TestSites(t *testing.T) { + srv := newServer(t) + // Four CPUs. + srv.write("proc/stat", "cpu 1000 0 0 800 0 0 0 0 0 0\ncpu0 1 2 3 4 5 6 7 8 0 0\ncpu1 1 2 3 4 5 6 7 8 0 0\ncpu2 1 2 3 4 5 6 7 8 0 0\ncpu3 1 2 3 4 5 6 7 8 0 0\n") + cs := []fakeContainer{ + {"a1", srv.home("example.com")}, + {"a2", srv.home("example.com")}, + {"b1", srv.home(".fly")}, + {"c1", "/srv/other"}, // not in the home folder + {"d1", ""}, // not a Compose container + {"e1", srv.home("example.com/nested/deep")}, // the parent is not the home folder + } + srv.cgroup("a1", 1000000, 300<<20, 100<<20) + srv.cgroup("a2", 2000000, 50<<20, 80<<20) // more inactive file than current: 0 + srv.cgroup("b1", 3000000, 400<<20, 0) + srv.cgroup("c1", 1, 1, 0) + srv.cgroup("e1", 1, 1, 0) + c := sitesCollector(t, srv, cs) + base := time.Now().Add(time.Second) + + first, err := c.Sample(base) + if err != nil { + t.Fatal(err) + } + if got := siteNames(first.Sites); got != "[.fly example.com]" { + t.Fatalf("sites = %s, want [.fly example.com]", got) + } + for _, site := range first.Sites { + if site.CPUPercent != nil { + t.Errorf("%s cpu = %v in the first sample, want null", site.Directory, *site.CPUPercent) + } + } + if m := first.Sites[1].MemoryUsedBytes; m == nil || *m != 200<<20 { + t.Errorf("example.com memory = %v, want %d", ptr(m), 200<<20) + } + + // One minute: a1 and a2 use 2.4 s of CPU, b1 uses 4.8 s. The server has + // 240 s of CPU in one minute. + srv.cgroup("a1", 1000000+1200000, 300<<20, 100<<20) + srv.cgroup("a2", 2000000+1200000, 50<<20, 80<<20) + srv.cgroup("b1", 3000000+4800000, 400<<20, 0) + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + want := map[string]float64{".fly": 2, "example.com": 1} + for _, site := range s.Sites { + if site.CPUPercent == nil || !near(*site.CPUPercent, want[site.Directory]) { + t.Errorf("%s cpu = %v, want %v", site.Directory, ptr(site.CPUPercent), want[site.Directory]) + } + } + if m := s.Sites[0].MemoryUsedBytes; m == nil || *m != 400<<20 { + t.Errorf(".fly memory = %v, want %d", ptr(m), 400<<20) + } +} + +func TestSitesLeaveOutAContainerThatRestarted(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1000000, 1, 0) + srv.cgroup("a2", 1000000, 1, 0) + engineCS := []fakeContainer{{"a1", srv.home("example.com")}, {"a2", srv.home("example.com")}} + c := sitesCollector(t, srv, engineCS) + base := time.Now().Add(time.Second) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + + // a2 restarted: it has a new id, and its CPU time starts again from 0. + srv.cgroup("a1", 1000000+600000, 1, 0) + srv.cgroup("a3", 50000000, 1, 0) + c.docker = dockerapi.New(testutil.NewEngine(t, testutil.EngineHandler("29.7.1", + containerJSON([]fakeContainer{{"a1", srv.home("example.com")}, {"a3", srv.home("example.com")}}))).Socket) + + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if len(s.Sites) != 1 || s.Sites[0].CPUPercent == nil || !near(*s.Sites[0].CPUPercent, 1) { + t.Errorf("sites = %+v, want example.com with 1%% from a1 only", s.Sites) + } +} + +func TestSitesWithTheCgroupfsDriver(t *testing.T) { + srv := newServer(t) + srv.cgroupAt("sys/fs/cgroup/docker/a1", 1000000, 10<<20, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + + s, err := c.Sample(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + if len(s.Sites) != 1 || s.Sites[0].MemoryUsedBytes == nil || *s.Sites[0].MemoryUsedBytes != 10<<20 { + t.Errorf("sites = %+v, want the memory of a1", s.Sites) + } +} + +func TestSitesNullAndEmpty(t *testing.T) { + t.Run("no project", func(t *testing.T) { + srv := newServer(t) + srv.cgroup("c1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"c1", "/srv/other"}}) + s, err := c.Sample(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + data, err := json.Marshal(s) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(data), `"sites":[]`) { + t.Errorf("sample JSON = %s, want \"sites\":[]: Docker runs, and no project matches", data) + } + }) + + for _, tt := range []struct { + name string + setup func(*testing.T, *server, *Collector) + }{ + {"Docker does not run", func(_ *testing.T, _ *server, c *Collector) { + c.docker = dockerapi.New("/nonexistent/docker.sock") + }}, + {"cgroup v1", func(t *testing.T, s *server, _ *Collector) { + if err := os.Remove(filepath.Join(s.root, "sys/fs/cgroup/cgroup.controllers")); err != nil { + t.Fatal(err) + } + }}, + {"no home folder", func(_ *testing.T, _ *server, c *Collector) { c.home = "" }}, + } { + t.Run(tt.name, func(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + tt.setup(t, srv, c) + s, err := c.Sample(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + data, err := json.Marshal(s) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(data), `"sites":null`) { + t.Errorf("sample JSON = %s, want \"sites\":null", data) + } + }) + } +} + +func TestSitesDiskIsSentOneTime(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + var walked []string + c.walker.walk = func(_ context.Context, root string) (uint64, int, error) { + walked = append(walked, root) + return 8192, 0, nil + } + base := time.Now().Add(time.Second) + + first, err := c.Sample(base) + if err != nil { + t.Fatal(err) + } + if first.Sites[0].DiskUsedBytes != nil { + t.Errorf("disk in the first sample = %d, want null: the walk runs in the background", *first.Sites[0].DiskUsedBytes) + } + waitWalk(t, c) + if len(walked) != 1 || walked[0] != srv.home("example.com") { + t.Errorf("walked %v, want the folder of example.com", walked) + } + + second, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if d := second.Sites[0].DiskUsedBytes; d == nil || *d != 8192 { + t.Errorf("disk in the second sample = %v, want 8192", ptr(d)) + } + + // One time only, and no new walk within one hour. + third, err := c.Sample(base.Add(2 * time.Minute)) + if err != nil { + t.Fatal(err) + } + waitWalk(t, c) + if third.Sites[0].DiskUsedBytes != nil || len(walked) != 1 { + t.Errorf("third sample disk = %v after %d walks, want null after 1 walk", ptr(third.Sites[0].DiskUsedBytes), len(walked)) + } +} + +func TestSitesDiskWalkThatFails(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + c.walker.walk = func(context.Context, string) (uint64, int, error) { return 0, 0, context.DeadlineExceeded } + + if _, err := c.Sample(time.Now().Add(time.Second)); err != nil { + t.Fatal(err) + } + waitWalk(t, c) + s, err := c.Sample(time.Now().Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.Sites[0].DiskUsedBytes != nil { + t.Errorf("disk = %d, want null after a walk that stopped", *s.Sites[0].DiskUsedBytes) + } +} + +// waitWalk waits until the disk walk that runs ends. +func waitWalk(t *testing.T, c *Collector) { + t.Helper() + c.walker.mu.Lock() + done := c.walker.done + c.walker.mu.Unlock() + if done == nil { + return + } + select { + case <-done: + case <-time.After(10 * time.Second): + t.Fatal("the disk walk did not end") + } +} + +func siteNames(sites []wire.Site) string { + var names []string + for _, s := range sites { + names = append(names, s.Directory) + } + return fmt.Sprint(names) +} diff --git a/internal/metrics/sys_linux.go b/internal/metrics/sys_linux.go index 2686e0d..b5dd8ec 100644 --- a/internal/metrics/sys_linux.go +++ b/internal/metrics/sys_linux.go @@ -1,6 +1,8 @@ package metrics import ( + "runtime" + "golang.org/x/sys/unix" ) @@ -32,3 +34,18 @@ func kernelRelease() string { return unix.ByteSliceToString(u.Release[:]) } + +// lowerIOPriority gives the calling goroutine the idle I/O class: its disk +// reads wait for all other reads. The class is a property of the thread, so +// the goroutine stays locked to its thread, and the thread ends with the +// goroutine. Without the right to set it, the priority does not change. +func lowerIOPriority() { + runtime.LockOSThread() + + const ( + whoProcess = 1 // IOPRIO_WHO_PROCESS: with id 0, the calling thread + classIdle = 3 // IOPRIO_CLASS_IDLE + classShift = 13 + ) + _, _, _ = unix.Syscall(unix.SYS_IOPRIO_SET, whoProcess, 0, classIdle< Date: Thu, 24 Sep 2026 12:03:45 +0600 Subject: [PATCH 2/4] fix(agent): keep the disk result of the sites, and leave out restarted containers - The sites come after each step of Sample that can fail, so a dropped sample does not take the hourly disk result with it. - Each folder walk runs in a goroutine of its own: a walk that hangs in a system call no longer stops the next walks after the 5 minute limit. - The CPU time of the containers uses the time of their own reads, not the time of the tick reading. - A container that restarted with the same id has a new cgroup: its CPU time is left out for that minute. - A home folder of / makes no site. --- internal/agent/outbox_test.go | 21 ++++++ internal/metrics/diskusage_other.go | 2 + internal/metrics/diskusage_unix.go | 9 +++ internal/metrics/metrics.go | 9 ++- internal/metrics/sites.go | 77 ++++++++++++++------ internal/metrics/sites_test.go | 104 ++++++++++++++++++++++++++++ 6 files changed, 197 insertions(+), 25 deletions(-) diff --git a/internal/agent/outbox_test.go b/internal/agent/outbox_test.go index ded606e..12eaa6e 100644 --- a/internal/agent/outbox_test.go +++ b/internal/agent/outbox_test.go @@ -272,6 +272,27 @@ func TestCleanSites(t *testing.T) { } } +func TestOutboxKeepsSitesNullAndEmptyApart(t *testing.T) { + dir := t.TempDir() + o := loadOutbox(dir, slog.New(slog.DiscardHandler)) + o.addSample(cleanSample(wire.Sample{Sites: nil})) + o.addSample(cleanSample(wire.Sample{Sites: []wire.Site{}})) + + o = loadOutbox(dir, slog.New(slog.DiscardHandler)) + if len(o.samples) != 2 { + t.Fatalf("samples = %d, want 2", len(o.samples)) + } + for i, want := range []string{`"sites":null`, `"sites":[]`} { + data, err := json.Marshal(o.samples[i]) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(data), want) { + t.Errorf("sample %d JSON = %s, want %s after the queue on disk", i, data, want) + } + } +} + func TestCleanEventDropsACommandIDThatIsNotAULID(t *testing.T) { if e := cleanEvent(wire.Event{CommandID: "not-a-ulid"}); e.CommandID != "" { t.Errorf("command_id = %q, want it removed", e.CommandID) diff --git a/internal/metrics/diskusage_other.go b/internal/metrics/diskusage_other.go index a903918..569424c 100644 --- a/internal/metrics/diskusage_other.go +++ b/internal/metrics/diskusage_other.go @@ -8,6 +8,8 @@ import ( "runtime" ) +func fileID(string) uint64 { return 0 } + func diskUsage(context.Context, string) (uint64, int, error) { return 0, 0, errors.New("disk use is not measured on " + runtime.GOOS) } diff --git a/internal/metrics/diskusage_unix.go b/internal/metrics/diskusage_unix.go index 4bcb7b5..583128d 100644 --- a/internal/metrics/diskusage_unix.go +++ b/internal/metrics/diskusage_unix.go @@ -9,6 +9,15 @@ import ( "syscall" ) +// fileID returns the inode of a path, or 0. +func fileID(path string) uint64 { + var st syscall.Stat_t + if err := syscall.Stat(path, &st); err != nil { + return 0 + } + return uint64(st.Ino) +} + // diskUsage returns the space that the files under root take on the disk: the // allocated blocks, as du shows them (st_blocks × 512). It does not follow a // symbolic link, does not go into an other file system, and counts a file diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 1febb53..b2c1b72 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -188,6 +188,7 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { if err != nil { return wire.Sample{}, err } + readAt := time.Now() // A tick right after the start has no minute to measure. Its reading // starts the next minute, which then has all its values. @@ -202,10 +203,6 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { } } - // The containers come right after the reading of the tick: their CPU - // times must be of the same moment. - sites := c.sites(context.Background(), cur) - c.refreshUpdates(context.Background()) var s wire.Sample @@ -238,7 +235,9 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { } setPeaks(&s, c.prev, c.readings, cur) - s.Sites = sites + // The sites come after each step that can fail: they take the result + // of the disk walk, which must not go with a sample that is dropped. + s.Sites = c.sites(context.Background(), cur, readAt) c.prev = &cur c.readings = nil diff --git a/internal/metrics/sites.go b/internal/metrics/sites.go index 722bf1b..e9edb2c 100644 --- a/internal/metrics/sites.go +++ b/internal/metrics/sites.go @@ -29,9 +29,12 @@ const ( ) // containerUsage is the CPU time of one container, in microseconds, at a tick. +// cgroup identifies the cgroup folder of the container: a restart keeps the +// id of the container, but makes a new cgroup whose CPU time starts again. type containerUsage struct { project string usec uint64 + cgroup uint64 } // containerReadings are the CPU times of the containers at one tick. @@ -54,12 +57,14 @@ func homeDir() string { } // sites measures each Docker Compose project in the home folder at the tick -// cur (contract v0.5.0). It returns nil when the agent cannot read Docker or +// cur, read at readAt (contract v0.5.0). It returns nil when the agent cannot read Docker or // the cgroups: Docker does not run, the socket refuses the agent, or the // server has cgroup v1. It returns an empty, non-nil slice when no project // matches. -func (c *Collector) sites(ctx context.Context, cur reading) []wire.Site { - if c.home == "" { +func (c *Collector) sites(ctx context.Context, cur reading, readAt time.Time) []wire.Site { + // Without a home folder, or with "/" as the home, no folder is a project + // of a site. + if c.home == "" || filepath.Dir(c.home) == c.home { return nil } if _, err := os.Stat(c.file("sys/fs/cgroup/cgroup.controllers")); err != nil { @@ -79,9 +84,13 @@ func (c *Collector) sites(ctx context.Context, cur reading) []wire.Site { memOK bool } projects := map[string]*project{} - now := containerReadings{bootID: cur.BootID, at: cur.At, usage: map[string]containerUsage{}} + // The CPU times are read now, some time after the tick reading: the + // request to Docker and apt-check come before. The time of the minute + // is the time between two such reads. + at := cur.At.Add(time.Since(readAt)) + now := containerReadings{bootID: cur.BootID, at: at, usage: map[string]containerUsage{}} prev := c.prevContainers - usable := prev != nil && prev.bootID == cur.BootID && cur.At.After(prev.at) && cur.At.Sub(prev.at) <= maxAge + usable := prev != nil && prev.bootID == cur.BootID && at.After(prev.at) && at.Sub(prev.at) <= maxAge for _, ct := range containers { dir := filepath.Clean(ct.Labels[workingDirLabel]) @@ -100,10 +109,12 @@ func (c *Collector) sites(ctx context.Context, cur reading) []wire.Site { continue } if usec, err := parseFile(c, filepath.Join(cg, "cpu.stat"), parseCPUStat); err == nil { - now.usage[ct.ID] = containerUsage{project: name, usec: usec} - // Only a container that both ticks saw, with the same id: a - // container that started or restarted in the minute is left out. - if old, ok := prev.lookup(ct.ID); usable && ok && old.project == name && usec >= old.usec { + id := fileID(c.file(cg)) + now.usage[ct.ID] = containerUsage{project: name, usec: usec, cgroup: id} + // Only a container that both ticks saw, with the same id and + // the same cgroup: a container that started or restarted in + // the minute is left out. + if old, ok := prev.lookup(ct.ID); usable && ok && old.project == name && old.cgroup == id && usec >= old.usec { p.cpuUsec += usec - old.usec p.cpuOK = true } @@ -128,7 +139,7 @@ func (c *Collector) sites(ctx context.Context, cur reading) []wire.Site { p := projects[name] site := wire.Site{Directory: name} if p.cpuOK && cpuErr == nil { - us := float64(cur.At.Sub(prev.at).Microseconds()) * float64(cpus) + us := float64(at.Sub(prev.at).Microseconds()) * float64(cpus) v := min(float64(p.cpuUsec)/us*100, 100) site.CPUPercent = &v } @@ -226,12 +237,13 @@ type diskWalker struct { results map[string]uint64 // by directory, not yet sent // done is closed when the running walk ends. Tests wait for it. done chan struct{} - // walk measures one folder. Tests replace it. - walk func(ctx context.Context, root string) (uint64, int, error) + // walk measures one folder, and timeout limits it. Tests replace them. + walk func(ctx context.Context, root string) (uint64, int, error) + timeout time.Duration } func newDiskWalker() *diskWalker { - return &diskWalker{walk: diskUsage} + return &diskWalker{walk: diskUsage, timeout: sitesDiskTimeout} } // start begins a walk of the folders, by directory, when no walk runs and the @@ -250,10 +262,6 @@ func (w *diskWalker) start(folders map[string]string, log interface { go func() { defer close(done) - // The walk runs at the lowest I/O priority, so that it does not - // slow the sites. The priority is a property of the thread: the - // thread stays locked, and ends with the goroutine. - lowerIOPriority() results := map[string]uint64{} for _, name := range slices.Sorted(func(yield func(string) bool) { @@ -263,9 +271,7 @@ func (w *diskWalker) start(folders map[string]string, log interface { } } }) { - ctx, cancel := context.WithTimeout(context.Background(), sitesDiskTimeout) - used, skipped, err := w.walk(ctx, folders[name]) - cancel() + used, skipped, err := w.walkOne(folders[name]) if err != nil { log.Warn("cannot measure the disk use of a site; sending null", "directory", name, "error", err) continue @@ -282,6 +288,37 @@ func (w *diskWalker) start(folders map[string]string, log interface { }() } +// walkOne measures one folder in a goroutine of its own, and stops waiting +// for it after the timeout. A walk that hangs in a system call, for example on +// a network mount that does not answer, stays behind; the other folders and +// the next walks go on. +func (w *diskWalker) walkOne(root string) (uint64, int, error) { + ctx, cancel := context.WithTimeout(context.Background(), w.timeout) + defer cancel() + + type result struct { + used uint64 + skipped int + err error + } + ch := make(chan result, 1) + go func() { + // The walk runs at the lowest I/O priority, so that it does not + // slow the sites. The priority is a property of the thread: the + // thread stays locked, and ends with the goroutine. + lowerIOPriority() + used, skipped, err := w.walk(ctx, root) + ch <- result{used, skipped, err} + }() + + select { + case r := <-ch: + return r.used, r.skipped, r.err + case <-ctx.Done(): + return 0, 0, ctx.Err() + } +} + // take returns the results of the last walk that ended, one time. func (w *diskWalker) take() map[string]uint64 { w.mu.Lock() diff --git a/internal/metrics/sites_test.go b/internal/metrics/sites_test.go index cbaa40b..e32b8eb 100644 --- a/internal/metrics/sites_test.go +++ b/internal/metrics/sites_test.go @@ -3,6 +3,7 @@ package metrics import ( "context" "encoding/json" + "errors" "fmt" "os" "path/filepath" @@ -272,6 +273,109 @@ func TestSitesDiskWalkThatFails(t *testing.T) { } } +func TestSitesLeaveOutAContainerThatRestartedWithTheSameID(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1000000, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + base := time.Now().Add(time.Second) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + + // docker restart keeps the id, but makes a new cgroup whose CPU time + // starts again. The new run already used more than the old one. + scope := filepath.Join(srv.root, "sys/fs/cgroup/system.slice/docker-a1.scope") + if err := os.Rename(scope, scope+".old"); err != nil { + t.Fatal(err) + } + srv.cgroup("a1", 1600000, 1, 0) + + s, err := c.Sample(base.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.Sites[0].CPUPercent != nil { + t.Errorf("cpu = %v, want null: the container restarted in the minute", *s.Sites[0].CPUPercent) + } +} + +func TestSitesWithTheRootAsHome(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", "/"}, {"a2", "/srv"}}) + c.home = "/" + s, err := c.Sample(time.Now().Add(time.Second)) + if err != nil { + t.Fatal(err) + } + if s.Sites != nil { + t.Errorf("sites = %+v, want null: with / as the home, no folder is a site", s.Sites) + } +} + +func TestSitesDiskResultSurvivesASampleThatFails(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + base := time.Now().Add(time.Second) + if _, err := c.Sample(base); err != nil { + t.Fatal(err) + } + waitWalk(t, c) + + // The tick that would carry the disk use fails. + statfs := c.statfs + c.statfs = func(string) (uint64, uint64, error) { return 0, 0, errors.New("statfs failed") } + if _, err := c.Sample(base.Add(time.Minute)); err == nil { + t.Fatal("Sample() = nil error, want the statfs error") + } + c.statfs = statfs + + s, err := c.Sample(base.Add(2 * time.Minute)) + if err != nil { + t.Fatal(err) + } + if d := s.Sites[0].DiskUsedBytes; d == nil || *d != 4096 { + t.Errorf("disk = %v, want 4096 in the next sample that goes", ptr(d)) + } +} + +func TestSitesDiskWalkThatHangs(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + srv.cgroup("b1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("a.com")}, {"b1", srv.home("b.com")}}) + // The walk of a.com hangs in a system call and does not see its + // context. The walk of b.com works. + hang := make(chan struct{}) + t.Cleanup(func() { close(hang) }) + c.walker.timeout = 50 * time.Millisecond + c.walker.walk = func(_ context.Context, root string) (uint64, int, error) { + if strings.HasSuffix(root, "a.com") { + <-hang + } + return 4096, 0, nil + } + + if _, err := c.Sample(time.Now().Add(time.Second)); err != nil { + t.Fatal(err) + } + waitWalk(t, c) + s, err := c.Sample(time.Now().Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if s.Sites[0].DiskUsedBytes != nil || s.Sites[1].DiskUsedBytes == nil { + t.Errorf("disk = %v, %v; want null for a.com and 4096 for b.com", ptr(s.Sites[0].DiskUsedBytes), ptr(s.Sites[1].DiskUsedBytes)) + } + c.walker.mu.Lock() + running := c.walker.running + c.walker.mu.Unlock() + if running { + t.Error("the walk still runs: a walk that hangs must not stop the next walks") + } +} + // waitWalk waits until the disk walk that runs ends. func waitWalk(t *testing.T, c *Collector) { t.Helper() From 6a67da118c3a721657d9d97649b2bafc524cc893 Mon Sep 17 00:00:00 2001 From: nabil1440 <52530910+nabil1440@users.noreply.github.com> Date: Thu, 24 Sep 2026 12:50:01 +0600 Subject: [PATCH 3/4] fix(agent): warn one time for each site folder with files that cannot be read On a FlyWP server the databases in ~/.fly belong to the container user, so the walk skips them each hour. The agent now warns one time for each directory in a process, and logs the later walks at debug level. --- internal/metrics/sites.go | 20 +++++++++--- internal/metrics/sites_test.go | 59 ++++++++++++++++++++++++++++++++++ 2 files changed, 74 insertions(+), 5 deletions(-) diff --git a/internal/metrics/sites.go b/internal/metrics/sites.go index e9edb2c..9447742 100644 --- a/internal/metrics/sites.go +++ b/internal/metrics/sites.go @@ -5,6 +5,7 @@ import ( "bytes" "context" "fmt" + "log/slog" "os" "os/user" "path/filepath" @@ -235,6 +236,9 @@ type diskWalker struct { running bool last time.Time // the start of the last walk results map[string]uint64 // by directory, not yet sent + // skipWarned holds the directories whose skipped files were logged as a + // warning. Only the walk goroutine uses it. + skipWarned map[string]bool // done is closed when the running walk ends. Tests wait for it. done chan struct{} // walk measures one folder, and timeout limits it. Tests replace them. @@ -243,14 +247,12 @@ type diskWalker struct { } func newDiskWalker() *diskWalker { - return &diskWalker{walk: diskUsage, timeout: sitesDiskTimeout} + return &diskWalker{walk: diskUsage, timeout: sitesDiskTimeout, skipWarned: map[string]bool{}} } // start begins a walk of the folders, by directory, when no walk runs and the // last walk started one hour ago or more. -func (w *diskWalker) start(folders map[string]string, log interface { - Warn(msg string, args ...any) -}) { +func (w *diskWalker) start(folders map[string]string, log *slog.Logger) { w.mu.Lock() defer w.mu.Unlock() if w.running || len(folders) == 0 || (!w.last.IsZero() && time.Since(w.last) < sitesDiskEvery) { @@ -277,7 +279,15 @@ func (w *diskWalker) start(folders map[string]string, log interface { continue } if skipped > 0 { - log.Warn("some files of a site cannot be read; its disk use is lower than the real use", "directory", name, "skipped", skipped) + // The same folders are skipped each hour, for example the + // databases in ~/.fly that belong to the container user: + // warn one time for each directory, then log at debug. + level := slog.LevelDebug + if !w.skipWarned[name] { + w.skipWarned[name] = true + level = slog.LevelWarn + } + log.Log(context.Background(), level, "some files of a site cannot be read; its disk use is lower than the real use", "directory", name, "skipped", skipped) } results[name] = used } diff --git a/internal/metrics/sites_test.go b/internal/metrics/sites_test.go index e32b8eb..221665c 100644 --- a/internal/metrics/sites_test.go +++ b/internal/metrics/sites_test.go @@ -5,9 +5,11 @@ import ( "encoding/json" "errors" "fmt" + "log/slog" "os" "path/filepath" "strings" + "sync" "testing" "time" @@ -376,6 +378,63 @@ func TestSitesDiskWalkThatHangs(t *testing.T) { } } +func TestSitesWarnOneTimeForFilesThatCannotBeRead(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home(".fly")}}) + rec := &levels{} + c.log = slog.New(rec) + c.walker.walk = func(context.Context, string) (uint64, int, error) { return 4096, 9, nil } + + for i := range 3 { + c.walker.mu.Lock() + c.walker.last = time.Time{} // the hour passed + c.walker.mu.Unlock() + if _, err := c.Sample(time.Now().Add(time.Duration(i+1) * time.Minute)); err != nil { + t.Fatal(err) + } + waitWalk(t, c) + } + + msg := "some files of a site cannot be read; its disk use is lower than the real use" + if got := rec.count(slog.LevelWarn, msg); got != 1 { + t.Errorf("%d warnings for 3 walks, want 1", got) + } + if got := rec.count(slog.LevelDebug, msg); got != 2 { + t.Errorf("%d debug lines for 3 walks, want 2", got) + } +} + +// levels is a slog handler that keeps the level and the message of each +// record. +type levels struct { + mu sync.Mutex + records []slog.Record +} + +func (l *levels) Enabled(context.Context, slog.Level) bool { return true } +func (l *levels) WithAttrs([]slog.Attr) slog.Handler { return l } +func (l *levels) WithGroup(string) slog.Handler { return l } + +func (l *levels) Handle(_ context.Context, r slog.Record) error { + l.mu.Lock() + defer l.mu.Unlock() + l.records = append(l.records, r) + return nil +} + +func (l *levels) count(level slog.Level, msg string) int { + l.mu.Lock() + defer l.mu.Unlock() + n := 0 + for _, r := range l.records { + if r.Level == level && r.Message == msg { + n++ + } + } + return n +} + // waitWalk waits until the disk walk that runs ends. func waitWalk(t *testing.T, c *Collector) { t.Helper() From 10544b594dea31098881854169bfc852863f1ad6 Mon Sep 17 00:00:00 2001 From: nabil1440 <52530910+nabil1440@users.noreply.github.com> Date: Thu, 24 Sep 2026 14:09:42 +0600 Subject: [PATCH 4/4] fix(agent): start the CPU times of the sites at a short first minute A first minute without a sample now also reads the containers, so the next sample has the CPU of each site, not null. --- internal/metrics/metrics.go | 4 ++++ internal/metrics/sites_test.go | 21 +++++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index b2c1b72..7f48fbe 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -196,6 +196,10 @@ func (c *Collector) Sample(now time.Time) (wire.Sample, error) { c.fromStart = false if d := cur.At.Sub(c.prev.At); d >= 0 && d < c.minFirst { c.prev, c.readings = &cur, nil + // The CPU times of the containers start the next minute too, + // so that its sites have a CPU value. A new process has no + // result of a disk walk to lose. + c.sites(context.Background(), cur, readAt) if err := statefile.Write(c.path, cur); err != nil { c.log.Warn("saving the counters", "error", err) } diff --git a/internal/metrics/sites_test.go b/internal/metrics/sites_test.go index 221665c..ec95d1d 100644 --- a/internal/metrics/sites_test.go +++ b/internal/metrics/sites_test.go @@ -13,6 +13,7 @@ import ( "testing" "time" + "github.com/flywp/server-cli/internal/agent" "github.com/flywp/server-cli/internal/agent/wire" "github.com/flywp/server-cli/internal/dockerapi" "github.com/flywp/server-cli/internal/testutil" @@ -435,6 +436,26 @@ func (l *levels) count(level slog.Level, msg string) int { return n } +func TestSitesHaveCPUAfterAShortFirstMinute(t *testing.T) { + srv := newServer(t) + srv.cgroup("a1", 1000000, 1, 0) + c := sitesCollector(t, srv, []fakeContainer{{"a1", srv.home("example.com")}}) + c.minFirst = minFirstMinute + + first := time.Now().Add(time.Second) + if _, err := c.Sample(first); !errors.Is(err, agent.ErrNoSample) { + t.Fatalf("Sample() 1 s after the start = %v, want ErrNoSample", err) + } + srv.cgroup("a1", 1600000, 1, 0) + s, err := c.Sample(first.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + if len(s.Sites) != 1 || s.Sites[0].CPUPercent == nil || !near(*s.Sites[0].CPUPercent, 1) { + t.Errorf("sites = %+v, want example.com with 1%%: the short minute starts the CPU times too", s.Sites) + } +} + // waitWalk waits until the disk walk that runs ends. func waitWalk(t *testing.T, c *Collector) { t.Helper()